Orchestration Pipelines DSL reference

בדף הזה מופיע חומר עזר בנושא Orchestration Pipelines DSL.

מגבלה בתצוגה המקדימה

במהלך תקופת התצוגה המקדימה, יש למערכת Orchestration Pipelines את המגבלות הבאות:

  • לפעולות pyspark ו-notebook:

    • יש תמיכה רק בקובץ requirements.txt אחד לכל הפעולות pyspark ו-notebook.
    • פלטפורמת Windows לא נתמכת ליצירת חבילות באמצעות הכלי uv.
    • יש תמיכה רק בחבילות Python עם קבצים בינאריים מוכנים מראש.
  • לפעולות sql:

    • ההגדרה inline במפתח query לא נתמכת.

מידע על הפורמט והערכים

צינורות מוגדרים בפורמט YAML, והם צריכים להיות מאוחסנים בקבצים נפרדים, קובץ אחד לכל צינור, במאגר.

ב-Orchestration Pipelines יש כמה דרכים להשתמש במשתנים בהגדרות של צינורות העברת הנתונים ובהגדרות הפריסה. לדוגמה, אפשר להגדיר משתנים מותאמים אישית, להשתמש בסודות של GitHub ולהחליף ערכי משתנים בשורת הפקודה. מידע נוסף זמין במאמר Variables, secrets, and substitution.

מידע נוסף על הוספת צינורות נוספים לחבילת הצינורות זמין במאמר הוספת צינור נוסף.

דוגמאות לקוד

המאגר orchestration-pipelines ב-GitHub כולל את דוגמאות הקוד העדכניות למגוון שילובים של פעולות ומנועים של צינורות. מומלץ להשתמש בדוגמאות האלה כנקודת התחלה כדי לבדוק את היכולות של Orchestration Pipelines.

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

הגדרת צינור עיבוד נתונים כוללת את המפתחות הבאים ברמה העליונה:

  • modelVersion: הגרסה של מודל הגדרת הצינור. הגרסה העדכנית של המודל היא 1.0.

  • pipelineId: מזהה ייחודי של הצינור. המזהה הזה נשאר עקבי בכמה פריסות ובכמה גרסאות, וכך מאפשר מעקב וניהול של ישות הפייפליין הלוגית.

  • description: תיאור צינור עיבוד הנתונים, שממופה לתיאור של Airflow DAG בסביבת ההפעלה.

  • owner: הבעלים של הפייפליין.

  • tags: מזהי מחרוזות שמוחלים על הפייפליין, ומשמשים לסינון הפייפליינים.

  • notifications: התראות על אירועים בפייפליין. סוגי ההתראות הנתמכים:

    • onPipelineFailure: אימייל על כשלים בצינור.

    כדי לקבל התראות, צריך להגדיר את שירותי האימייל של SendGrid בסביבת ה-Runner. הוראות מופיעות במאמר בנושא הגדרת התראות באימייל.

    דוגמה:

    notifications:
      onPipelineFailure:
        email: ["user1@example.com", "user2@example.com"]
    
  • runner: מציין את מנוע התזמור של היעד. שמור לשימוש בעתיד. מגדירים את הערך הזה כ-airflow.

  • defaults: הגדרת ערכי ברירת מחדל למאפיינים כמו project_id, ‏ location ו-executionConfig שחלים על כל הפעולות, אלא אם הם מוחלפים בפעולה ספציפית. אפשר לשנות את הערכים של המאפיינים project_id ו-location באמצעות מאפיינים של פעולות ספציפיות. אי אפשר לשנות את מאפיין executionConfig בפעולות נפרדות, והוא מציין את מספר הניסיונות החוזרים לכל הפעולות בצינור בשדה retries.

  • triggers: הגדרה של אופן ההפעלה של צינור עיבוד הנתונים:

    • אין ערך. עדיין אפשר להפעיל את צינור הנתונים באופן ידני.

    • schedule. הפעלת צינור עיבוד הנתונים לפי לוח זמנים, באמצעות ביטויי cron.

      דוגמה ללוח זמנים:

      triggers:
        - schedule:
            interval: "0 5 * * *"
            startTime: "2025-10-01T00:00:00"
            endTime: "2026-10-01T00:00:00"
            catchup: false
            timezone: "UTC"
      
  • actions

    מיפוי של משימות לביצוע. כל רשומה במיפוי מתאימה לפעולה אחת. מידע נוסף על פעולות

