תבנית של מחולל נתונים בסטרימינג

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

התבנית תומכת ביעדים הבאים:

  • נושא ב-Apache Kafka
  • טבלה ב-BigQuery
  • קטגוריה של Cloud Storage
  • נקודת קצה של Java Database Connectivity ‏ (JDBC)
  • נושא Pub/Sub
  • טבלת Spanner

הנה כמה תרחישי שימוש אפשריים:

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

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

הגדרת סכימת הרשומה

התבנית מספקת סכמה מוגדרת מראש לנתונים שנוצרו. כדי להשתמש בסכימה הזו, צריך להגדיר את פרמטר התבנית schemaTemplate ל-GAME_EVENT.

לחלופין, אפשר לספק סכימת נתונים משלכם באופן הבא:

  1. יוצרים קובץ סכימה שמכיל תבנית JSON לנתונים שנוצרו. התבנית הזו משתמשת בספרייה JSON Data Generator, שתומכת בפונקציות שונות להוספת אקראיות לנתונים. לדוגמה:

    {
      "id": {{integer(0,1000)}},
      "name": "{{uuid()}}",
      "isInStock": {{bool()}}
    }

    מידע נוסף זמין ב מסמכי התיעוד של json-data-generator.

  2. מעלים את קובץ הסכימה לקטגוריה של Cloud Storage.
  3. מגדירים את פרמטר התבנית schemaLocation ל-URI של Cloud Storage של קובץ התבנית.

ציון פורמט הפלט

כברירת מחדל, התבנית יוצרת נתוני JSON. במקומות מסוימים, התבנית תומכת גם בפורמטים Avro או Parquet:

  • ‫Avro: נתמך ב-Cloud Storage, ב-Apache Kafka וב-Pub/Sub
  • ‫Parquet: נתמך ב-Cloud Storage.

כדי להפיק פלט בפורמט Avro או Parquet:

  1. מגדירים את פרמטר התבנית outputType בתור AVRO לפורמט Avro, או בתור PARQUET לפורמט Parquet.
  2. יוצרים קובץ סכימה של Avro.
  3. מעלים את קובץ הסכימה ל-Cloud Storage.
  4. מגדירים את פרמטר התבנית avroSchemaLocation ל-URI של Cloud Storage של קובץ הסכימה.

ציון יעד הנתונים

בקטעים הבאים מוסבר איך להגדיר את התבנית לכל סוג של יעד.

נושא ב-Apache Kafka

כדי לכתוב לנושא Kafka, מגדירים את פרמטרי התבנית הבאים:

  • sinkType: KAFKA.
  • bootstrapServer: כתובת האתחול של אשכול Kafka.
  • kafkaTopic: נושא Kafka שאליו רוצים לכתוב.

אם אתם כותבים לאשכול של שירות מנוהל של Google Cloud ל-Apache Kafka, צריך להעניק ל חשבון השירות של העובד את התפקיד 'לקוח מנוהל של Kafka' (roles/managedkafka.client).

טבלה ב-BigQuery

כדי לכתוב לטבלה ב-BigQuery, צריך להגדיר את פרמטרים התבנית הבאים:

  • sinkType: BIGQUERY.
  • outputTableSpec: הטבלה ב-BigQuery שאליה רוצים לכתוב. צריך לפרמט את הפרמטר הזה באופן הבא: PROJECT_ID:DATASET.TABLE.

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

  • outputDeadletterTable: שם הטבלה שבה צינור הנתונים כותב רשומות שנכשלו. אם לא מציינים שם, צינור העיבוד יוצר טבלה בשם OUTPUT_TABLE_error_records, כאשר OUTPUT_TABLE הוא שם טבלת הפלט.
  • writeDisposition: מציין איך לכתוב לטבלה קיימת. אלה הערכים הנתמכים:

    • WRITE_APPEND. מוסיפים שורות לטבלה הקיימת.
    • WRITE_TRUNCATE. חיתוך השורות הקיימות.
    • WRITE_EMPTY. כותבים רק אם הטבלה ריקה. אם כבר יש נתונים בטבלה, המשימה תיכשל.

    ערך ברירת המחדל הוא WRITE_APPEND.

נותנים ל חשבון השירות של העובד את התפקיד BigQuery Data Editor (עריכת נתונים ב-BigQuery) ‏(roles/bigquery.dataEditor).

Cloud Storage

כדי לכתוב לקטגוריה של Cloud Storage, מגדירים את פרמטרי התבנית הבאים:

  • sinkType: GCS.
  • outputDirectory: הנתיב של תיקיית Cloud Storage שאליה רוצים לכתוב.

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

  • numShards: המספר המקסימלי של רסיסים. ערך גבוה יותר יכול לאפשר תפוקה גבוהה יותר, אבל עלול להוביל לעלויות גבוהות יותר של צבירת נתונים. אם הערך הוא 0, Dataflow בוחר את מספר הרסיסים. ערך ברירת המחדל הוא 0.
  • outputFilenamePrefix: הקידומת של שם הקובץ. ערך ברירת המחדל הוא output-.
  • windowDuration: המרווח שבו צינור הנתונים כותב קבצים ל-Cloud Storage. הפורמטים המותרים הם: Ns (שניות), Nm (דקות) ו- Nh (שעות). ערך ברירת המחדל הוא 1m (דקה אחת).

