【发布时间】: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