פעולות

פעולות בצינור מגדירות שלבים נפרדים בהרצת צינור. לכל פעולה צריך לציין מנוע או מסגרת. המנוע או המסגרת קובעים באילו משאבים נעשה שימוש כדי להפעיל את הפעולה.

‫Orchestration Pipelines תומך בפעולות הבאות:

  • ‫Pyspark ‏ (pyspark): הפעלת סקריפט PySpark.
  • ‫Notebook (notebook): הפעלת קובץ notebook.
  • שאילתת SQL‏ (notebook): הפעלת שאילתת SQL.
  • ‫Python (python): הרצת סקריפט Python.
  • פייפליין (pipeline): הפעלה של צינורות עיבוד נתונים.

‫Orchestration Pipelines תומך במנועים ובמסגרות הבאים:

  • dataprocOnGce > existingCluster: Managed Service for Apache Spark cluster identified by clusterName, project and location.

  • dataprocOnGce > ephemeral: Managed Service for Apache Spark cluster נוצר ונמחק אחרי הפעלת העבודה.

  • dataprocServerless: Managed Service for Apache Spark batch submission.

  • bigquery: משימה ב-BigQuery.

  • python > local: סקריפט Python שהופעל ב-Airflow Worker בסביבת ההפעלה.

  • dbt > airflowWorker: מודלים של dbt שהופעלו ב-worker של Airflow בסביבת runner באמצעות dbt-core.

  • dataform > airflowWorker: תהליכי עבודה של Dataform שמופעלים ב-worker של Airflow בסביבת ההרצה באמצעות dataform core cli.

  • dataform > dataformService: תהליכי עבודה של Dataform שמופעלים בשירות Dataform.

בטבלה הבאה מפורטים שילובי סוגי הפעולות, המנועים ומסגרות העבודה האפשריים. דוגמאות לקוד פעולה מופיעות בתיאורים של המנועים והמסגרות.

פעולה מנוע או Framework פלט אל
pyspark dataprocOnGce > existingCluster יומני משימות של Managed Service for Apache Spark
pyspark dataprocOnGce > ephemeralCluster יומני משימות של Managed Service for Apache Spark
pyspark dataprocServerless Managed Service for Apache Spark Batch logs
notebook dataprocOnGce > existingCluster קטגוריית הרצים, בספרייה composer_declarative_dags_resources
notebook dataprocOnGce > ephemeralCluster יומני משימות של Managed Service for Apache Spark
notebook dataprocServerless קטגוריית הרצים, בספרייה composer_declarative_dags_resources
sql bigquery הטבלה שצוינה בפרמטר destinationTable
sql dataprocServerless יומני עיבוד ברצף (batch processing) של Managed Service for Apache Spark.
python local (הפעלה מקומית) יומנים
pipeline dbt > airflowWorker יומנים ו-BigQuery
pipeline dataform > airflowWorker טבלה שצוינה ב-BigQuery
pipeline dataform > dataformService ב-Dataform

לכל הפעולות יש את המפתחות הנפוצים הבאים. מקשים אחרים תלויים בסוג הפעולה.

  • name: שם הפעולה. השם הזה ממופה לשם המשימה ב-Airflow בסביבת ההפעלה. אם פעולה מסוימת דורשת יותר ממשימת Airflow אחת, השם הזה ממופה לקבוצת המשימות.

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

  • executionTimeout: הזמן הקצוב לתפוגה להרצת הפעולה. דוגמאות: 1h, 30m, 40s.