מקצים ל חשבון השירות של העובד את התפקיד 'אדמין אובייקט אחסון' (roles/storage.objectAdmin).

נקודת קצה של JDBC

כדי לכתוב לנקודת קצה של JDBC, מגדירים את פרמטרים התבנית הבאים:

  • sinkType: JDBC.
  • driverClassName: מחלקת מנהל ההתקן של JDBC שבה ייעשה שימוש. דוגמה: com.mysql.jdbc.Driver.
  • connectionUrl: מחרוזת החיבור לחיבור למקור JDBC.
  • statement: הצהרת ה-INSERT INTOSQL שמשמשת לכתיבה במסד הנתונים. בהצהרה צריך לציין אילו עמודות בטבלה לכתוב, עם תווים של placeholder‏ '?' עבור הסעיף VALUES. הצינור מחליף את ה-placeholders בערכים המתאימים של השדות מנתוני ה-JSON.

    דוגמה: INSERT INTO tableName (column1, column2) VALUES (?,?).

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

  • username: שם המשתמש לחיבור JDBC.
  • password: הסיסמה לחיבור JDBC.
  • connectionProperties: מחרוזת מאפיינים לחיבור JDBC. דוגמה: unicode=true;characterEncoding=UTF-8.

נושא Pub/Sub

כדי לכתוב לנושא Pub/Sub, מגדירים את הפרמטרים הבאים של התבנית:

  • sinkType: PUBSUB.
  • topic: נושא ה-Pub/Sub שרוצים לכתוב אליו.

מקצים ל חשבון השירות של העובד את התפקיד 'פרסום הודעות ב-Pub/Sub' (roles/pubsub.publisher).

טבלת Spanner

כדי לכתוב לטבלת Spanner, מגדירים את הפרמטרים הבאים של התבנית:

  • sinkType: SPANNER.
  • projectId: מזהה הפרויקט שמכיל את טבלת Spanner.
  • spannerInstanceName: השם של מכונת Spanner.
  • spannerDatabaseName: השם של מסד הנתונים ב-Spanner.
  • spannerTableName: השם של טבלת Spanner.

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

  • maxNumMutations: המספר המקסימלי של תאים שעברו שינוי בכל אצווה.
  • maxNumRows: המספר המקסימלי של שורות שעברו שינוי בכל אצווה.
  • batchSizeBytes: המספר המקסימלי של בייטים שניתן לשנות בכל קבוצת פעולות.
  • commitDeadlineSeconds: תאריך היעד לקריאה ל-API של השליחה, בשניות.

