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

אפשר לשנות את נתוני האירועים באמצעות כתיבת ביטויי שינוי באמצעות CEL. לדוגמה, אתם יכולים לשנות את מטען הייעודי (payload) של אירועים כדי לעמוד בדרישות של חוזה 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: a supported Eventarc Advanced location

      לחלופין, אפשר להגדיר את מאפיין המיקום של ה-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: a supported Eventarc Advanced location

      לחלופין, אפשר להגדיר את מאפיין המיקום של ה-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, הפונקציה נכשלת. שימו לב שאפשר להשתמש ב-toJsonString כדי להמיר את הנתונים למחרוזת JSON.

לדוגמה: 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

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

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"}]}"}'

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