python

פעולות מסוג python. להריץ סקריפטים של Python.

מקשים ספציפיים לסוג הפעולה:

  • mainFilePath: נתיב יחסי לקובץ סקריפט Python.
  • pythonCallable: השם של הפונקציה הניתנת להפעלה ב-Python שרוצים להריץ בסקריפט Python.
  • opKwargs: מיפוי של ארגומנטים של מילות מפתח לאופרטור.
  • (אופציונלי) environment: מריצים את הסקריפט בסביבה וירטואלית של Python שנוצרה באופן דינמי.

    • requirements: דרישות לסביבה הווירטואלית. הדרישות נבדקות בזמן הריצה.

      • inline: הדרישות מצוינות בשורה.

        • list: רשימת הדרישות. מפרטים את הדרישות האישיות בהתאם ל-PEP-508.

          דוגמה:

          environment:
            requirements:
              inline:
                list: ["pandas>=2.0.0"]
          
      • (חלופה) path: הנתיב לקובץ עם הדרישות. הדרישות בקובץ הזה צריכות להיות מפורטות בהתאם ל-PEP-508.

        דוגמה:

        environment:
          requirements:
            path: "scripts/requirements.txt"
        
    • systemSitePackages: אם הערך הוא true, הסביבה הווירטואלית מקבלת בירושה חבילות מהספרייה site-packages של העובד (worker) של Airflow. אתם יכולים להתקין חבילות PyPI בהתאמה אישית בסביבת ה-Runner.

  • engine:

    • local: ביצוע מקומי בסביבת ההרצה

דוגמה:

מקומי

modelVersion: "1.0"
pipelineId: "python-virtual-env"
description: "A pipeline that runs Python script on isolated virtual environment with specified dependencies."
runner: "airflow"
owner: "data-eng-team"
tags: ["python_virtual_environment", "example"]

defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 1

triggers:
  - schedule:
      interval: "0 5 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"

actions:
  - python:
      name: "first_script_run"
      executionTimeout: "30m"
      mainFilePath: "scripts/venv_test_script_1.py"
      pythonCallable: "main"
      engine:
        local: {}
      environment:
        systemSitePackages: true
        requirements:
          path: "scripts/requirements.txt"

  - python:
      name: "second_script_run"
      executionTimeout: "30m"
      mainFilePath: "scripts/venv_test_script_1.py"
      pythonCallable: "main"
      engine:
        local: {}
      environment:
        requirements:
          inline:
            list: ["pandas>=2.0.0"]
        systemSitePackages: true
      dependsOn: ["first_script_run"]

  - python:
      name: "third_script_run"
      executionTimeout: "30m"
      mainFilePath: "scripts/venv_test_script_2.py"
      pythonCallable: "main"
      engine:
        local: {}
      opKwargs:
        api_endpoint: "https://api.my-vendor.example.com/v1/status"
        api_key_secret_name: "my-vendor-api-key"
      dependsOn: ["first_script_run"]

pyspark

פעולות מסוג pyspark. להריץ סקריפטים של PySpark.

מקשים ספציפיים לסוג הפעולה:

  • mainFilePath: נתיב יחסי לסקריפט PySpark.
  • archiveUris: רשימה של מזהי URI של ארכיונים לשימוש בפעולה הזו.
  • stagingBucket: קטגוריית Cloud Storage לשימוש בפעולה הזו.
  • pyFiles: רשימה של קובצי Python לשימוש בעבודת Spark הזו.
  • environment: הגדרת סביבת Python.

    • requirements: קובץ הדרישות של Python שבו רוצים להשתמש.

      • path: הנתיב לקובץ עם הדרישות. הדרישות בקובץ הזה צריכות להיות מפורטות בהתאם ל-PEP-508.
  • engine:

    • dataprocOnGce > existingCluster
    • dataprocOnGce > ephemeralCluster
    • dataprocServerless

דוגמאות:

existingCluster

