טרנספורמציה של אירועים שהתקבלו

אתם יכולים לשנות את נתוני האירועים באמצעות כתיבת ביטויי שינוי באמצעות CEL. לדוגמה, אתם יכולים לשנות את מטען הנתונים של האירוע כדי לעמוד בדרישות של חוזה API ספציפי של יעד.

שימו לב: האירועים מועברים תמיד בפורמט CloudEvents באמצעות בקשת HTTP במצב תוכן בינארי, אלא אם מציינים message binding.

הגדרת פורמטים של נתוני קלט ופלט

בנוסף לכתיבת ביטוי טרנספורמציה ב-CEL, אפשר גם לציין את פורמט הנתונים של נתוני האירועים הנכנסים. כך Eventarc Advanced יודע איך לנתח את המטען הייעודי (payload) של האירוע. אפשר גם להמיר את הנתונים מפורמט אחד לפורמט אחר.

הפורמטים הנתמכים: Avro,‏ JSON ו-Protobuf. מידע נוסף זמין במאמר בנושא פורמט של אירועים שהתקבלו.

ביטויי טרנספורמציה

כשמבצעים טרנספורמציה של אירועים, אפשר לגשת לכל מאפייני האירוע בביטוי CEL כמשתנים דרך אובייקט message מוגדר מראש. המשתנים האלה מאוכלסים בערכים שמבוססים על נתוני האירועים בזמן הריצה. לדוגמה:

  • ‫message.id מחזירה את המאפיין id של האירוע
  • ‫message.data מחזירה ייצוג של מטען הייעודי (payload) של האירוע ב-CEL
  • ‫message.data.some-key מחזירה את התוכן של שדה בשם some-key מתוך המטען הייעודי (payload) של האירוע

השדות ב-message.data תמיד מיוצגים כסוגים של String, והערכים ממופים מהאירוע המקורי באמצעות הסכימה שצוינה כשמגדירים את פורמט נתוני הקלט.

ביטוי הטרנספורמציה צריך לבטא אירוע מלא שכולל את מאפייני הקשר של האירוע ואת מטען הנתונים של האירוע. הביטויים נכתבים ב-JSON, אבל יש תמיכה בפונקציות, בפקודות מאקרו ובאופרטורים מוגדרים מראש של CEL, וגם בביטויים רגולריים באמצעות RE2. בנוסף, Eventarc Advanced תומך בפונקציות הרחבה מסוימות שאפשר להשתמש בהן כדי לשנות את נתוני האירועים.

הנה שתי דוגמאות לשימוש בביטויי CEL כדי לשנות את נתוני האירועים. כאן אפשר לראות עוד תרחישים לדוגמה.

דוגמה: עיצוב ערכי מאפיינים

בדוגמה הבאה מוצגות פונקציות של ביטויים רגולריים שמשמשות לעיצוב ערכי המאפיין phone_number. (הושמטו מאפיינים אחרים).

  // Input:
  // {
  //   "data":
  //   {
  //     "email_address": "charlie@altostrat.com",
  //     "phone_number": "8005550100",
  //   }
  // }
  // Output:
  // {
  //    "data":
  //    {
  //      "email_domain": "altostrat.com",
  //      "phone_number": "(800) 555-0100",
  //      "area_code": "800",
  //      "local_number": "5550100",
  //    }
  // }

  {
    "data":
    {
      "email_domain": re.capture(
                        message.data.email_address,
                        "\\S+@(\\S+)"),

      "phone_number": re.extract(
                        message.data.phone_number,
                        "^(\\d{3})(\\d{3})(\\d{4})", "(\\1) \\2-\\3"
                      ),

    }.merge ( re.captureN(message.data.phone_number,
                        "^(?P\d{3})[\w\-)(]*(?P\d{7})"
                      )
    )
  }

