Bigtable change streams to BigQuery template

התבנית Bigtable change streams to BigQuery היא פייפליין של סנכרון שינויים בזרמי נתונים, שמעביר רשומות של שינויים בנתוני Bigtable וכותב אותן בטבלאות BigQuery באמצעות Dataflow.

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

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

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

שורות בטבלת יומן השינויים שלא ניתן לכתוב ל-BigQuery בגלל שגיאה מתמשכת, מועברות באופן קבוע לספרייה של תור הודעות שלא עברו עיבוד (dead-letter queue) ב-Cloud Storage, לבדיקה על ידי משתמש או לעיבוד נוסף על ידי המשתמש.

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

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

סכימת טבלת הפלט ב-BigQuery

שם עמודה סוג Nullable תיאור
row_key STRING או BYTES לא מפתח השורה של השורה שהשתנתה. אם אפשרות הצינור writeRowkeyAsBytes מוגדרת כ-true, סוג העמודה חייב להיות BYTES. אחרת, משתמשים בסוג STRING.
mod_type STRING לא סוג השינוי בשורה. משתמשים באחד מהערכים הבאים: SET_CELL, DELETE_CELLS או DELETE_FAMILY.
column_family STRING לא קבוצת העמודות שהושפעה משינוי השורה.
column STRING כן מגדיר העמודה שהושפע משינוי השורה. בסוג המוטציה DELETE_FAMILY, מגדירים את הערך NULL.
commit_timestamp TIMESTAMP לא השעה שבה Bigtable מחיל את השינוי.
big_query_commit_timestamp TIMESTAMP כן אופציונלי: מציין את השעה שבה BigQuery כותב את השורה לטבלת פלט. השדה לא יאוכלס אם שם העמודה מופיע בערך של אפשרות הפייפליין bigQueryChangelogTableFieldsToIgnore.
timestamp TIMESTAMP או INT64 כן ערך חותמת הזמן של התא שהמוטציה משפיעה עליו. אם אפשרות הצינור writeNumericTimestamps מוגדרת ל-true, סוג העמודה צריך להיות INT64. אחרת, צריך להשתמש בסוג TIMESTAMP. בסוגי המוטציות DELETE_CELLS ו-DELETE_FAMILY, צריך להגדיר את הערך NULL.
timestamp_from TIMESTAMP או INT64 כן תיאור של התחלה כוללת של מרווח חותמות הזמן לכל התאים שנמחקו על ידי המוטציה DELETE_CELLS. עבור סוגי מוטציות אחרים, הערך שמוגדר הוא NULL.
timestamp_to TIMESTAMP או INT64 כן מתאר את סוף טווח הזמן הבלעדי של חותמת הזמן לכל התאים שנמחקו על ידי המוטציה DELETE_CELLS. עבור סוגי מוטציות אחרים, הערך הוא NULL.
is_gc BOOL לא אופציונלי: כאשר המוטציה מופעלת על ידי מדיניות איסוף אשפה, הגדר ל-true. בכל שאר המקרים, מגדירים את הערך false. השדה לא מאוכלס כשהשם של העמודה מופיע בערך של אפשרות הפייפליין bigQueryChangelogTableFieldsToIgnore.
source_instance STRING לא אופציונלי: תיאור של השם של מופע Bigtable שממנו מגיעה המוטציה. השדה לא מאוכלס כשהשם של העמודה מופיע בערך של אפשרות הפייפליין bigQueryChangelogTableFieldsToIgnore.
source_cluster STRING לא אופציונלי: תיאור של שם אשכול Bigtable שממנו מגיעה המוטציה. השדה לא מאוכלס כששם העמודה מופיע בערך של אפשרות הצינור bigQueryChangelogTableFieldsToIgnore.
source_table STRING לא אופציונלי: תיאור של שם טבלת Bigtable שהמוטציה חלה עליה. הערך בעמודה הזו יכול להיות שימושי אם כמה טבלאות Bigtable מעבירות שינויים בסטרימינג לאותה טבלה ב-BigQuery. השדה לא מאוכלס כשהשם של העמודה מופיע בערך של אפשרות הפייפליין bigQueryChangelogTableFieldsToIgnore.
tiebreaker INT64 לא אופציונלי: כששני שינויים נרשמים בו-זמנית על ידי אשכולות שונים של Bigtable, השינוי עם הערך tiebreaker הגבוה ביותר מוחל על טבלת המקור. מוטציות עם ערכים נמוכים יותר של tiebreaker נפסלות. השדה לא מאוכלס כשהשם של העמודה מופיע בערך של אפשרות הפייפליין bigQueryChangelogTableFieldsToIgnore.
value STRING או BYTES כן הערך החדש שמוגדר על ידי המוטציה. אם אפשרות הצינור writeValuesAsBytes מוגדרת לערך true, סוג העמודה צריך להיות BYTES. אחרת, צריך להשתמש בסוג STRING. הערך מוגדר למוטציות מסוג SET_CELL. עבור סוגי מוטציות אחרים, הערך מוגדר ל-NULL.

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

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

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

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

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

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

  • writeRowkeyAsBytes: האם לכתוב את מפתחות השורות כ-BigQuery BYTES. אם הערך הוא true, מפתחות השורות נכתבים בעמודה BYTES. אחרת, מפתחות השורות נכתבים בעמודה STRING. ברירת המחדל היא false.
  • writeValuesAsBytes: אם הערך הוא true, הערכים נכתבים בעמודה מסוג BYTES. אחרת, הם נכתבים בעמודה מסוג STRING. ערך ברירת המחדל הוא false.
  • writeNumericTimestamps: האם לכתוב את חותמת הזמן של Bigtable כ-INT64 ב-BigQuery. כשמגדירים את הערך true, הערכים נכתבים בעמודה INT64. אחרת, הערכים נכתבים בעמודה TIMESTAMP. העמודות המושפעות: timestamp,‏ timestamp_from וtimestamp_to. ברירת המחדל היא false. אם הערך הוא true, הזמן נמדד במיקרו-שניות מאז ראשית זמן יוניקס (1 בינואר 1970 ב-UTC).
  • bigQueryProjectId: מזהה הפרויקט של מערך הנתונים ב-BigQuery. ברירת המחדל היא הפרויקט של משימת Dataflow.
  • bigQueryChangelogTableName: שם טבלת היעד ב-BigQuery. אם לא מציינים ערך, המערכת משתמשת בערך bigtableReadTableId + "_changelog". ברירת המחדל היא ריק.
  • bigQueryChangelogTablePartitionGranularity: מציינת את רמת הפירוט לחלוקת טבלת יומן השינויים למחיצות. אם מגדירים את האפשרות הזו, הטבלה מחולקת למחיצות. משתמשים באחד מהערכים הנתמכים הבאים: HOUR,‏ DAY,‏ MONTH או YEAR. כברירת מחדל, הטבלה לא מחולקת למחיצות.
  • bigQueryChangelogTablePartitionExpirationMs: הגדרת זמן התפוגה של מחיצות בטבלת יומן השינויים, באלפיות השנייה. אם מגדירים את הערך true, מחיצות שהגיל שלהן גדול ממספר אלפיות השנייה שצוין נמחקות. כברירת מחדל, לא מוגדר תאריך תפוגה.
  • bigQueryChangelogTableFieldsToIgnore: רשימה מופרדת בפסיקים של עמודות ביומן השינויים, שאם מציינים אותן, הן לא נוצרות ולא מאוכלסות. אפשר להשתמש באחד מהערכים הנתמכים הבאים: is_gc,‏ source_instance,‏ source_cluster,‏ source_table,‏ tiebreaker או big_query_commit_timestamp. כברירת מחדל, כל העמודות מאוכלסות.
  • dlqDirectory: הספרייה שבה ישמש תור של הודעות שלא ניתן להעביר (DLQ). רשומות שלא ניתן לעבד נשמרות בספרייה הזו. ברירת המחדל היא ספרייה במיקום הזמני של משימת Dataflow. ברוב המקרים, אפשר להשתמש בנתיב ברירת המחדל.
  • bigtableChangeStreamMetadataInstanceId: מזהה מופע של מטא-נתונים של Bigtable change streams. ברירת המחדל היא ריק.
  • bigtableChangeStreamMetadataTableTableId: המזהה של טבלת המטא-נתונים של מחבר Bigtable change streams. אם לא מספקים את המידע הזה, נוצרת באופן אוטומטי טבלת מטא-נתונים של מחבר Bigtable לשינויים בזמן אמת במהלך ההפעלה של צינור הנתונים. ברירת המחדל היא ריק.
  • bigtableChangeStreamCharset: שם ערכת התווים של סנכרון שינויים בזרמי נתונים של Bigtable. ברירת המחדל היא: 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, בוחרים בתבנית Bigtable change streams to 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 \
bigtableReadInstanceId=BIGTABLE_INSTANCE_ID,\
bigtableReadTableId=BIGTABLE_TABLE_ID,\
bigtableChangeStreamAppProfile=BIGTABLE_APPLICATION_PROFILE_ID,\
bigQueryDataset=BIGQUERY_DESTINATION_DATASET

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

  • 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.
  • BIGQUERY_DESTINATION_DATASET: השם של מערך הנתונים של היעד ב-BigQuery

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",
        "bigQueryDataset": "BIGQUERY_DESTINATION_DATASET"
    }
  }
}

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

  • 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.
  • BIGQUERY_DESTINATION_DATASET: השם של מערך הנתונים של היעד ב-BigQuery

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