modelVersion: "1.0"
pipelineId: "dataproc-existing-cluster-script-pipeline-pyfiles"
description: "A pipeline with a Dataproc Existing cluster running a PySpark task with additional python files."
runner: "airflow"
owner: "data-eng-team"
tags: ["dataproc_create_batch", "example", "pyfiles"]
defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0
triggers:
  - schedule:
      interval: "0 4 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"
actions:
  - pyspark:
      name: "run-pyspark-with-pyfiles-on-existing-cluster"
      engine:
        dataprocOnGce:
          existingCluster:
            clusterName: "example-dataproc-cluster"
            location: "us-central1"
            projectId: "example-project"
            impersonationChain: "example-account@example-project.iam.gserviceaccount.com"
      mainFilePath: "scripts/my_spark_job_with_pyfiles.py"
      pyFiles:
        - "data/lib1.py"

ephemeralCluster

pipelineId: "dataproc-ephemeral-inline-pyspark"
description: "A pipeline with a Dataproc Ephemeral cluster running the PySpark task."
runner: "airflow"
owner: "data-eng-team"
tags: ["dataproc_create_batch", "example"]
modelVersion: "1.0"

defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0

triggers:
  - schedule:
      interval: "0 5 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"

actions:
  - pyspark:
      name: "run_dataproc_ephemeral"
      executionTimeout: "1h"
      mainFilePath: "scripts/my_spark_job_with_pyfiles.py"
      pyFiles:
        - "data/lib1.py"
      engine:
        dataprocOnGce:
          ephemeralCluster:
            projectId: "example-project"
            location: "us-central1"
            clusterName: "ephemeral-cluster-inline"
            resourceProfile:
              inline:
                config:
                  masterConfig:
                    numInstances: 1
                    machineTypeUri: "n1-standard-4"
                    diskConfig:
                      bootDiskType: "pd-standard"
                      bootDiskSizeGb: 1024
                  workerConfig:
                    numInstances: 2
                    machineTypeUri: "n1-standard-4"
                    diskConfig:
                      bootDiskType: "pd-standard"
                      bootDiskSizeGb: 1024
            properties:
              spark.submit.deployMode: "cluster"

dataprocServerless

modelVersion: "1.0"
pipelineId: "dataproc-create-batch-pipeline-pyspark"
description: "A pipeline with a Batch task."
runner: "airflow"
owner: "data-eng-team"
tags: ["dataproc_create_batch", "example"]

defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0

actions:
  - pyspark:
      name: "run-pyspark-on-dataproc-serverless"
      mainFilePath: "scripts/my_spark_job_with_pyfiles.py"
      pyFiles:
        - "data/lib1.py"
      engine:
        dataprocServerless:
          location: "us-central1"
          resourceProfile:
            inline:
              environmentConfig:
                executionConfig:
                  serviceAccount: "example-account@example-project.iam.gserviceaccount.com"
                  networkUri: "projects/example-project/global/networks/default"
              runtimeConfig:
                version: "2.3"
                properties:
                  spark.app.name: "run-notebook-on-dataproc-serverless"
                  spark.executor.instances: "2"
                  spark.driver.cores: "4"

notebook

פעולות מסוג notebook. מריצים מחברת ‎ .ipynb דרך Papermill.

מקשים ספציפיים לסוג הפעולה:

  • mainFilePath: נתיב יחסי לקובץ ה-Notebook.
  • archiveUris: רשימה של מזהי URI של ארכיונים לשימוש בפעולה הזו.
  • stagingBucket: קטגוריית Cloud Storage לשימוש בפעולה הזו.
  • environment: הגדרת סביבת Python.

    • requirements: קובץ הדרישות של Python שבו רוצים להשתמש.

      • path: הנתיב לקובץ עם הדרישות. הדרישות בקובץ הזה צריכות להיות מפורטות בהתאם ל-PEP-508.
  • engine:

    • dataprocOnGce > existingCluster
    • dataprocOnGce > ephemeralCluster
    • dataprocServerless