אלה פונקציות הביטוי הרגולרי שבהן נעשה שימוש בדוגמה הקודמת:

  • ‫re.capture: הערך הראשון של קבוצה ללא שם או עם שם. הארגומנטים הם:
    • ‫target: מחרוזת שצריך לנתח
    • ‫regex: ביטוי רגולרי שמשמש ללכידת ערכים

    הפונקציה מחזירה מחרוזת של הערך של הקבוצה הראשונה שתועדה.

  • ‫re.captureN: מבצע התאמה מלאה למחרוזת ולביטוי הרגולרי שצוינו. הארגומנטים הם:
    • ‫target: מחרוזת שצריך לנתח
    • ‫regex: ביטוי רגולרי שמשמש ללכידת ערכים

    מחזירה מפה עם זוגות של מפתח וערך לקבוצה עם שם (שם הקבוצה, המחרוזת שנתפסה) או לקבוצה ללא שם (אינדקס הקבוצה, המחרוזת שנתפסה).

  • ‫re.extract: מתאימה ערכי קבוצות ממחרוזת היעד הנתונה וכותבת מחדש את המחרוזת. הארגומנטים הם:
    • target: מחרוזת שצריך לנתח
    • ‫regex: ביטוי רגולרי שמשמש לחילוץ ערכים
    • ‫rewrite: ביטוי רגולרי לפורמט שבו התוצאה צריכה להיות

    הפונקציה מחזירה מחרוזת של הערכים שחולצו, בפורמט שנקבע לפי הארגומנט rewrite.

דוגמה: מיפוי מערך למערך של אובייקטים

בדוגמה הבאה ממופה מערך של מספרים שלמים למערך של אובייקטים. (הושמטו מאפיינים אחרים).

  // Input:
  // {
  //   "data":
  //   {
  //        "product_ids": [1, 2, 3]
  //   }
  // }
  // Output:
  // {
  //    "data":
  //    {
  //             "products": [
  //                {
  //                   "name": "apple",
  //                   "price": 70
  //                },
  //                {
  //                    "name": "orange",
  //                    "price":  80
  //                },
  //                {
  //                    "name": "Product(3)",
  //                    "price": 0
  //                },
  //                {
  //                     "name": "apple",
  //                     "price": 70
  //                }
  //            ]
  //    }
  // }

  {
    "data":
    {
      "products":  message.data.product_ids.map(product_id,
              product_id == 1?
              {
                "name": "apple",
                "price": 70
              } :
              product_id == 2?
              {
                "name": "orange",
                "price":  80
              } :
              // Default:
              {
                "name": "Product(" + string(product_id) + ")",
                "price": 0
              }
          )
    }
  }

הגדרת צינור עיבוד נתונים לשינוי אירועים

אפשר להגדיר צינור לעיבוד נתונים כדי לשנות נתוני אירועים במסוף Google Cloud או באמצעות ה-CLI של gcloud.

חשוב לדעת: אפשר להשתמש רק בתהליך בחירת רשת אחד לכל צינור.

המסוף

  1. במסוף Google Cloud , נכנסים לדף Eventarc > Pipelines.

    מעבר אל Pipelines

  2. אפשר ליצור צינור או ללחוץ על שם הצינור כדי לעדכן אותו.

  3. בדף פרטי צינור, לוחצים על עריכה.

  4. בחלונית Event mediation:

    1. מסמנים את תיבת הסימון החלת טרנספורמציה.
    2. ברשימה Inbound format (פורמט נכנס), בוחרים את הפורמט הרלוונטי.

      מידע נוסף זמין במאמר בנושא פורמט של אירועים שהתקבלו.

    3. בשדה CEL expression (ביטוי CEL), כותבים ביטוי המרה בפורמט JSON. יש תמיכה בפונקציות, בפקודות מאקרו ובאופרטורים מוגדרים מראש של CEL, וגם בביטויים רגולריים. לדוגמה:

      {
      "id": message.id,
      "datacontenttype": "application/json",
      "data": "{ \"scrubbed\": \"true\" }"
      }

      בדוגמה שלמעלה:

      • כל המאפיינים מוסרים מהאירוע המקורי, חוץ מid
      • הגדרת המאפיין datacontenttype לערך application/json
      • החלפת המטען הייעודי (payload) של האירוע במחרוזת JSON סטטית
    4. לוחצים על Continue.

  5. בחלונית Destination (יעד):

    1. אם רלוונטי, בוחרים פורמט מהרשימה Outbound format (פורמט יוצא).

      מידע נוסף זמין במאמר בנושא פורמט של אירועים שהתקבלו.

    2. אפשר גם להחיל קישור הודעה. מידע נוסף מופיע בקטע הגדרת קישור של הודעה במסמך הזה.

  6. לוחצים על Save.

    יכולות לעבור כמה דקות עד שצינור יתעדכן.

