תבנית Cloud Storage Text to BigQuery עם פונקציית UDF ב-Python

צינור הנתונים Cloud Storage Text to BigQuery with Python UDF הוא צינור נתונים של אצווה שקורא קובצי טקסט שמאוחסנים ב-Cloud Storage, מבצע בהם טרנספורמציה באמצעות פונקציה מוגדרת על ידי המשתמש (UDF) ב-Python, ומצרף את התוצאה לטבלה ב-BigQuery.

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

  • יוצרים קובץ JSON שמתאר את הסכימה של BigQuery.

    מוודאים שיש מערך JSON ברמה העליונה בשם BigQuery Schema ושהתוכן שלו תואם לתבנית {"name": "COLUMN_NAME", "type": "DATA_TYPE"}.

    תבנית האצווה Cloud Storage Text to BigQuery לא תומכת בייבוא נתונים לשדות STRUCT (Record) בטבלת היעד ב-BigQuery.

    קובץ ה-JSON הבא מתאר סכימה לדוגמה של BigQuery:

    {
      "BigQuery Schema": [
        {
          "name": "name",
          "type": "STRING"
        },
        {
          "name": "age",
          "type": "INTEGER"
        },
      ]
    }
  • יוצרים קובץ Python ‏ (.py) עם פונקציית ה-UDF שמספקת את הלוגיקה להמרת שורות הטקסט. הפונקציה צריכה להחזיר מחרוזת JSON.

    לדוגמה, הפונקציה הזו מפצלת כל שורה בקובץ CSV ומחזירה מחרוזת JSON אחרי שינוי הערכים.

    import json
    def process(value):
      data = value.split(',')
      obj = { 'name': data[0], 'age': int(data[1]) }
      return json.dumps(obj)

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

פרמטר תיאור
JSONPath gs:// הנתיב לקובץ ה-JSON שמגדיר את הסכימה של BigQuery, שמאוחסן ב-Cloud Storage. לדוגמה, gs://path/to/my/schema.json.
pythonExternalTextTransformGcsPath כתובת ה-URI של Cloud Storage של קובץ קוד Python שמגדיר את הפונקציה בהגדרת המשתמש (UDF) שרוצים להשתמש בה. לדוגמה, gs://my-bucket/my-udfs/my_file.py.
pythonExternalTextTransformFunctionName השם של פונקציית Python בהגדרת המשתמש (UDF) שרוצים להשתמש בה.
inputFilePattern הנתיב gs:// לטקסט ב-Cloud Storage שרוצים לעבד. לדוגמה, gs://path/to/my/text/data.txt.
outputTable שם הטבלה ב-BigQuery שרוצים ליצור כדי לאחסן בה את הנתונים המעובדים. אם משתמשים מחדש בטבלה קיימת ב-BigQuery, הנתונים מתווספים לטבלת היעד. לדוגמה, my-project-name:my-dataset.my-table.
bigQueryLoadingTemporaryDirectory הספרייה הזמנית לתהליך הטעינה של BigQuery. לדוגמה, gs://my-bucket/my-files/temp_dir.
useStorageWriteApi אופציונלי: אם true, צינור הנתונים משתמש ב- BigQuery Storage Write API ‏ (gRPC). ערך ברירת המחדל הוא false. מידע נוסף זמין במאמר בנושא שימוש ב-Storage Write API‏ (gRPC).
useStorageWriteApiAtLeastOnce אופציונלי: כשמשתמשים ב-Storage Write API‏ (gRPC), מציינים את סמנטיקת הכתיבה. כדי להשתמש ב, צריך להגדיר את הפרמטר הזה ל-true. כדי להשתמש בסמנטיקה של שליחה בדיוק פעם אחת, צריך להגדיר את הפרמטר לערך false. הפרמטר הזה רלוונטי רק כש-useStorageWriteApi הוא true. ערך ברירת המחדל הוא false.

פונקציה בהגדרת המשתמש

אפשר גם להרחיב את התבנית הזו על ידי כתיבת פונקציה בהגדרת המשתמש (UDF). התבנית קוראת ל-UDF לכל רכיב קלט. מטענים ייעודיים של רכיבים עוברים סריאליזציה כמחרוזות JSON. מידע נוסף זמין במאמר בנושא יצירת פונקציות מוגדרות על ידי המשתמש לתבניות Dataflow.

מפרט הפונקציה

המאפיינים של פונקציית UDF:

  • קלט: שורה של טקסט מקובץ קלט של Cloud Storage.
  • ‫Output: מחרוזת JSON שתואמת לסכימה של טבלת היעד ב-BigQuery.

הרצת התבנית