דוגמה:

dataprocServerless

modelVersion: "1.0"
pipelineId: "dataproc-create-batch-pipeline-resource-profile-gcs-overrides"
description: "A pipeline with a DataprocCreateBatch task."
runner: "airflow"
owner: "data-eng-team"
tags: ["dataproc_create_batch", "example"]

defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0

triggers:
  - schedule:
      interval: "0 5 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"

actions:
  - notebook:
      name: "run-notebook-on-dataproc-serverless"
      mainFilePath: "scripts/notebookWithArchivesCheck.ipynb"
      archiveUris:
        - "gs://example-bucket-additional-data/custom_venv.tar.gz"
      staging_bucket: "example-bucket-additional-data-additional-data"
      engine:
        dataprocServerless:
          location: "us-central1"
          resourceProfile:
            externalConfigPath: "gs://example-runner-bucket/data/batch-resource-profile.yml"
            overrides:
              runtimeConfig:
                properties:
                  spark.app.name: "run-notebook-on-dataproc-serverless"
                  spark.executor.instances: "2"
                  spark.driver.cores: "4"

sql

פעולות מסוג sql. הפעלת שאילתות SQL.

מקשים ספציפיים לסוג הפעולה:

  • query: מגדיר שאילתה.

    • path: השאילתה מוגדרת בקובץ שנמצא בנתיב היחסי לקובץ הגדרות הפריסה.
    • inline: השאילתה מוגדרת בשורה.

  • engine:

    • bigQuery
    • dataprocServerless
    • dataprocOnGce > existingCluster
    • dataprocOnGce > ephemeralCluster

bigQuery

modelVersion: "1.0"
pipelineId: "sql-on-bigquery"
description: "A pipeline with a BigQueryInsertJob task."
runner: "airflow"
tags: ["bigquery", "example"]
owner: "data-eng-team"

defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0

triggers:
  - schedule:
      interval: "0 5 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"

actions:
  - sql:
      name: "run_bigquery_insert_job_create"
      query:
        inline: "CREATE TABLE IF NOT EXISTS `example-project.example_dataset.example_table` (id INT64, name STRING, timestamp TIMESTAMP );"
      engine:
        bigquery:
          location: "US"
  - sql:
      name: "run_bigquery_insert_job_select"
      query:
        path: "sql-scripts/count_rows.sql"
      engine:
        bigquery:
          location: "US"
          destinationTable: "example-project.example_dataset.example_table_query_results"
      dependsOn:
        - "run_bigquery_insert_job_create"

dataprocServerless

modelVersion: "1.0"
pipelineId: "sql-on-dataproc-serverless"
description: "A pipeline with a cluster running the PySpark task."
runner: "airflow"
owner: "data-eng-team"
tags: ["dataproc-serverless", "example"]

defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0

triggers:
  - schedule:
      interval: "0 5 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"

actions:
  - sql:
      name: "run-sql-on-dataproc"
      query:
        path: "sql-scripts/test_query.sql"
      engine:
        dataprocServerless:
          location: "us-central1"
          impersonationChain: "example-account@example-project.iam.gserviceaccount.com"
          resourceProfile:
            inline:
              runtimeConfig:
                properties:
                  spark.sql.catalog.bigquery: "com.google.cloud.spark.bigquery.BigQueryCatalog"
                  spark.sql.catalog.bigquery.project: "example-project"
                  dataproc.sparkBqConnector.version: "0.42.3"

existingCluster

modelVersion: "1.0"
pipelineId: "sql-on-dataproc-gce-existing"
description: "A pipeline running a SQL query on an existing Dataproc cluster."
runner: "airflow"
owner: "data-eng-team"
tags: ["dataproc-gce", "example", "existing-cluster"]

defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0

triggers:
  - schedule:
      interval: "0 5 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"

