התבנית 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: הספרייה שבה יש להשתמש לתור של הודעות שלא ניתן להעביר. רשומות שלא עוברות עיבוד נשמרות בספרייה הזו. ברירת המחדל היא ספרייה במיקום הזמני של עבודת Dataflow. ברוב המקרים אפשר להשתמש בנתיב ברירת המחדל.
- 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.
הרצת התבנית
המסוף
- עוברים לדף Create job from template (יצירת משימה מתבנית) ב-Dataflow. כניסה לדף Create job from template
- בשדה שם המשימה, מזינים שם ייחודי למשימה.
- אופציונלי: בשדה Regional endpoint (נקודת קצה אזורית), בוחרים ערך מהתפריט הנפתח. אזור ברירת המחדל הוא
us-central1.רשימה של אזורים שבהם אפשר להריץ משימת Dataflow מופיעה במאמר מיקומי Dataflow.
- בתפריט הנפתח Dataflow template, בוחרים בתבנית Bigtable change streams to BigQuery.
- בשדות הפרמטרים שמופיעים, מזינים את ערכי הפרמטרים.
- לוחצים על הפעלת העבודה.
gcloud
במעטפת או בטרמינל, מריצים את התבנית:
gcloud dataflow flex-template run JOB_NAME \ --region=REGION_NAME \ --template-file-gcs-location=gs://dataflow-templates-REGION_NAME/VERSION/flex/Bigtable_Change_Streams_to_BigQuery \ --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/Bigtable_Change_Streams_to_BigQuery", "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