תבנית של פונקציות UDF ב-Python מ-Pub/Sub ל-MongoDB

התבנית Pub/Sub to MongoDB with Python UDFs (‏Pub/Sub ל-MongoDB עם פונקציות מוגדרות על ידי המשתמש ב-Python) היא צינור עיבוד נתונים של סטרימינג שקורא הודעות מקודדות ב-JSON ממינוי Pub/Sub וכותב אותן ל-MongoDB כמסמכים. במקרה הצורך, צינור הנתונים הזה תומך בטרנספורמציות נוספות שאפשר לכלול באמצעות פונקציית Python שהוגדרה על ידי המשתמש (UDF).

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

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

  • מינוי Pub/Sub חייב להתקיים, וההודעות צריכות להיות מקודדות בפורמט JSON תקין.
  • קלאסטר MongoDB צריך להתקיים ולהיות נגיש ממכונות העובדים של Dataflow.

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

פרמטר תיאור
inputSubscription השם של המינוי ל-Pub/Sub. לדוגמה: projects/my-project-id/subscriptions/my-subscription-id
mongoDBUri רשימת שרתי MongoDB שמופרדים בפסיקים. לדוגמה: 192.285.234.12:27017,192.287.123.11:27017
database מסד נתונים ב-MongoDB לאחסון האוסף. לדוגמה: my-db.
collection השם של האוסף במסד הנתונים של MongoDB. לדוגמה: my-collection.
deadletterTable טבלה ב-BigQuery שבה מאוחסנות הודעות בגלל כשלים (סכימה לא תואמת, JSON לא תקין וכו'). לדוגמה: project-id:dataset-name.table-name.
pythonExternalTextTransformGcsPath אופציונלי: ה-URI של Cloud Storage של קובץ קוד Python שמגדיר את הפונקציה בהגדרת המשתמש (UDF) שרוצים להשתמש בה. לדוגמה, gs://my-bucket/my-udfs/my_file.py.
pythonExternalTextTransformFunctionName אופציונלי: השם של פונקציה בהגדרת המשתמש (UDF) ב-Python שרוצים להשתמש בה.
batchSize אופציונלי: גודל האצווה שמשמש להוספת מסמכים ל-MongoDB באצווה. ברירת מחדל: 1000.
batchSizeBytes אופציונלי: גודל אצווה בבייטים. ברירת מחדל: 5242880.
maxConnectionIdleTime אופציונלי: משך הזמן המקסימלי של חוסר פעילות בשניות לפני שחלף הזמן הקצוב לתפוגת החיבור. ברירת מחדל: 60000.
sslEnabled אופציונלי: ערך בוליאני שמציין אם חיבור ל-MongoDB מופעל באמצעות SSL. ברירת מחדל: true.
ignoreSSLCertificate אופציונלי: ערך בוליאני שמציין אם להתעלם מאישור ה-SSL. ברירת מחדל: true.
withOrdered אופציונלי: ערך בוליאני שמאפשר הוספות בכמות גדולה ל-MongoDB לפי סדר. ברירת מחדל: true.
withSSLInvalidHostNameAllowed אופציונלי: ערך בוליאני שמציין אם מותר להשתמש בשם מארח לא תקין לחיבור SSL. ברירת מחדל: true.

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

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

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

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

  • קלט: שורה אחת מקובץ קלט CSV.
  • פלט: מסמך JSON שהומר למחרוזת להוספה ל-MongoDB.

הרצת התבנית

המסוף

  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 to MongoDB with Python UDFs (מ-Pub/Sub ל-MongoDB עם פונקציות מוגדרות על ידי המשתמש ב-Python).
  6. בשדות הפרמטרים שמופיעים, מזינים את ערכי הפרמטרים.
  7. לוחצים על הפעלת העבודה.

gcloud

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

gcloud dataflow flex-template run JOB_NAME \
    --project=PROJECT_ID \
    --region=REGION_NAME \
    --template-file-gcs-location=gs://dataflow-templates-REGION_NAME/VERSION/flex/Cloud_PubSub_to_MongoDB_Xlang \
    --parameters \
inputSubscription=INPUT_SUBSCRIPTION,\
mongoDBUri=MONGODB_URI,\
database=DATABASE,
collection=COLLECTION,
deadletterTable=UNPROCESSED_TABLE
  

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

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

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

    • latest כדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית ההורה ללא תאריך בדלי – gs://dataflow-templates-REGION_NAME/latest/
    • שם הגרסה, כמו 2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית, שאפשר למצוא אותה בתיקיית האב המתאימה עם התאריך בדלי – gs://dataflow-templates-REGION_NAME/
  • INPUT_SUBSCRIPTION: המינוי ל-Pub/Sub (לדוגמה, projects/my-project-id/subscriptions/my-subscription-id)
  • MONGODB_URI: כתובות שרת MongoDB (לדוגמה, 192.285.234.12:27017,192.287.123.11:27017)
  • DATABASE: השם של מסד הנתונים של MongoDB (לדוגמה, users)
  • COLLECTION: השם של אוסף MongoDB (לדוגמה, profiles)
  • UNPROCESSED_TABLE: השם של הטבלה ב-BigQuery (לדוגמה, your-project:your-dataset.your-table-name)

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",
      "parameters": {
          "inputSubscription": "INPUT_SUBSCRIPTION",
          "mongoDBUri": "MONGODB_URI",
          "database": "DATABASE",
          "collection": "COLLECTION",
          "deadletterTable": "UNPROCESSED_TABLE"
      },
      "containerSpecGcsPath": "gs://dataflow-templates-LOCATION/VERSION/flex/Cloud_PubSub_to_MongoDB_Xlang",
   }
}
  

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

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

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

    • latest כדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית ההורה ללא תאריך בדלי – gs://dataflow-templates-REGION_NAME/latest/
    • שם הגרסה, כמו 2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית, שאפשר למצוא אותה בתיקיית האב המתאימה עם התאריך בדלי – gs://dataflow-templates-REGION_NAME/
  • INPUT_SUBSCRIPTION: המינוי ל-Pub/Sub (לדוגמה, projects/my-project-id/subscriptions/my-subscription-id)
  • MONGODB_URI: כתובות שרת MongoDB (לדוגמה, 192.285.234.12:27017,192.287.123.11:27017)
  • DATABASE: השם של מסד הנתונים של MongoDB (לדוגמה, users)
  • COLLECTION: השם של אוסף MongoDB (לדוגמה, profiles)
  • UNPROCESSED_TABLE: השם של הטבלה ב-BigQuery (לדוגמה, your-project:your-dataset.your-table-name)

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