【问题标题】:Orchestrate Dataflow job using Dataflowpythonoperator使用 Dataflowpythonoperator 编排 Dataflow 作业
【发布时间】:2021-07-29 11:45:21
【问题描述】:

在气流上运行 Dataflow 作业时需要一些建议。我可以在本地运行 beam 作业,也可以在 dataflow 运行器上运行使用 Python sdk 的 wordcount 示例。但无法使用 DataflowPythonOperator 通过气流编排作业(不确定是否已弃用)。我使用 Dataflowjavaoperator 编排 jar 文件没有任何问题。

example_composer_..._operator = DataFlowPythonOperator(

    gcp_conn_id='gcp_default',
    task_id='composer_dataflow_python_...',
    py_file='gs://dataflow.../WordCountPython.py',
    job_name='Airflow ... Job ',
    py_options=None,
    dataflow_default_options=None,
    options=None,

关于我应该如何解决这个问题以及我是否做错了什么的任何建议。我应该使用 Python 运算符来调用 wordcountpython.py 文件吗

干杯

【问题讨论】:

  • 您是否遇到任何错误?运行它的结果是什么?

标签: airflow dataflow google-cloud-composer


【解决方案1】:

Dataflowjavaoperator 已弃用,您应该使用 DataflowCreateJavaJobOperator

您可以从提供商处导入运算符。使用示例:

from airflow.providers.google.cloud.operators.dataflow import DataflowCreateJavaJobOperator
task = DataflowCreateJavaJobOperator(
    gcp_conn_id="gcp_default",
    task_id="normalize-cal",
    jar="{{var.value.gcp_dataflow_base}}pipeline-ingress-cal-normalize-1.0.jar",
    options={
        "autoscalingAlgorithm": "BASIC",
        "maxNumWorkers": "50",
        "start": "{{ds}}",
        "partitionType": "DAY",
    },
    dag=dag,
)

如果你想使用 python 版本,那么你应该使用DataflowCreatePythonJobOperator。使用示例:

from airflow.providers.google.cloud.operators.dataflow import DataflowCreatePythonJobOperator

start_python_job = DataflowCreatePythonJobOperator(
    task_id="start-python-job",
    py_file=GCS_PYTHON,
    py_options=[],
    job_name='{{task.task_id}}',
    options={
        'output': GCS_OUTPUT,
    },
    py_requirements=['apache-beam[gcp]==2.21.0'],
    py_interpreter='python3',
    py_system_site_packages=False,
    location='europe-west3',
)

如果您运行的是 Airflow Google backport provider 来获取这些运算符:

pip install apache-airflow-backport-providers-google

如果您运行的是 Airflow >= 2.0.0,则可以通过安装 Google provider 来获取这些运算符:

pip install apache-airflow-providers-google

【讨论】:

  • 嗨,Elad 我可以使用 Python 运算符运行该作业,并将所需的参数保存在 config.py 文件中。在 Dag 中进行可调用的函数调用,谢谢
猜你喜欢
  • 2018-12-14
  • 2021-04-02
  • 1970-01-01
  • 2013-01-30
  • 2021-08-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多