gcloud

  1. פותחים טרמינל.

  2. אפשר ליצור צינור או לעדכן צינור באמצעות הפקודה gcloud eventarc pipelines update:

    gcloud eventarc pipelines update PIPELINE_NAME \
        --location=REGION \
        --mediations=transformation_template=\
    '
     {
        TRANSFORMATION_EXPRESSION
     }
    '

    מחליפים את מה שכתוב בשדות הבאים:

    • ‫PIPELINE_NAME: המזהה של צינור העיבוד או שם מוגדר במלואו
    • ‫REGION: מיקום נתמך של Eventarc Advanced

      לחלופין, אפשר להגדיר את מאפיין המיקום של ה-CLI של gcloud:

      gcloud config set eventarc/location REGION
      
    • ‫TRANSFORMATION_EXPRESSION: ביטוי שנכתב ב-JSON. יש תמיכה בפונקציות, בפקודות מאקרו ובאופרטורים מוגדרים מראש של CEL, וגם בביטויים רגולריים. הדגל mediations משמש להחלת מפתח transformation_template.

    יכולות לעבור כמה דקות עד שצינור יתעדכן.

    דוגמה:

    gcloud eventarc pipelines update my-pipeline \
        --location=us-central1 \
        --mediations=transformation_template=\
    '
     {
        "id": message.id,
        "datacontenttype": "application/json",
        "data": "{ \"scrubbed\": \"true\" }"
     }
    '

    בדוגמה שלמעלה:

    • כל המאפיינים מוסרים מהאירוע המקורי, חוץ מid
    • הגדרת המאפיין datacontenttype לערך application/json
    • החלפת המטען הייעודי (payload) של האירוע במחרוזת JSON סטטית

פונקציות של תוספים

התכונה Eventarc Advanced תומכת בפונקציות ההרחבה הבאות, שאפשר להשתמש בהן כדי לשנות את נתוני האירועים שמתקבלים דרך האוטובוס.

תפקיד תיאור
denormalize

הפונקציה מבצעת דה-נורמליזציה של מפה או רשימה על ידי הוספת נתונים מיותרים כדי לשפר את ביצועי הקריאה. שמות השדות במפה שמתקבלת מופרדים באמצעות נקודה (.). אינדקס הרשימה מומר למפתח מחרוזת, החל מ-0.

הערה: אי אפשר להשתמש בנקודה (.) בשמות של שדות ב-Avro וב-Protobuf, ולכן הפונקציה הזו מתאימה רק לנתוני JSON.

לדוגמה: map.() -> map(string, dyn) או list() -> map(string, dyn)

merge

הפונקציה מחברת בין שני שדות ומחזירה את השדה המשולב. שדות עם שמות כפולים ימוזגו.

לדוגמה: message.(message) -> message

removeFields

הסרה של שדות ספציפיים מאירוע. שמות השדות נפתרים כנתיבים. התו נקודה (.) משמש כתו מפריד.

שימו לב שצריך להזין JSON גולמי. אם מבצעים המרה של ה-JSON, יכול להיות שהטרנספורמציה תוחל על מחרוזת JSON ותגרום לשגיאה.

לדוגמה: message.(list(string)) -> message

setField

הוספה או החלפה של שדה באירוע עם מפתח נתון. שם השדה מפורש כנתיב. התו נקודה (.) משמש כתו מפריד.

לדוגמה: message.(string, dyn) -> message

דוגמה: הוספת מאפיין למטען הייעודי (payload) של אירוע בלי לשנות נתונים אחרים

// Input:
// {
//   "data": 
//   {
//        "credit_card_number": "XXXX-XXXX-XXXX-XXXX"
//   }
// }
// Output:
// {
//    "data":
//    {
//        "credit_card_number": "XXXX-XXXX-XXXX-XXXX",
//        "card_type": "credit"
//    }
// }
{
  "data": message.data.merge(
    {
      "card_type": "credit"
    }
  )
}

דוגמה: ביטול הנורמליזציה של רשימת פריטים ממטען ייעודי (payload) של אירוע

// Input:
//{
//"data": 
//   {
//        "products": [
//          {
//            "number": 021774,
//            "type": "perishable",
//            "price": 2.00
//          },
//          {
//            "number": 95602,
//            "type": "diy",
//            "price": 120.00
//          },
//          {
//            "number": 568302,
//            "type": "toys",
//            "price": 12.00
//          }
//        ]
//   }
//}
//
// Output:
//{
//"data":
//    {
//        "products": {
//            "0.number": 021774,
//            "0.type": "perishable",
//            "0.price": 2.00,
//            "1.number": 95602,
//            "1.type": "diy",
//            "1.price": 120.00,
//            "2.number": 568302,
//            "2.type": "toys",
//            "2.price": 12.00
//          }
//   }
//}
//
//
message.setField("data.products", message.data.products.denormalize())

