תבנית Pub/Sub ל-MongoDB

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

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

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

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

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

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

  • inputSubscription: שם המינוי ל-Pub/Sub. לדוגמה, projects/your-project-id/subscriptions/your-subscription-name.
  • mongoDBUri: רשימה של שרתי MongoDB מופרדת בפסיקים. לדוגמה, host1:port,host2:port,host3:port.
  • database: מסד הנתונים ב-MongoDB שבו יאוחסן האוסף. לדוגמה, my-db.
  • collection: שם האוסף במסד הנתונים של MongoDB. לדוגמה, my-collection.
  • deadletterTable: טבלה ב-BigQuery שבה מאוחסנות הודעות שנגרמו כתוצאה מכשלים, כמו סכימה לא תואמת, JSON לא תקין וכו'. לדוגמה, your-project-id:your-dataset.your-table-name.

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

  • batchSize: גודל האצווה שמשמש להוספה של מסמכים ל-MongoDB. ברירת המחדל היא 1,000.
  • batchSizeBytes: גודל האצווה בבייטים. ברירת המחדל היא: 5242880.
  • maxConnectionIdleTime: משך הזמן המרבי של חוסר פעילות שמותר בשניות לפני שמתרחש פסק זמן לחיבור. ברירת המחדל היא 60,000.
  • sslEnabled: ערך בוליאני שמציין אם חיבור ה-SSL ל-MongoDB מופעל. ברירת המחדל היא true.
  • ignoreSSLCertificate: ערך בוליאני שמציין אם להתעלם מאישור ה-SSL. ברירת המחדל היא true.
  • withOrdered: ערך בוליאני שמאפשר הוספות בכמות גדולה ל-MongoDB לפי סדר. ברירת המחדל היא true.
  • withSSLInvalidHostNameAllowed: ערך בוליאני שמציין אם מותר שם מארח לא חוקי לחיבור SSL. ברירת המחדל היא true.
  • javascriptTextTransformGcsPath: ה-URI של Cloud Storage של קובץ ה-‎ .js שמגדיר את הפונקציה בהגדרת המשתמש (UDF) ב-JavaScript שבה רוצים להשתמש. לדוגמה, gs://my-bucket/my-udfs/my_file.js.
  • javascriptTextTransformFunctionName: השם של פונקציית JavaScript בהגדרת המשתמש (UDF) שבה רוצים להשתמש. לדוגמה, אם קוד פונקציית ה-JavaScript הוא myTransform(inJson) { /*...do stuff...*/ }, אז שם הפונקציה הוא myTransform. דוגמאות ל-UDF ב-JavaScript זמינות במאמר UDF Examples (https://github.com/GoogleCloudPlatform/DataflowTemplates#udf-examples).
  • javascriptTextTransformReloadIntervalMinutes: מציין את התדירות שבה יש לטעון מחדש את הפונקציה המוגדרת על ידי המשתמש, בדקות. אם הערך גדול מ-0, מערכת Dataflow בודקת מעת לעת את קובץ ה-UDF ב-Cloud Storage, ומטעינה מחדש את ה-UDF אם הקובץ משתנה. הפרמטר הזה מאפשר לכם לעדכן את ה-UDF בזמן שהצינור פועל, בלי שתצטרכו להפעיל מחדש את העבודה. אם הערך הוא 0, טעינה מחדש של פונקציות UDF מושבתת. ערך ברירת המחדל הוא 0.

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

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

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

  • 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)

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