תבנית של סנכרון שינויים בזרמי נתונים מ-Bigtable ל-Pub/Sub

תבנית Bigtable change streams to Pub/Sub היא צינור עיבוד נתונים בסטרימינג שמעביר בסטרימינג רשומות של שינויים בנתוני Bigtable ומפרסם אותן בנושא Pub/Sub באמצעות Dataflow.

זרם שינויים ב-Bigtable מאפשר להירשם לשינויים בנתונים ברמת הטבלה. כשנרשמים לסנכרון שינויים בזרמי נתונים בטבלה, חלים האילוצים הבאים:

  • מוחזרים רק תאים ששונו ותיאורים של פעולות מחיקה.
  • מוחזר רק הערך החדש של תא ששונה.

כשמפרסמים רשומות של שינויי נתונים בנושא Pub/Sub, יכול להיות שההודעות יוכנסו שלא לפי הסדר בהשוואה לסדר המקורי של חותמות הזמן של ביצוע השינויים ב-Bigtable.

רשומות של שינויים בנתונים ב-Bigtable שלא ניתן לפרסם בנושאים ב-Pub/Sub מועברות באופן זמני לספרייה של תור הודעות שלא ניתן לעיבוד (dead-letter queue) ב-Cloud Storage. אחרי מספר הניסיונות החוזרים המקסימלי שלא צלחו, הרשומות האלה מוצבות ללא הגבלת זמן באותו מדריך של תור הודעות מתות לבדיקה אנושית או לעיבוד נוסף על ידי המשתמש.

הצינור דורש שנושא היעד ב-Pub/Sub יתקיים. יכול להיות שהנושא של היעד מוגדר לאימות הודעות באמצעות סכימה. אם נושא Pub/Sub מציין סכימה, הצינור יתחיל לפעול רק אם הסכימה תקינה. בהתאם לסוג הסכימה, משתמשים באחת מהגדרות הסכימה הבאות לנושא היעד:

מאגרי אחסון לפרוטוקולים

syntax = "proto2";

package com.google.cloud.teleport.bigtable;

option java_outer_classname = "ChangeLogEntryProto";

message ChangelogEntryProto{
  required bytes rowKey = 1;
  enum ModType {
    SET_CELL = 0;
    DELETE_FAMILY = 1;
    DELETE_CELLS = 2;
    UNKNOWN = 3;
  }
  required ModType modType = 2;
  required bool isGC = 3;
  required int32 tieBreaker = 4;
  required int64 commitTimestamp = 5;
  required string columnFamily = 6;
  optional bytes column = 7;
  optional int64 timestamp = 8;
  optional int64 timestampFrom = 9;
  optional int64 timestampTo = 10;
  optional bytes value = 11;
  required string sourceInstance = 12;
  required string sourceCluster = 13;
  required string sourceTable = 14;
}
  

Avro

{
    "name" : "ChangelogEntryMessage",
    "type" : "record",
    "namespace" : "com.google.cloud.teleport.bigtable",
    "fields" : [
      { "name" : "rowKey", "type" : "bytes"},
      {
        "name" : "modType",
        "type" : {
          "name": "ModType",
          "type": "enum",
          "symbols": ["SET_CELL", "DELETE_FAMILY", "DELETE_CELLS", "UNKNOWN"]}
      },
      { "name": "isGC", "type": "boolean" },
      { "name": "tieBreaker", "type": "int"},
      { "name": "columnFamily", "type": "string"},
      { "name": "commitTimestamp", "type" : "long"},
      { "name" : "sourceInstance", "type" : "string"},
      { "name" : "sourceCluster", "type" : "string"},
      { "name" : "sourceTable", "type" : "string"},
      { "name": "column", "type" : ["null", "bytes"]},
      { "name": "timestamp", "type" : ["null", "long"]},
      { "name": "timestampFrom", "type" : ["null", "long"]},
      { "name": "timestampTo", "type" : ["null", "long"]},
      { "name" : "value", "type" : ["null", "bytes"]}
   ]
}
    

JSON

משתמשים בסכימת Protobuf הבאה עם קידוד הודעות JSON:

syntax = "proto2";

package com.google.cloud.teleport.bigtable;

option java_outer_classname = "ChangelogEntryMessageText";

message ChangelogEntryText{
  required string rowKey = 1;
  enum ModType {
    SET_CELL = 0;
    DELETE_FAMILY = 1;
    DELETE_CELLS = 2;
    UNKNOWN = 3;
  }
  required ModType modType = 2;
  required bool isGC = 3;
  required int32 tieBreaker = 4;
  required int64 commitTimestamp = 5;
  required string columnFamily = 6;
  optional string column = 7;
  optional int64 timestamp = 8;
  optional int64 timestampFrom = 9;
  optional int64 timestampTo = 10;
  optional string value = 11;
  required string sourceInstance = 12;
  required string sourceCluster = 13;
  required string sourceTable = 14;
}
    