actions:
  - sql:
      name: "run-sql-on-existing-cluster"
      query:
        path: "sql-scripts/test_query.sql"
      engine:
        dataprocOnGce:
          existingCluster:
            clusterName: "cluster-sql"
            projectId: "example-project"
            location: "us-central1"
            impersonationChain: "example-account@example-project.iam.gserviceaccount.com"
            properties:
              spark.sql.catalog.bigquery: "com.google.cloud.spark.bigquery.BigQueryCatalog"
              spark.sql.catalog.bigquery.project: "example-project"
              # This field is needed and important for the Spark-BigQuery connector.
              spark.jars: "gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.42.3.jar"

ephemeralCluster

modelVersion: "1.0"
pipelineId: "sql-on-dataproc-gce-ephemeral"
description: "A pipeline running a SQL query on an ephemeral Dataproc cluster."
runner: "airflow"
owner: "data-eng-team"
tags: ["dataproc-gce", "example", "ephemeral-cluster"]

defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0

triggers:
  - schedule:
      interval: "0 5 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"

actions:
  - sql:
      name: "run-sql-on-ephemeral-cluster"
      query:
        path: "sql-scripts/test_query.sql"
      engine:
        dataprocOnGce:
          ephemeralCluster:
            clusterName: "example-ephemeral-cluster"
            projectId: "example-project"
            location: "us-central1"
            impersonationChain: "example-account@example-project.iam.gserviceaccount.com"
            properties:
              spark.sql.catalog.bigquery: "com.google.cloud.spark.bigquery.BigQueryCatalog"
              spark.sql.catalog.bigquery.project: "example-project"
              # This field is needed and important for the Spark-BigQuery connector.
              spark.jars: "gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.42.3.jar"
            resourceProfile:
              inline:
                clusterConfig:
                  gceClusterConfig:
                    zoneUri: "us-central1-a"
                    metadata: # This metadata field is needed and important for the Spark-BigQuery connector.
                      SPARK_BQ_CONNECTOR_URL: "gs://spark-lib/bigquery/spark-3.5-bigquery-0.43.1.jar"
                  masterConfig:
                    numInstances: 1
                    machineTypeUri: "n1-standard-4"
                  workerConfig:
                    numInstances: 2
                    machineTypeUri: "n1-standard-4"

פייפליין

פעולות מסוג pipeline. הפעלת צינור לעיבוד נתונים.

מקשים ספציפיים לסוג הפעולה:

  • framework:

    • dbt
    • dataform > airflowWorker
    • dataform > dataformService

דוגמאות:

dbt

modelVersion: "1.0"
pipelineId: "dbt-pipeline"
description: "A pipeline that runs a dbt workflow."
runner: "airflow"
owner: "example-owner"
defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0
triggers:
  - schedule:
      interval: "0 5 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"
actions:
  - pipeline:
      name: "dbt-action"
      framework:
        dbt:
          airflowWorker:
            projectDirectoryPath: "dbt_project"
            selectModels: ["model_1", "model_2"]

dataform>airflowWorker

modelVersion: "1.0"
pipelineId: "dataform-pipeline-local"
description: "A pipeline that runs a Dataform workflow."
runner: "airflow"
owner: "example-owner"
defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0
triggers:
  - schedule:
      interval: "0 5 * * *"
      startTime: "2025-10-01T00:00:00"
      endTime: "2026-10-01T00:00:00"
      catchup: false
      timezone: "UTC"
actions:
  - pipeline:
      name: "run_dataform"
      framework:
        dataform:
          airflowWorker:
            projectDirectoryPath: "dataform_local"

dataform>dataformService

modelVersion: "1.0"
pipelineId: dataform-service
description: "A pipeline that runs a Dataform workflow."
owner: "example-owner"
tags: ["dataform", "example"]
runner: "airflow"

defaults:
  projectId: "example-project"
  location: "us-central1"
  executionConfig:
    retries: 0