דוגמה: הסרת שדה ממטען ייעודי (payload) של אירוע

// Input:
// {
//   "data": 
//   {
//     "payment": {
//       "card_number": "XXXX-XXXX-XXXX-XXXX",
//       "card_type": "credit",
//     }
//   }
// }
// Output:
// {
//   "data":
//   {
//     "payment": {
//       "card_type": "credit"
//     }
//   }
// }
message.removeFields(["data.payment.card_number"])

הגדרת קישור של הודעה

כברירת מחדל, האירועים מועברים תמיד ליעד בפורמט CloudEvents באמצעות בקשת HTTP במצב תוכן בינארי. אם רוצים, אפשר לשנות את שיטת הפעולה הזאת על ידי הגדרת קישור להודעה ויצירת בקשת HTTP חדשה.

כל כותרות ה-HTTP שנוספו על ידי כללי מדיניות או אמצעי בקרה אחרים (לדוגמה, אסימוני OAuth או OIDC) נשמרות וממוזגות עם הכותרות שנוצרות מהביטוי המחייב.

אפשר להגדיר קישור הודעה כשמגדירים פייפליין בGoogle Cloud מסוף או באמצעות ה-CLI של gcloud.

המסוף

  1. במסוף Google Cloud , נכנסים לדף Eventarc > Pipelines.

    מעבר אל Pipelines

  2. אפשר ליצור צינור או ללחוץ על שם הצינור כדי לעדכן אותו.

    הערה: יכול להיות שייקח יותר מ-10 דקות לעדכן את הצינור.

  3. בדף פרטי צינור, לוחצים על עריכה.

  4. בחלונית יעד, מחילים קישור להודעה, שהוא ביטוי CEL שנכתב ב-JSON. התוצאה היא בקשת HTTP חדשה שנוצרת ונשלחת ליעד של צינור הנתונים.

    מידע נוסף זמין בקטעים גישה להודעות נכנסות ויצירת בקשות HTTP במאמר הזה.

  5. לוחצים על Save.

gcloud

  1. פותחים טרמינל.

  2. אפשר ליצור צינור או לעדכן צינור באמצעות הפקודה gcloud eventarc pipelines update:

    ‫
    gcloud eventarc pipelines update PIPELINE_NAME \
        --location=REGION \
        --destinations=http_endpoint_message_binding_template='MESSAGE_BINDING'

    מחליפים את מה שכתוב בשדות הבאים:

    • ‫PIPELINE_NAME: המזהה של צינור העיבוד או שם מוגדר במלואו
    • ‫REGION: מיקום נתמך של Eventarc Advanced

      לחלופין, אפשר להגדיר את מאפיין המיקום של ה-CLI של gcloud:

      gcloud config set eventarc/location REGION
      
    • ‫MESSAGE_BINDING: ביטוי CEL שנכתב ב-JSON ויוצר בקשת HTTP חדשה שנשלחת ליעד של צינור הנתונים.

      מידע נוסף זמין בקטעים גישה להודעות נכנסות ויצירת בקשות HTTP במאמר הזה.

    דוגמה:

    gcloud eventarc pipelines create my-pipeline \
        --location=us-central1 \
        --destinations=http_endpoint_uri='https://example-endpoint.com', \
    http_endpoint_message_binding_template='{"headers":{"new-header-key": "new-header-value"}}'

    שימו לב: אם אתם משתמשים במפתח http_endpoint_message_binding_template, אתם צריכים להגדיר גם את המפתח http_endpoint_uri.

גישה להודעות נכנסות

אתם יכולים להשתמש בביטוי CEL כדי לגשת להודעת CloudEvents נכנסת באופן הבא:

  • משתמשים בערך message.data כדי לגשת לשדה data של ההודעה הנכנסת.
  • משתמשים בערכים message.key (כאשר key הוא שם המאפיין) כדי לגשת למאפיינים של ההודעה הנכנסת.
  • משתמשים במשתנה headers כדי לגשת לכל הכותרות שנוספו לבקשת ה-HTTP על ידי תהליכי גישור קודמים בשרשרת העיבוד. המשתנה הזה מגדיר מיפוי של צמדי מפתח/ערך שמתאימים לכותרות HTTP נוספות, ולא לכותרות המקוריות של הבקשה הנכנסת הראשונית.

    לדוגמה, אפשר להשתמש בביטוי ה-CEL הבא כדי ליצור בקשת HTTP עם כותרות בלבד, על ידי הוספת כותרת נוספת לכותרות שנוספו בתיווכים הקודמים של צינור העיבוד:

    {"headers": headers.merge({"new-header-key": "new-header-value"})}

