תבנית Pub/Sub Proto to BigQuery

תבנית ה-proto מ-Pub/Sub ל-BigQuery היא צינור להעברת נתונים בזמן אמת, שקולט נתוני proto ממנוי Pub/Sub לטבלה ב-BigQuery. כל השגיאות שמתרחשות במהלך הכתיבה לטבלה ב-BigQuery מועברות בסטרימינג לנושא לא מעובד ב-Pub/Sub.

אפשר לספק פונקציה בהגדרת המשתמש (UDF) ב-JavaScript כדי לבצע טרנספורמציה של נתונים. שגיאות במהלך ההפעלה של ה-UDF יכולות להישלח לנושא נפרד ב-Pub/Sub או לאותו נושא לא מעובד כמו השגיאות ב-BigQuery.

לפני שמריצים צינור Dataflow לתרחיש הזה, כדאי לשקול אם מינוי ל-Pub/Sub BigQuery עם UDF עונה על הדרישות שלכם.

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

  • המינוי לקלט Pub/Sub חייב להתקיים.
  • קובץ הסכימה של רשומות ה-Proto צריך להיות קיים ב-Cloud Storage.
  • נושא ה-Pub/Sub של הפלט חייב להתקיים.
  • מערך הנתונים ב-BigQuery שבו יאוחסנו התוצאות חייב להתקיים.
  • אם הטבלה ב-BigQuery קיימת, היא צריכה לכלול סכימה שתואמת לנתוני הפרוטו, ללא קשר לערך של createDisposition.

פרמטרים של תבניות

פרמטרים נדרשים

  • protoSchemaPath (נתיב Cloud Storage לקובץ סכימת Proto): נתיב Cloud Storage לקובץ של קבוצת תיאורים עצמאיים. דוגמה: gs://MyBucket/schema.pb. אפשר ליצור את schema.pb על ידי הוספת --descriptor_set_out=schema.pb לפקודה protoc שמקמפלת את קובצי ה-proto. אפשר להשתמש בדגל --include_imports כדי להבטיח שהקובץ יהיה עצמאי.
  • fullMessageName (השם המלא של הודעת Proto): השם המלא של ההודעה (לדוגמה: package.name.MessageName). אם ההודעה מוטמעת בתוך הודעה אחרת, צריך לכלול את כל ההודעות עם התו המפריד '.' (לדוגמה: package.name.OuterMessage.InnerMessage). הערך של 'package.name' צריך להיות מההצהרה package ולא מההצהרה java_package.
  • inputSubscription (מינוי קלט של Pub/Sub): מינוי Pub/Sub לקריאת הקלט, בפורמט 'projects/your-project-id/subscriptions/your-subscription-name' (דוגמה: projects/your-project-id/subscriptions/your-subscription-name).
  • outputTableSpec (טבלת פלט של BigQuery): מיקום הטבלה ב-BigQuery שאליה ייכתב הפלט. השם צריך להיות בפורמט <project>:<dataset>.<table_name>. הסכימה של הטבלה צריכה להתאים לאובייקטים של הקלט.
  • outputTopic (נושא פלט Pub/Sub): שם הנושא שאליו יש לפרסם את הנתונים, בפורמט 'projects/your-project-id/topics/your-topic-name' (דוגמה: projects/your-project-id/topics/your-topic-name).