מקצים ל חשבון השירות של העובד את התפקיד 'משתמש במסד נתונים ב-Cloud Spanner' (roles/spanner.databaseUser).

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

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

  • qps: מציין את קצב ההודעות לשנייה שיתפרסמו ב-Pub/Sub.

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

  • schemaTemplate: תבנית סכימה קיימת לשימוש. הערך צריך להיות אחד מהערכים הבאים: [GAME_EVENT].
  • schemaLocation: הנתיב ב-Cloud Storage של מיקום הסכימה. לדוגמה, gs://<bucket-name>/prefix.
  • topic: השם של הנושא שאליו צינור הנתונים צריך לפרסם נתונים. לדוגמה, projects/<project-id>/topics/<topic-name>.
  • messagesLimit: מציין את המספר המקסימלי של הודעות הפלט שייווצרו. הערך 0 מציין שההגדרה לא מוגבלת. ברירת המחדל היא 0.
  • outputType: סוג הפלט של ההודעה. ברירת המחדל היא JSON.
  • avroSchemaLocation: נתיב ב-Cloud Storage למיקום של סכימת Avro. חובה אם סוג הפלט הוא AVRO או PARQUET. לדוגמה, gs://your-bucket/your-path/schema.avsc.
  • sinkType: סוג ה-Sink של ההודעה. ברירת המחדל היא PUBSUB.
  • outputTableSpec: טבלת הפלט ב-BigQuery. חובה כש-sinkType הוא BIGQUERY. לדוגמה, <project>:<dataset>.<table_name>.
  • writeDisposition: BigQuery WriteDisposition. לדוגמה, WRITE_APPEND, ‏ WRITE_EMPTY או WRITE_TRUNCATE. ברירת המחדל היא: WRITE_APPEND.
  • outputDeadletterTable: הודעות שלא הצליחו להגיע לטבלת הפלט מכל מיני סיבות (למשל, סכימה לא תואמת, JSON לא תקין) נכתבות בטבלה הזו. אם הוא לא קיים, הוא ייווצר במהלך ההרצה של צינור הנתונים. לדוגמה, your-project-id:your-dataset.your-table-name.
  • windowDuration: משך החלון או גודל החלון שבו הנתונים ייכתבו ל-Cloud Storage. הפורמטים המותרים הם: Ns (לשניות, לדוגמה: 5s), ‏ Nm (לדקות, לדוגמה: 12m), ‏ Nh (לשעות, לדוגמה: 2h). לדוגמה, 1m. ברירת המחדל היא: ‎1m.
  • outputDirectory: הנתיב והקידומת של שם הקובץ לכתיבת קובצי פלט. חייב להסתיים בלוכסן. הפורמט של התאריך והשעה משמש לניתוח נתיב הספרייה. לדוגמה, gs://your-bucket/your-path/.
  • outputFilenamePrefix: הקידומת שתוצב בכל קובץ עם חלון. לדוגמה, output-. ברירת המחדל היא: output-.
  • numShards: המספר המקסימלי של רסיסי פלט שנוצרים בזמן כתיבה. מספר גבוה יותר של רסיסים פירושו תפוקה גבוהה יותר לכתיבה ב-Cloud Storage, אבל יכול להיות שעלות צבירת הנתונים תהיה גבוהה יותר ברסיסים כשמעבדים קבצים של פלט Cloud Storage. ערך ברירת המחדל נקבע על ידי Dataflow.
  • driverClassName: שם מחלקת מנהל ההתקן של JDBC לשימוש. לדוגמה, com.mysql.jdbc.Driver.
  • connectionUrl: מחרוזת חיבור של כתובת URL לחיבור למקור JDBC. לדוגמה, jdbc:mysql://some-host:3306/sampledb.
  • username: שם המשתמש שישמש לחיבור JDBC.
  • password: הסיסמה שתשמש לחיבור JDBC.
  • connectionProperties: מחרוזת מאפיינים לשימוש בחיבור JDBC. הפורמט של המחרוזת צריך להיות [propertyName=property;]*. לדוגמה, unicode=true;characterEncoding=UTF-8.
  • statement: הצהרת SQL שתופעל כדי לכתוב במסד הנתונים. בהצהרה צריך לציין את שמות העמודות בטבלה, בכל סדר. רק הערכים של שמות העמודות שצוינו ייקראו מ-JSON ויתווספו להצהרה. לדוגמה, INSERT INTO tableName (column1, column2) VALUES (?,?).
  • projectId: מזהה הפרויקט ב-GCP שבו נמצאת טבלת Spanner.
  • spannerInstanceName: שם מכונת Cloud Spanner.
  • spannerDatabaseName: שם מסד הנתונים ב-Cloud Spanner.
  • spannerTableName: שם הטבלה ב-Cloud Spanner.
  • maxNumMutations: מציינת את מגבלת המוטציות בתאים (המספר המקסימלי של תאים שעברו מוטציה בכל אצווה). ערך ברירת המחדל הוא 5,000.
  • maxNumRows: מציינת את מגבלת השינויים בשורות (מספר השורות המקסימלי שניתן לשנות בכל אצווה). ערך ברירת המחדל הוא 1,000.
  • batchSizeBytes: מציין את מגבלת גודל האצווה (מספר הבייטים המקסימלי שניתן לשנות בכל אצווה). ערך ברירת המחדל הוא 1MB.
  • commitDeadlineSeconds: מציין את תאריך היעד בשניות לקריאה ל-Commit API.
  • bootstrapServer: שרת ה-Bootstrap של Kafka. לדוגמה, localhost:9092.
  • kafkaTopic: נושא Kafka לכתיבה. לדוגמה, topic.

הרצת התבנית

המסוף

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

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

  5. בתפריט הנפתח תבנית של העברת נתונים, בוחרים בתבנית Streaming Data Generator.
  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/ \
    --parameters \
schemaLocation=SCHEMA_LOCATION,\
qps=QPS,\
topic=PUBSUB_TOPIC
  

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

  • PROJECT_ID: מזהה Google Cloud הפרויקט שבו רוצים להריץ את משימת Dataflow
  • 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/
  • SCHEMA_LOCATION: הנתיב לקובץ הסכימה ב-Cloud Storage. לדוגמה: gs://mybucket/filename.json.
  • QPS: מספר ההודעות שיפורסמו בכל שנייה
  • PUBSUB_TOPIC: נושא הפלט של Pub/Sub. לדוגמה: projects/my-project-id/topics/my-topic-id.

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": {
          "schemaLocation": "SCHEMA_LOCATION",
          "qps": "QPS",
          "topic": "PUBSUB_TOPIC"
      },
      "containerSpecGcsPath": "gs://dataflow-templates-LOCATION/VERSION/flex/",
   }
}
  

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

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

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

    • latest כדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית ההורה ללא תאריך בדלי – gs://dataflow-templates-REGION_NAME/latest/‎
    • שם הגרסה, כמו 2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית, שאפשר למצוא אותה בתיקיית האב המתאימה עם התאריך בדלי – gs://dataflow-templates-REGION_NAME/
  • SCHEMA_LOCATION: הנתיב לקובץ הסכימה ב-Cloud Storage. לדוגמה: gs://mybucket/filename.json.
  • QPS: מספר ההודעות שיפורסמו בכל שנייה
  • PUBSUB_TOPIC: נושא הפלט של Pub/Sub. לדוגמה: projects/my-project-id/topics/my-topic-id.

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