actions:
  - pipeline:
      name: "run_dataform"
      framework:
        dataform:
          dataformService:
            location: "us-central1"
            projectId: "example-project"
            repositoryId: "example-repository"
            workflowInvocation:
              workflowConfig: "projects/example-project/locations/us-central1/repositories/example-repository/workflowConfigs/example-test"

  - pipeline:
      name: "run_dataform_compilation"
      framework:
        dataform:
          dataformService:
            location: "us-central1"
            projectId: "example-project"
            repositoryId: "example-repository"
            workflowInvocation:
              compilationResult: "projects/example-project/locations/us-central1/repositories/tets/compilationResults/8027975f-9692-42da-a13b-2267a4909f76"
              invocationConfig:
                serviceAccount: "example-account@example-project.iam.gserviceaccount.com"

מנועי חיפוש

מנועים שמשמשים בפעולות.

dataprocOnGce > existingCluster

ההרצה מתבצעת באשכול קיים של Managed Service for Apache Spark שמזוהה לפי clusterName, ‏ project ו-location.

אפשר לנהל את האשכול שצוין בהגדרות הפריסה או באופן ידני ב-Managed Service for Apache Spark. מומלץ לשדרג את האשכול באופן קבוע.

מפתחות:

  • clusterName: שם האשכול
  • location: האזור שבו נמצא האשכול
  • projectId: מזהה הפרויקט שבו נמצא האשכול
  • properties: מיפוי של מאפייני עבודות Spark.

דוגמה:

engine:
  dataprocOnGce:
    existingCluster:
      clusterName: "example-dataproc-cluster"
      location: "us-central1"
      projectId: "example-project"
      impersonationChain: "example-account@example-project.iam.gserviceaccount.com"

dataprocOnGce > ephemeralCluster

ההרצה מתבצעת באשכול זמני של Managed Service for Apache Spark, שנוצר ונמחק אחרי הרצת העבודה.

מפתחות:

  • clusterName: שם האשכול
  • location: האזור שבו נמצא האשכול
  • projectId: מזהה הפרויקט שבו נמצא האשכול
  • impersonationChain: שרשרת התחזות לחשבון שירות לשימוש בהרצת הפעולה.
  • resourceProfile: פרופיל משאב של אשכול Managed Service for Apache Spark.

    תיאור של השדות הזמינים מופיע במאמר ClusterConfig במאמרי העזרה של Managed Service for Apache Spark.

    אפשר לציין פרופיל משאב בדרכים הבאות:

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

    אפשר להחיל שינויים על פרופיל המשאב שצוין באמצעות המפתח override. ההחלפות מוחלות באמצעות מיזוג עמוק על פרופיל המקור שצוין.

  • properties: מיפוי של מאפייני עבודות Spark.

דוגמה:

engine:
  dataprocOnGce:
    ephemeralCluster:
      projectId: "example-project"
      location: "us-central1"
      clusterName: "example-ephemeral-cluster"
      resourceProfile:
        inline:
          config:
            masterConfig:
              numInstances: 1
              machineTypeUri: "n1-standard-4"
              diskConfig:
                bootDiskType: "pd-standard"
                bootDiskSizeGb: 1024
            workerConfig:
              numInstances: 2
              machineTypeUri: "n1-standard-4"
              diskConfig:
                bootDiskType: "pd-standard"
                bootDiskSizeGb: 1024
      properties:
        spark.submit.deployMode: "cluster"

dataprocServerless

הפעלה ב-Managed Service for Apache Spark שליחת עיבוד ברצף.

מפתחות:

  • location: האזור שבו צריך להריץ את משימת Spark.
  • impersonationChain:שרשרת התחזות לחשבון שירות לשימוש בהרצת הפעולה.
  • resourceProfile: פרופיל משאבים של Managed Service for Apache Spark.

    אפשר לציין פרופיל משאב בדרכים הבאות:

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

    המפתחות הבאים מציינים את ההגדרה של פרופיל המשאב:

    • environmentConfig: הגדרת הסביבה
    • runtimeConfig: הגדרת זמן ריצה

    תיאור של השדות הזמינים מופיע במאמרים בנושא RuntimeConfig וEnvironmentConfig במאמרי העזרה של Managed Service for Apache Spark.

    אפשר להחיל שינויים על פרופיל המשאב שצוין באמצעות המפתח override. השינויים חלים על פרופיל המשאב שצוין באמצעות מיזוג עמוק.