פרמטרים אופציונליים

  • preserveProtoFieldNames (שמירה של שמות שדות פרוטו): דגל שקובע אם לשמור את שמות שדות הפרוטו או להמיר אותם ל-lowerCamelCase. אם הטבלה כבר קיימת, צריך להשתמש בערך שמתאים לסכימה של הטבלה. אחרת, הוא יקבע את שמות העמודות של הטבלה שנוצרה. True כדי לשמור על snake_case בפרוטו. הערך False ימיר את השדות ל-lowerCamelCase. (ברירת מחדל: false).
  • bigQueryTableSchemaPath (נתיב לסכימת טבלה ב-BigQuery): נתיב ב-Cloud Storage לקובץ ה-JSON של סכימת BigQuery. אם לא מגדירים את זה, הסכימה נגזרת מסכימת ה-Proto. (דוגמה: gs://MyBucket/bq_schema.json).
  • udfOutputTopic (נושא פלט של Pub/Sub לבעיות ב-UDF): נושא פלט אופציונלי לשליחת בעיות ב-UDF. אם האפשרות הזו לא מוגדרת, הכשלים ייכתבו לאותו נושא כמו הכשלים ב-BigQuery. (דוגמה: projects/your-project-id/topics/your-topic-name).
  • writeDisposition (הגדרת כתיבה לשימוש ב-BigQuery): הגדרת כתיבה ב-BigQuery. לדוגמה, WRITE_APPEND, ‏ WRITE_EMPTY או WRITE_TRUNCATE. ברירת המחדל היא: WRITE_APPEND.
  • createDisposition (יצירת הגדרת מיקום לשימוש ב-BigQuery): הגדרת מיקום ליצירה ב-BigQuery. לדוגמה: CREATE_IF_NEEDED,‏ CREATE_NEVER. ברירת המחדל היא: CREATE_IF_NEEDED.
  • javascriptTextTransformGcsPath (נתיב Cloud Storage לקוד המקור של פונקציית UDF ב-JavaScript): תבנית הנתיב ב-Cloud Storage לקוד JavaScript שמכיל את הפונקציות בהגדרת המשתמש. (דוגמה: gs://your-bucket/your-function.js).
  • javascriptTextTransformFunctionName (שם פונקציית JavaScript בהגדרת המשתמש): השם של הפונקציה לקריאה מקובץ ה-JavaScript. אפשר להשתמש רק באותיות, ספרות וקווים תחתונים. (דוגמה: 'transform' או 'transform_udf1').
  • javascriptTextTransformReloadIntervalMinutes (מרווח הזמן לטעינה מחדש אוטומטית של פונקציית JavaScript UDF (בדקות)): הגדרה של מרווח הזמן שבו העובדים יכולים לבדוק אם יש שינויים בפונקציית JavaScript UDF כדי לטעון מחדש את הקבצים. ברירת המחדל היא 0.
  • useStorageWriteApi (שימוש ב-BigQuery Storage Write API): אם הערך הוא true, צינור הנתונים משתמש ב-Storage Write API כשכותבים את הנתונים ל-BigQuery (ראו https://cloud.google.com/blog/products/data-analytics/streaming-data-into-bigquery-using-storage-write-api). ערך ברירת המחדל הוא False. כשמשתמשים ב-Storage Write API במצב 'פעם אחת בדיוק', צריך להגדיר את הפרמטרים הבאים: 'מספר הזרמים ל-BigQuery Storage Write API' ו'תדירות ההפעלה בשניות ל-BigQuery Storage Write API'. אם מפעילים את מצב 'לפחות פעם אחת' ב-Dataflow או מגדירים את הפרמטר useStorageWriteApiAtLeastOnce לערך true, לא צריך להגדיר את מספר הסטרימינג או את תדירות ההפעלה.
  • useStorageWriteApiAtLeastOnce (שימוש בסמנטיקה של 'לפחות פעם אחת' ב-BigQuery Storage Write API): הפרמטר הזה תקף רק אם האפשרות 'שימוש ב-BigQuery Storage Write API' מופעלת. אם האפשרות הזו מופעלת, נעשה שימוש בסמנטיקה של מסירה לפחות פעם אחת עבור Storage Write API. אחרת, נעשה שימוש בסמנטיקה של מסירה חד-פעמית. ברירת המחדל היא: false.
  • numStorageWriteApiStreams (מספר הזרמים של BigQuery Storage Write API): מספר הזרמים מגדיר את המקביליות של פעולת הכתיבה של BigQueryIO, והוא תואם בערך למספר הזרמים של Storage Write API שייעשה בהם שימוש בצינור. הערכים המומלצים מפורטים בכתובת https://cloud.google.com/blog/products/data-analytics/streaming-data-into-bigquery-using-storage-write-api. ברירת המחדל היא 0.
  • storageWriteApiTriggeringFrequencySec (תדירות ההפעלה בשניות של BigQuery Storage Write API): תדירות ההפעלה תקבע כמה זמן יעבור עד שהנתונים יהיו זמינים להרצת שאילתות ב-BigQuery. הערכים המומלצים מפורטים בכתובת https://cloud.google.com/blog/products/data-analytics/streaming-data-into-bigquery-using-storage-write-api.

פונקציה בהגדרת המשתמש

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

מפרט הפונקציה

המאפיינים של פונקציית UDF:

  • קלט: שדה הנתונים של הודעת Pub/Sub, שעבר סריאליזציה כמחרוזת JSON.
  • Output: מחרוזת JSON שתואמת לסכימה של טבלת היעד ב-BigQuery.
  • הרצת התבנית

    המסוף

    1. עוברים לדף Create job from template (יצירת משימה מתבנית) ב-Dataflow.
    2. כניסה לדף Create job from template
    3. בשדה שם המשימה, מזינים שם ייחודי למשימה.
    4. אופציונלי: בשדה Regional endpoint (נקודת קצה אזורית), בוחרים ערך מהתפריט הנפתח. אזור ברירת המחדל הוא us-central1.

      רשימת האזורים שבהם אפשר להריץ משימת Dataflow מופיעה במאמר בנושא מיקומי Dataflow.

    5. בתפריט הנפתח Dataflow template (תבנית Dataflow), בוחרים בתבנית Pub/Sub Proto to BigQuery (פרוטו Pub/Sub ל-BigQuery).
    6. בשדות הפרמטרים שמופיעים, מזינים את ערכי הפרמטרים.
    7. לוחצים על הפעלת העבודה.

    gcloud

    במעטפת או בטרמינל, מריצים את התבנית:

    gcloud dataflow flex-template run JOB_NAME \
        --region=REGION_NAME \
        --template-file-gcs-location=gs://dataflow-templates-REGION_NAME/VERSION/flex/ \
        --parameters \
    schemaPath=SCHEMA_PATH,\
    fullMessageName=PROTO_MESSAGE_NAME,\
    inputSubscription=SUBSCRIPTION_NAME,\
    outputTableSpec=BIGQUERY_TABLE,\
    outputTopic=UNPROCESSED_TOPIC
      

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

    • JOB_NAME: שם ייחודי של המשימה לפי בחירתכם
    • REGION_NAME: האזור שבו רוצים לפרוס את עבודת Dataflow, לדוגמה: us-central1
    • VERSION: הגרסה של התבנית שבה רוצים להשתמש

      אפשר להשתמש בערכים הבאים:

      • latest כדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית האב ללא תאריך בדלי – gs://dataflow-templates-REGION_NAME/latest/
      • שם הגרסה, כמו 2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית, שאפשר למצוא אותה בתיקיית האב המתאימה עם התאריך בדלי – gs://dataflow-templates-REGION_NAME/
    • SCHEMA_PATH: הנתיב ב-Cloud Storage לקובץ סכמת ה-Proto (לדוגמה, gs://MyBucket/file.pb)
    • PROTO_MESSAGE_NAME: שם הודעת ה-Proto (לדוגמה, package.name.MessageName)
    • SUBSCRIPTION_NAME: שם המינוי לקלט Pub/Sub
    • BIGQUERY_TABLE: שם טבלת הפלט ב-BigQuery
    • UNPROCESSED_TOPIC: נושא ה-Pub/Sub שבו רוצים להשתמש לתור של פריטים שלא עברו עיבוד

    API

    כדי להריץ את התבנית באמצעות API בארכיטקטורת REST, שולחים בקשת HTTP POST. מידע נוסף על ה-API ועל היקפי ההרשאות שלו זמין במאמר projects.templates.launch.

    POST https://dataflow.googleapis.com/v1b3/projects/PROJECT_ID/locations/LOCATION/flexTemplates:launch
    {
       "launch_parameter": {
          "jobName": "JOB_NAME",
          "containerSpecGcsPath": "gs://dataflow-templates-REGION_NAME/VERSION/flex/",
          "parameters": {
              "schemaPath": "SCHEMA_PATH",
              "fullMessageName": "PROTO_MESSAGE_NAME",
              "inputSubscription": "SUBSCRIPTION_NAME",
              "outputTableSpec": "BIGQUERY_TABLE",
              "outputTopic": "UNPROCESSED_TOPIC"
          }
       }
    }
      

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

    • PROJECT_ID: מזהה הפרויקט שבו רוצים להריץ את משימת Dataflow Google Cloud
    • JOB_NAME: שם ייחודי של המשימה לפי בחירתכם
    • LOCATION: האזור שבו רוצים לפרוס את עבודת Dataflow, לדוגמה: us-central1
    • VERSION: הגרסה של התבנית שבה רוצים להשתמש

      אפשר להשתמש בערכים הבאים:

      • latest כדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית האב ללא תאריך בדלי – gs://dataflow-templates-REGION_NAME/latest/
      • שם הגרסה, כמו 2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית, שאפשר למצוא אותה בתיקיית האב המתאימה עם התאריך בדלי – gs://dataflow-templates-REGION_NAME/
    • SCHEMA_PATH: הנתיב ב-Cloud Storage לקובץ סכמת ה-Proto (לדוגמה, gs://MyBucket/file.pb)
    • PROTO_MESSAGE_NAME: שם הודעת ה-Proto (לדוגמה, package.name.MessageName)
    • SUBSCRIPTION_NAME: שם המינוי לקלט Pub/Sub
    • BIGQUERY_TABLE: שם טבלת הפלט ב-BigQuery
    • UNPROCESSED_TOPIC: נושא ה-Pub/Sub שבו רוצים להשתמש לתור של פריטים שלא עברו עיבוד

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