בדף הזה מופיע חומר עזר בנושא 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>existingClusterdataprocOnGce>ephemeralClusterdataprocServerless
דוגמאות:
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>existingClusterdataprocOnGce>ephemeralClusterdataprocServerless
דוגמה:
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:bigQuerydataprocServerlessdataprocOnGce>existingClusterdataprocOnGce>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:dbtdataform>airflowWorkerdataform>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"