פתרון בעיות בהעברה

בקטע הזה מפורטים טיפים לפתרון בעיות נפוצות שיכולות לקרות כשיוצרים ומנהלים צינורות של Datastream ו-Dataflow.

פתרון בעיות בפרופיל חיבור של Datastream

במהלך תהליך ההעברה צריך ליצור שני פרופילים של חיבורים ל-Datastream: אחד לקריאת נתונים ממסד נתונים תואם ל-MongoDB, ואחד לכתיבת הנתונים לקטגוריה של Cloud Storage.

השלבים האלה מתבססים על הפקודה gcloud datastream connection-profiles create. כשיוצרים פרופיל חיבור של Datastream באמצעות הפקודה הזו, מוחזרים מטא-נתונים שדומים לדוגמה הבאה:

metadata:
  '@type': type.googleapis.com/google.cloud.datastream.v1.OperationMetadata
  apiVersion: v1
  createTime: '2025-05-15T21:49:05.022509533Z'
  requestedCancellation: false
  target: projects/PROJECT_ID/locations/LOCATION/connectionProfiles/SRC_CONNECTION_PROFILE_NAME
  verb: create
name: projects/PROJECT_ID/locations/LOCATION/operations/operation-1747345744961-63533a26d9ee6-b2386fbf-204c28d6

אפשר להשתמש במזהה המודגש שמתחיל ב-operation- כדי לאחזר את הסטטוס של פעולת Datastream שצוינה. בדוגמה של הפלט שלמעלה, הפקודה הבאה של ה-CLI של gcloud מאחזרת פרטים על בקשת היצירה של פרופיל החיבור:

gcloud datastream operations describe \
operation-1747345744961-63533a26d9ee6-b2386fbf-204c28d6 \
--location="$LOCATION"

אפשר לבדוק את פרופילי החיבור של Datastream ואת הארטיפקטים האחרים של Datastream ב- Google Cloud console.

נכנסים לדף Datastream במסוף Google Cloud :

מעבר אל Datastream

פתרון בעיות בסטרימינג של Datastream

כשיוצרים מקור נתונים של Datastream באמצעות הפקודה gcloud datastream streams create, מוחזרים מטא-נתונים שדומים לדוגמה הבאה:

metadata:
  '@type': type.googleapis.com/google.cloud.datastream.v1.OperationMetadata
  apiVersion: v1
  createTime: '2025-05-14T19:31:20.209503095Z'
  requestedCancellation: false
  target: projects/PROJECT_ID/locations/LOCATION/streams/DATASTREAM_NAME
  verb: create
  name: projects/PROJECT_ID/locations/LOCATION/operations/operation-1747251080085-6351d97f63eb8-43204f78-35c87474

אפשר להשתמש במזהה המודגש שמתחיל ב-operation- כדי לאחזר את הסטטוס של פעולת Datastream שצוינה. בדוגמה של הפלט שלמעלה, הפקודה הבאה ב-CLI של gcloud מאחזרת פרטים על בקשת יצירת הזרם:

gcloud datastream operations describe \
operation-1747345744961-63533a26d9ee6-b2386fbf-204c28d6 \
--location="$LOCATION"

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

gcloud datastream streams describe $DATASTREAM_NAME --location=$LOCATION

אפשר לבדוק את הסטרימינג של Datastream ואת הארטיפקטים האחרים של Datastream במסוף Google Cloud .

נכנסים לדף Datastream במסוף Google Cloud :

מעבר אל Datastream

פתרון בעיות בצינור עיבוד הנתונים של Dataflow

אפשר לעקוב אחר ההפעלה של צינור עיבוד הנתונים של Dataflow ב Google Cloud מסוף.

נכנסים לדף Dataflow במסוף Google Cloud :

מעבר אל Dataflow

פתרון בעיות של שגיאות שאפשר לנסות שוב