המסוף

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

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

  5. בתפריט הנפתח Dataflow template (תבנית Dataflow), בוחרים את התבנית Text Files on Cloud Storage to BigQuery with Python UDF (Batch) (קובצי טקסט ב-Cloud Storage ל-BigQuery עם פונקציית UDF של Python) (באצווה).
  6. בשדות הפרמטרים שמופיעים, מזינים את ערכי הפרמטרים.
  7. לוחצים על הפעלת העבודה.

gcloud

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

gcloud dataflow flex-template run JOB_NAME \
    --template-file-gcs-location gs://dataflow-templates/VERSION/flex/GCS_Text_to_BigQuery_Xlang \
    --region REGION_NAME \
    --parameters \
pythonExternalTextTransformFunctionName=PYTHON_FUNCTION,\
JSONPath=PATH_TO_BIGQUERY_SCHEMA_JSON,\
pythonExternalTextTransformGcsPath=PATH_TO_PYTHON_UDF_FILE,\
inputFilePattern=PATH_TO_TEXT_DATA,\
outputTable=BIGQUERY_TABLE,\
bigQueryLoadingTemporaryDirectory=PATH_TO_TEMP_DIR_ON_GCS

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

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

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

    • latest כדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית ההורה ללא תאריך בדלי – gs://dataflow-templates/latest/
    • שם הגרסה, כמו 2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית. אפשר למצוא את הגרסה הזו בתיקיית ההורה המתאימה עם התאריך בדלי – gs://dataflow-templates/
  • ‫REGION_NAME: האזור שבו רוצים לפרוס את עבודת Dataflow, לדוגמה: us-central1
  • ‫PYTHON_FUNCTION: השם של פונקציית Python בהגדרת המשתמש (UDF) שרוצים להשתמש בה.
  • ‫PATH_TO_BIGQUERY_SCHEMA_JSON: הנתיב ב-Cloud Storage אל קובץ ה-JSON שמכיל את הגדרת הסכימה
  • ‫PATH_TO_PYTHON_UDF_FILE: ה-URI של Cloud Storage של קובץ קוד Python שמגדיר את הפונקציה בהגדרת המשתמש (UDF) שבה רוצים להשתמש. לדוגמה, gs://my-bucket/my-udfs/my_file.py.
  • ‫PATH_TO_TEXT_DATA: הנתיב שלכם ב-Cloud Storage למערך הנתונים של הטקסט
  • ‫BIGQUERY_TABLE: שם הטבלה ב-BigQuery
  • ‫PATH_TO_TEMP_DIR_ON_GCS: הנתיב של Cloud Storage לספריית הזמנית

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": {
        "pythonExternalTextTransformFunctionName": "PYTHON_FUNCTION",
        "JSONPath": "PATH_TO_BIGQUERY_SCHEMA_JSON",
        "pythonExternalTextTransformGcsPath": "PATH_TO_PYTHON_UDF_FILE",
        "inputFilePattern":"PATH_TO_TEXT_DATA",
        "outputTable":"BIGQUERY_TABLE",
        "bigQueryLoadingTemporaryDirectory": "PATH_TO_TEMP_DIR_ON_GCS"
      },
      "containerSpecGcsPath": "gs://dataflow-templates/VERSION/flex/GCS_Text_to_BigQuery_Xlang",
   }
}

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

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

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

    • latest כדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית ההורה ללא תאריך בדלי – gs://dataflow-templates/latest/
    • שם הגרסה, כמו 2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית. אפשר למצוא את הגרסה הזו בתיקיית ההורה המתאימה עם התאריך בדלי – gs://dataflow-templates/
  • ‫LOCATION: האזור שבו רוצים לפרוס את עבודת Dataflow, לדוגמה: us-central1
  • ‫PYTHON_FUNCTION: השם של פונקציית Python בהגדרת המשתמש (UDF) שרוצים להשתמש בה.
  • ‫PATH_TO_BIGQUERY_SCHEMA_JSON: הנתיב ב-Cloud Storage אל קובץ ה-JSON שמכיל את הגדרת הסכימה
  • ‫PATH_TO_PYTHON_UDF_FILE: ה-URI של Cloud Storage של קובץ קוד Python שמגדיר את הפונקציה בהגדרת המשתמש (UDF) שבה רוצים להשתמש. לדוגמה, gs://my-bucket/my-udfs/my_file.py.
  • ‫PATH_TO_TEXT_DATA: הנתיב שלכם ב-Cloud Storage למערך הנתונים של הטקסט
  • ‫BIGQUERY_TABLE: שם הטבלה ב-BigQuery
  • ‫PATH_TO_TEMP_DIR_ON_GCS: הנתיב של Cloud Storage לספריית הזמנית

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