יצירת בקשות HTTP

התוצאה של ביטוי CEL צריכה להיות מיפוי של צמדי מפתח/ערך, שבהם השדות headers ו-body משמשים ליצירת בקשת ה-HTTP באופן הבא.

ל-headers שדות:

  • אם מפת headers קיימת כתוצאה מביטוי CEL, זוגות המפתח-ערך שלה ממופים ישירות לכותרות של בקשת ה-HTTP, והערכים שלה נוצרים באמצעות קידוד מחרוזת קנוני של סוג הנתונים המתאים.
  • אם השדה headers לא קיים, בקשת ה-HTTP שמתקבלת שומרת את הכותרות שנוספו במהלך תיווך קודם של צינור הנתונים.

ל-body שדות:

  • אם שדה body קיים כתוצאה מביטוי CEL, הערך שלו ממופה ישירות לגוף של בקשת ה-HTTP.
  • אם ערך השדה body הוא מסוג bytes או string, הוא ישמש כגוף הבקשה ל-HTTP כמו שהוא. אחרת, הוא יומר למחרוזת JSON.
  • אם השדה body לא קיים, גוף בקשת ה-HTTP שמתקבלת הוא הגוף של הקישור הסופי של הודעת ה-HTTP של CloudEvents במצב תוכן בינארי.

המערכת מתעלמת מכל שדה אחר שנוצר כתוצאה מביטוי CEL.

פונקציות של תוספים

‫Eventarc Advanced תומך בפונקציות ההרחבה הבאות, שאפשר להשתמש בהן כדי לשנות את נתוני האירועים כשמציינים קשירת הודעות.

תפקיד תיאור
merge

הפונקציה ממזגת מפת CEL שהועברה למפת ה-CEL שאליה היא מוחלת. אם אותו מפתח קיים בשתי המפות, או אם הערך של המפתח הוא מסוג map, שתי המפות ימוזגו. אחרת, ייעשה שימוש בערך מהמפה שהועברה.

לדוגמה: map1.merge(map2) -> map3

toBase64

הפונקציה ממירה ערך CEL למחרוזת בקידוד Base64 לכתובת URL.

לדוגמה: map.toBase64() -> string

toCloudEventJsonWithPayloadFormat

הפונקציה ממירה הודעה למפה של CEL שתואמת לייצוג JSON של הודעת CloudEvents, ומחילת את toDestinationPayloadFormat על נתוני ההודעה. התג הזה גם מגדיר את הערך של datacontenttype של האירוע לפורמט היוצא שצוין (output_payload_format_*). אם לא מוגדר פורמט יוצא, נעשה שימוש בערך קיים של datacontenttype. אחרת, הערך של datacontenttype לא מוגדר. אם ההודעה לא עומדת במפרט של CloudEvents, הפונקציה נכשלת. שימו לב: כדי להמיר את הנתונים למחרוזת JSON, אפשר להשתמש ב-toJsonString.

לדוגמה: message.toCloudEventJsonWithPayloadFormat() -> map.toJsonString() -> string

toDestinationPayloadFormat

הפונקציה ממירה את message.data לפורמט היוצא שצוין (output_payload_format_*). אם לא מוגדר פורמט יוצא, הפונקציה מחזירה את message.data ללא שינוי.

לדוגמה: message.data.toDestinationPayloadFormat() -> string or bytes

toJsonString

הפונקציה ממירה ערך CEL למחרוזת JSON.

לדוגמה: map.toJsonString() -> string

toMap

הפונקציה ממירה רשימה של מפות CEL למפת CEL אחת.

לדוגמה: list(map).toMap() -> map

דוגמה: שמירת הכותרות, הוספת כותרת חדשה, הגדרת גוף ההודעה לפורמט היעד destination format

gcloud eventarc pipelines create my-pipeline \
    --location=us-central1 \
    --input-payload-format-json='{}' \
    --destinations=http_endpoint_uri='https://example-endpoint.com',http_endpoint_message_binding_template='{"headers": headers.merge({"content-type":"application/avro"}), "body": message.data.toDestinationPayloadFormat()}',output_payload_format_avro_schema_definition='{"schema_definition": "{"type":"record","name":"myrecord","fields":[{"name":"name","type":"string"},{"name":"account_late","type":"boolean"}]}"}'

המאמרים הבאים