שגיאות שאפשר לנסות שוב הן שגיאות זמניות שבסופו של דבר יצליחו. שגיאות נפוצות שאפשר לנסות שוב לפתור:

  • בעיות ברשת או בקישוריות
  • התנגשות בין עסקאות
  • חיבורים שנסגרו בגלל איזון עומסים

מסמכים שלא ניתן לכתוב ליעד בגלל שגיאות יישמרו בקטגוריה של Cloud Storage, במיקום שצוין בפרמטר deadLetterQueueDirectory של תבנית Dataflow. כברירת מחדל, המערכת תנסה שוב לפתור שגיאות שאפשר לנסות שוב לפתור אותן עד מספר הפעמים שצוין בפרמטר dlqMaxRetryCount.

אפשר לבדוק את תור ההודעות המתות (DLQ) ישירות מ-Cloud Storage. לשם כך, עוברים לנתיב שצוין בDLQ_LOCATION משתנה הסביבה. הנתיב יכיל היררכיה של תיקיות עם חותמות זמן, שמכילות רשומות JSON של מסמכים ועדכונים שלא ניתן היה לכתוב למסד הנתונים של Firestore עם תאימות ל-MongoDB.

בהנחה שתור ההודעות המתות מכיל רק מסמכים שנכשלו בגלל שגיאות שאפשר לנסות שוב, אפשר לנסות לרוקן את התור על ידי הפעלת תבנית Dataflow במצב שפועל רק בתור הזה. כדי לעשות זאת, מגדירים את הפרמטר runMode לערך retryDLQ, כמו בדוגמה הבאה:

DLQ_START_TIME="$(date +'%Y%m%d%H%M%S')"

gcloud dataflow flex-template run "dataflow-mongodb-to-firestore-$DLQ_START_TIME" \
--template-file-gcs-location gs://dataflow-templates-us-central1/latest/flex/Cloud_Datastream_MongoDB_to_Firestore \
--region $LOCATION \
--num-workers $NUM_WORKERS \
--temp-location $TEMP_OUTPUT_LOCATION \
--additional-user-labels "" \
--parameters inputFilePattern=$INPUT_FILE_LOCATION,\
inputFileFormat=avro,\
fileReadConcurrency=10,\
connectionUri=$FIRESTORE_CONNECTION_URI,\
databaseName=$FIRESTORE_DATABASE_NAME,\
shadowCollectionPrefix=shadow_,\
batchSize=500,\
deadLetterQueueDirectory=$DLQ_LOCATION,\
dlqRetryMinutes=10,\
dlqMaxRetryCount=500,\
processBackfillFirst=false,\
runMode=retryDLQ,\
directoryWatchDurationInMinutes=10,\
streamName=$DATASTREAM_NAME,\
stagingLocation=$STAGING_LOCATION,\
autoscalingAlgorithm=THROUGHPUT_BASED,\
maxNumWorkers=$MAX_WORKERS,\
workerMachineType=$WORKER_TYPE

פתרון בעיות שגיאות שלא ניתן לנסות שוב

שגיאות נפוצות שלא ניתן לנסות שוב:

  • סוגי BSON שלא נתמכים
  • נעשה שימוש בסוגי BSON שלא נתמכים בתור _id
  • ‫0L לא נתמך בתור _id
  • גודל המסמך גדול מהמגבלה של Firestore‏ (16MB)

שגיאות שלא ניתן לנסות שוב לפתור אותן יישמרו בקטגוריה של Cloud Storage, במיקום שצוין בפרמטר deadLetterQueueDirectory של תבנית Dataflow.

בדיקת תור ההודעות ללא מוצא (DLQ)

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

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