דוגמה (בתוך השורה):

engine:
  dataprocServerless:
    location: "us-central1"
    resourceProfile:
      inline:
        environmentConfig:
          executionConfig:
            serviceAccount: "example-account@example-project.iam.gserviceaccount.com"
            networkUri: "projects/example-project/global/networks/default"
        runtimeConfig:
          version: "2.3"
          properties:
            spark.app.name: "run-notebook-on-dataproc-serverless"
            spark.executor.instances: "2"
            spark.driver.cores: "4"

דוגמה (נתיב חיצוני ושינויים):

      engine:
        dataprocServerless:
          location: "us-central1"
          resourceProfile:
            externalConfigPath: "gs://example-runner-bucket/data/batch-resource-profile.yml"
            overrides:
              runtimeConfig:
                properties:
                    spark.app.name: "run-notebook-on-dataproc-serverless"
                    spark.executor.instances: "2"
                    spark.driver.cores: "4"

bigQuery

הפעלה כעבודה ב-BigQuery.

מפתחות:

  • location: האזור שבו נמצאת טבלת היעד.
  • destinationTable: טבלת BigQuery שאליה יועברו הנתונים
  • impersonationChain: שרשרת התחזות לחשבון שירות לשימוש בהרצת הפעולה.

דוגמה:

      engine:
        bigquery:
          location: "US"
          destinationTable: "example-project.example_dataset.example_table_query_results"

מקומי

ההרצה מתבצעת באופן מקומי בסביבת הרץ.

במאמר בנושא פעולת python מוסבר איך להגדיר את הסביבה הווירטואלית.

דוגמה:

    engine:
      local: {}

‫Frameworks

מסגרות שמשמשות בפעולות.

dbt > airflowWorker

הפעלת מודל dbt שהופעל ב-worker של Airflow בסביבת ההרצה באמצעות dbt-core.

מפתחות:

  • projectDirectoryPath: נתיב יחסי לתיקייה שמכילה את פרויקט DBT.
  • selectModels: רשימת המודלים שרוצים לכלול בהרצה לפי שם (שווה ל-dbt --select).
  • tags: רשימת המודלים שייכללו בהרצה לפי תג (שווה ל-dbt --select).

דוגמה:

framework:
  dbt:
    airflowWorker:
      projectDirectoryPath: "dbt_project"
      selectModels: ["model_1", "model_2"]

dataform > airflowWorker

תהליכי עבודה של Dataform שמופעלים על Airflow worker בסביבת ה-runner באמצעות dataform core cli.

מפתחות:

  • projectDirectoryPath: נתיב יחסי לתיקייה שמכילה את ההגדרות של תהליך העבודה ב-Dataform.

דוגמה:

framework:
  dataform:
    airflowWorker:
      projectDirectoryPath: "dataform_local"

dataform > dataformService

הפעלת תהליכי עבודה של Dataform שמופעלים בשירות Dataform.

מפתחות:

  • location: המיקום שבו נמצא מאגר Dataform.
  • projectId: הפרויקט שבו נמצא מאגר Dataform.
  • repositoryId: מזהה מאגר Dataform
  • workflowInvocation: הגדרה להפעלת תהליך העבודה, שמציינת אילו פעולות להריץ. ראו WorkflowInvocation.

דוגמה:

framework:
  dataform:
    dataformService:
      location: "us-central1"
      projectId: "example-project"
      repositoryId: "example-repository"
      workflowInvocation:
        compilationResult: "projects/example-project/locations/us-central1/repositories/example-repository/compilationResults/example-compilation-id"
        invocationConfig:
          serviceAccount: "example-account@example-project.iam.gserviceaccount.com"