כל הודעה חדשה ב-Pub/Sub כוללת רשומה אחת מרשומה של שינוי בנתונים שמוחזרת על ידי זרם השינויים מהשורה המתאימה בטבלת Bigtable. התבנית Pub/Sub משטחת את הרשומות בכל רשומה של שינוי נתונים לשינויים ברמת התא הבודד.

תיאור הודעת הפלט של Pub/Sub

שם השדה תיאור
rowKey מפתח השורה של השורה שהשתנתה. הנתונים מגיעים בצורה של מערך בייטים. כשמגדירים קידוד של הודעות JSON, מפתחות השורות מוחזרים כמחרוזות. אם מציינים את useBase64Rowkeys, מפתחות השורות מקודדים ב-Base64. אחרת, נעשה שימוש בערכת תווים שצוינה על ידי bigtableChangeStreamCharset כדי לפענח את הבייטים של מפתח השורה למחרוזת.
modType סוג השינוי בשורה. משתמשים באחד מהערכים הבאים: SET_CELL, DELETE_CELLS או DELETE_FAMILY.
columnFamily קבוצת העמודות שהושפעה משינוי השורה.
column מגדיר העמודה שהושפע משינוי השורה. עבור סוג המוטציה DELETE_FAMILY, שדה העמודה לא מוגדר. הנתונים מגיעים בצורה של מערך בייטים. כשמוגדר קידוד הודעות JSON, העמודות מוחזרות כמחרוזות. אם מציינים את useBase64ColumnQualifier, שדה העמודה מקודד ב-Base64. אחרת, נעשה שימוש בערכת תווים שצוינה על ידי bigtableChangeStreamCharset כדי לפענח את הבייטים של מפתח השורה למחרוזת.
commitTimestamp השעה שבה Bigtable מחיל את השינוי. הזמן נמדד במיקרו-שניות מאז ראשית זמן יוניקס (1 בינואר 1970 בשעה 00:00 לפי שעון UTC).
timestamp ערך חותמת הזמן של התא שהושפע מהשינוי. עבור סוגי המוטציות DELETE_CELLS ו-DELETE_FAMILY, חותמת הזמן לא מוגדרת. הזמן נמדד במיקרו-שניות מאז ראשית זמן יוניקס (1 בינואר 1970 בשעה UTC).
timestampFrom מתאר את ההתחלה של מרווח הזמן של חותמות הזמן, כולל, עבור כל התאים שנמחקו על ידי המוטציה DELETE_CELLS. בסוגים אחרים של מוטציות, הערך של timestampFrom לא מוגדר. הזמן נמדד במיקרו-שניות מאז ראשית זמן יוניקס (1 בינואר 1970 בשעה 00:00 לפי שעון UTC).
timestampTo מתאר את הסוף הבלעדי של מרווח חותמות הזמן של כל התאים שנמחקו על ידי המוטציה DELETE_CELLS. בסוגים אחרים של מוטציות, הערך של timestampTo לא מוגדר.
isGC ערך בוליאני שמציין אם הפעולה לשינוי נתונים נוצרה על ידי מנגנון איסוף של Bigtable.
tieBreaker כששני שינויים נרשמים בו-זמנית על ידי אשכולות שונים של Bigtable, השינוי עם הערך tiebreaker הגבוה ביותר מוחל על טבלת המקור. מוטציות עם ערכי tiebreaker נמוכים יותר נפסלות.
value הערך החדש שהוגדר על ידי השינוי. אם לא מגדירים את האפשרות stripValues של צינור העיבוד, הערך מוגדר למוטציות SET_CELL. בסוגים אחרים של מוטציות, הערך לא מוגדר. הנתונים מגיעים בצורה של מערך בייטים. כשמוגדר קידוד של הודעות JSON, הערכים מוחזרים כמחרוזות. אם מציינים את useBase64Values, הערך מקודד ב-Base64. אחרת, נעשה שימוש בערכת תווים שצוינה על ידי bigtableChangeStreamCharset כדי לפענח את בייטים של הערך למחרוזת.
sourceInstance השם של מופע Bigtable שרשם את השינוי. יכול להיות שזה יקרה כשמספר צינורות מעבירים שינויים ממופעים שונים לאותו נושא Pub/Sub.
sourceCluster השם של אשכול Bigtable שרשם את השינוי. יכול להיות שיהיה צורך להשתמש באפשרות הזו כשכמה צינורות מעבירים שינויים בסטרימינג ממופעים שונים לאותו נושא Pub/Sub.
sourceTable השם של טבלת Bigtable שקיבלה את השינוי. יכול להיות שיהיה שימוש בפרמטר הזה במקרים שבהם כמה צינורות מעבירים שינויים מטבלאות שונות לאותו נושא Pub/Sub.

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

  • מופע המקור של Bigtable שצוין.
  • טבלת המקור שצוינה ב-Bigtable. צריך להפעיל את סנכרון שינויים בזרמי נתונים בטבלה.
  • פרופיל האפליקציה שצוין ב-Bigtable.
  • צריך לציין נושא Pub/Sub קיים.

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

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

  • pubSubTopic: השם של נושא היעד ב-Pub/Sub.
  • bigtableChangeStreamAppProfile: מזהה פרופיל האפליקציה של Bigtable. פרופיל האפליקציה צריך להשתמש בניתוב של אשכול יחיד ולאפשר עסקאות של שורה אחת.
  • bigtableReadInstanceId: המזהה של מכונת Bigtable של המקור.
  • bigtableReadTableId: מזהה טבלת Bigtable של המקור.

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

  • messageEncoding: הקידוד של ההודעות שיפורסמו בנושא Pub/Sub. כשסכימת נושא היעד מוגדרת, קידוד ההודעה נקבע לפי הגדרות הנושא. הערכים הנתמכים הם: BINARY ו-JSON. ברירת המחדל היא JSON.
  • messageFormat: הקידוד של ההודעות לפרסום בנושא Pub/Sub. כשסכימת נושא היעד מוגדרת, קידוד ההודעה נקבע לפי הגדרות הנושא. הערכים הנתמכים הם: AVRO,‏ PROTOCOL_BUFFERS ו-JSON. ערך ברירת המחדל הוא JSON. כשמשתמשים בפורמט JSON, השדות rowKey,‏ column ו-value של ההודעה הם מחרוזות, והתוכן שלהן נקבע לפי אפשרויות הצינור useBase64Rowkeys,‏ useBase64ColumnQualifiers,‏ useBase64Values ו-bigtableChangeStreamCharset.
  • stripValues: אם הערך מוגדר כ-true, מוטציות SET_CELL מוחזרות בלי ערכים חדשים. ברירת המחדל היא false. הפרמטר הזה שימושי כשלא צריך שיופיע ערך חדש, או כשערכים גדולים מאוד וחורגים ממגבלות הגודל של הודעות Pub/Sub.
  • dlqDirectory: הספרייה של תור ההודעות המתות. רשומות שלא עוברות עיבוד מאוחסנות בספרייה הזו. ברירת המחדל היא ספרייה במיקום הזמני של משימת Dataflow. ברוב המקרים אפשר להשתמש בנתיב ברירת המחדל.
  • dlqRetryMinutes: מספר הדקות בין ניסיונות חוזרים של תור הודעות מתות. ברירת המחדל היא 10.
  • dlqMaxRetries: המספר המקסימלי של ניסיונות חוזרים של הודעות שהועברו לתור של הודעות שלא ניתן למסור. ברירת המחדל היא 5.
  • useBase64Rowkeys: משמש עם קידוד הודעות JSON. אם הערך הוא true, השדה rowKey הוא מחרוזת בקידוד Base64. אחרת, הפונקציה rowKey נוצרת באמצעות bigtableChangeStreamCharset כדי לפענח בייטים למחרוזת. ברירת המחדל היא false.
  • pubSubProjectId: מזהה הפרויקט ב-Bigtable. ברירת המחדל היא הפרויקט של משימת Dataflow.
  • useBase64ColumnQualifiers: משמש עם קידוד הודעות JSON. אם הערך הוא true, השדה column הוא מחרוזת בקידוד Base64. אחרת, העמודה נוצרת באמצעות bigtableChangeStreamCharset כדי לפענח בייטים למחרוזת. ברירת המחדל היא false.
  • useBase64Values: משמש עם קידוד הודעות JSON. אם הערך הוא true, שדה הערך הוא מחרוזת בקידוד Base64. אחרת, הערך נוצר באמצעות bigtableChangeStreamCharset כדי לפענח בייטים למחרוזת. ברירת המחדל היא false.
  • disableDlqRetries: מציין אם להשבית את הניסיונות החוזרים עבור תור ההודעות המתות. ברירת המחדל היא false.
  • bigtableChangeStreamMetadataInstanceId: מזהה מופע של מטא-נתונים של Bigtable change streams. ברירת המחדל היא ריקה.
  • bigtableChangeStreamMetadataTableTableId: המזהה של טבלת המטא-נתונים של מחבר Bigtable change streams. אם לא מספקים את המידע הזה, נוצרת באופן אוטומטי טבלת מטא-נתונים של מחבר Bigtable לשינויים בזמן אמת במהלך ההפעלה של צינור הנתונים. ברירת המחדל היא ריקה.
  • bigtableChangeStreamCharset: שם ערכת התווים של Bigtable change streams. ברירת המחדל היא UTF-8.
  • bigtableChangeStreamStartTimestamp: חותמת הזמן (https://tools.ietf.org/html/rfc3339) ההתחלתית (כולל) לשימוש לקריאת סנכרון שינויים בזרמי נתונים. לדוגמה, 2022-05-05T07:59:59Z. ברירת המחדל היא חותמת הזמן של שעת ההתחלה של הצינור.
  • bigtableChangeStreamIgnoreColumnFamilies: רשימה מופרדת בפסיקים של שינויים בשמות של משפחות עמודות שצריך להתעלם מהם. ברירת המחדל היא ריקה.
  • bigtableChangeStreamIgnoreColumns: רשימה מופרדת בפסיקים של שינויים בשמות עמודות שצריך להתעלם מהם. דוגמה: "cf1:col1,cf2:col2". ברירת המחדל היא ריקה.
  • bigtableChangeStreamName: שם ייחודי של צינור הנתונים של הלקוח. מאפשרת לכם להמשיך את העיבוד מהנקודה שבה צינור קודם שהופעל נעצר. ברירת המחדל היא שם שנוצר באופן אוטומטי. הערך שבו נעשה שימוש מופיע ביומני המשימות של Dataflow.
  • bigtableChangeStreamResume: אם הערך הוא true, צינור חדש ימשיך את העיבוד מהנקודה שבה צינור שפעל בעבר עם אותו ערך bigtableChangeStreamName הפסיק. אם הצינור עם הערך bigtableChangeStreamName שצוין לא הופעל אף פעם, לא יופעל צינור חדש. כשמגדירים את האפשרות false, מתחיל צינור חדש. אם צינור עם אותו ערך bigtableChangeStreamName כבר הופעל עבור המקור הנתון, לא יופעל צינור חדש. ברירת המחדל היא false.
  • bigtableReadChangeStreamTimeoutMs: הזמן הקצוב לתפוגה של בקשות Bigtable ReadChangeStream באלפיות השנייה.
  • bigtableReadProjectId: מזהה הפרויקט ב-Bigtable. ברירת המחדל היא הפרויקט של משימת Dataflow.

הרצת התבנית

המסוף

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

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

  5. בתפריט הנפתח Dataflow template (תבנית של Dataflow), בוחרים בתבנית Bigtable change streams to Pub/Sub (העברת סנכרון שינויים בזרמי נתונים ב-Bigtable ל-Pub/Sub).
  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 \
bigtableReadInstanceId=BIGTABLE_INSTANCE_ID,\
bigtableReadTableId=BIGTABLE_TABLE_ID,\
bigtableChangeStreamAppProfile=BIGTABLE_APPLICATION_PROFILE_ID,\
pubSubTopic=PUBSUB_TOPIC

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

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

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

    • latest כדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית ההורה ללא תאריך בדלי – gs://dataflow-templates-REGION_NAME/latest/
    • שם הגרסה, כמו 2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית, שאפשר למצוא אותה בתיקיית האב המתאימה עם התאריך בדלי – gs://dataflow-templates-REGION_NAME/
  • REGION_NAME: האזור שבו רוצים לפרוס את עבודת Dataflow, לדוגמה: us-central1
  • BIGTABLE_INSTANCE_ID: מזהה מופע Bigtable.
  • BIGTABLE_TABLE_ID: מזהה הטבלה ב-Bigtable.
  • BIGTABLE_APPLICATION_PROFILE_ID: מזהה פרופיל האפליקציה של Bigtable.
  • PUBSUB_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": {
        "bigtableReadInstanceId": "BIGTABLE_INSTANCE_ID",
        "bigtableReadTableId": "BIGTABLE_TABLE_ID",
        "bigtableChangeStreamAppProfile": "BIGTABLE_APPLICATION_PROFILE_ID",
        "pubSubTopic": "PUBSUB_TOPIC"
    }
  }
}

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

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

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

    • latest כדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית ההורה ללא תאריך בדלי – gs://dataflow-templates-REGION_NAME/latest/
    • שם הגרסה, כמו 2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית, שאפשר למצוא אותה בתיקיית האב המתאימה עם התאריך בדלי – gs://dataflow-templates-REGION_NAME/
  • LOCATION: האזור שבו רוצים לפרוס את עבודת Dataflow, לדוגמה: us-central1
  • BIGTABLE_INSTANCE_ID: מזהה מופע Bigtable.
  • BIGTABLE_TABLE_ID: מזהה הטבלה ב-Bigtable.
  • BIGTABLE_APPLICATION_PROFILE_ID: מזהה פרופיל האפליקציה של Bigtable.
  • PUBSUB_TOPIC: שם נושא היעד ב-Pub/Sub

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