אפשר לבדוק את תור ההודעות להמתנה (DLQ) ישירות מ-Cloud Storage:

  1. מנווטים למיקום ב-Cloud Storage שצוין במשתנה הסביבה DLQ_LOCATION.
  2. בעץ של תיקיות מקוננות שנבנה לפי חותמות הזמן, מרחיבים את התיקיות כדי לראות את התוכן האחרון בתור. יכול להיות שהקבצים פוצלו והם ייראו כמו בדוגמה הבאה:

    קבצים בתור להודעות שלא ניתן למסור (DLQ)

  3. בודקים כל קובץ כדי לראות את המסמכים והעדכונים שלא הוחלו על מסד הנתונים של היעד. כל הודעה תכיל מטען ייעודי (payload) של האירוע המקורי והחריגה שהתרחשה.

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

  4. אפשר לבחור להפסיק את ההעברה ולהמשיך את תנועת האפליקציה למסד הנתונים המקורי, לטפל במגבלות הנתונים של Firestore במסד הנתונים המקורי ולהפעיל מחדש את כל תהליך ההעברה, או לטפל בבעיות הנתונים בכל אירוע של DLQ ולנסות שוב את העיבוד של האירועים האלה.

לכל שורה בכל קובץ של תור להודעות שלא ניתן למסור (DLQ) אפשר:

  • כדי להמיר את סוגי הנתונים הלא נתמכים לסוג נתונים חלופי, צריך לערוך את מטען הנתונים של המסמך. המטען הייעודי (payload) של המסמך יהיה בפורמט JSON קנוני מורחב. מבטאים את סוג הנתונים החדש בפורמט JSON קנוני מורחב.

  • מקצים מחדש את מסמך 0L‏ _id ל-Long חדש או לסוג נתונים אחר, כמו String. יכול להיות שיהיה צורך לבצע שינויים בלוגיקה של האפליקציה כדי להקצות מחדש לסוג נתונים אחר, אם הציפיות של הקוד הקיים הן שכל הערכים של _id הם מסוג Long.

  • אפשר לחלק מסמכים שגדולים מהמגבלה של Firestore (16MB) למסמכים קטנים יותר, או לדחוס את התוכן שלהם. עם זאת, כדי לטפל בנתונים, צריך לבצע שינויים בלוגיקה של האפליקציה.

  • מתעלמים מהאירוע על ידי מחיקת השורה מקובץ ה-DLQ.

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

אחרי שמבצעים את השינויים בקובצי ה-DLQ, אפשר להפעיל מחדש את האירועים האלה באמצעות תבנית Dataflow להעברה במצב retryDLQ.

הפקודה הבאה מפעילה פייפליין חדש של Dataflow שיעבד רק את הפריטים בתור Dead Letter Queue:

DLQ_START_TIME="$(date +'%Y%m%d%H%M%S')"

gcloud dataflow flex-template run "dataflow-mongodb-to-firestore-$DLQ_START_TIME" \
--template-file-gcs-location gs://dataflow-templates-us-central1/latest/flex/Cloud_Datastream_MongoDB_to_Firestore \
--region $LOCATION \
--num-workers $NUM_WORKERS \
--temp-location $TEMP_OUTPUT_LOCATION \
--additional-user-labels "" \
--parameters inputFilePattern=$INPUT_FILE_LOCATION,\
inputFileFormat=avro,\
fileReadConcurrency=10,\
connectionUri=$FIRESTORE_CONNECTION_URI,\
databaseName=$FIRESTORE_DATABASE_NAME,\
shadowCollectionPrefix=shadow_,\
batchSize=500,\
deadLetterQueueDirectory=$DLQ_LOCATION,\
dlqRetryMinutes=10,\
dlqMaxRetryCount=500,\
processBackfillFirst=false,\
runMode=retryDLQ,\
directoryWatchDurationInMinutes=10,\
streamName=$DATASTREAM_NAME,\
stagingLocation=$STAGING_LOCATION,\
autoscalingAlgorithm=THROUGHPUT_BASED,\
maxNumWorkers=$MAX_WORKERS,\
workerMachineType=$WORKER_TYPE

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