【问题标题】:Dataflow job through Airflow DAG通过 Airflow DAG 的数据流作业
【发布时间】:2020-01-25 07:22:39
【问题描述】:

我正在尝试通过 Airflow 中的 BashOperator 使用数据流运行器执行 apache 光束管道 python 文件。我知道如何将参数动态传递给 python 文件。我期待优化参数 - 避免单独发送所有参数。 示例 sn-p:

text_context.py

import sys

def run_awc_orders(*args, **kwargs):
    print("all arguments -> ",  args)

if __name__ == "__main__":
    print("all params -> ", sys.argv)
    run_awc_orders( sys.argv[1],  sys.argv[2], sys.argv[3])

my_dag.py

test_DF_job = BashOperator(
    task_id='test_DF_job',
    provide_context=True,
    bash_command="python /usr/local/airflow/dags/test_context.py {{ execution_date }} {{ next_execution_date }} {{ params.db_params.new_text }}  --runner DataflowRunner --key path_to_creds_json_file --project project_name --staging_location staging_gcp_bucket_location --temp_location=temp_gcp_bucket_location --job_name test-job",
    params={
              'db_params': {
                'new_text': 'Hello World'
              }
            },
    dag=dag
)

所以,这就是我们可以在气流 UI 的日志中看到的内容。

[2019-09-25 06:44:44,103] {bash_operator.py:128} INFO - all params ->  ['/usr/local/airflow/dags/test_context.py', '2019-09-23T00:00:00+00:00', '2019-09-24T00:00:00+00:00', '127.0.0.1']
[2019-09-25 06:44:44,103] {bash_operator.py:128} INFO - all arguments ->  ('2019-09-23T00:00:00+00:00', '2019-09-24T00:00:00+00:00', '127.0.0.1')
[2019-09-25 06:44:44,106] {bash_operator.py:132} INFO - Command exited with return code 0

【问题讨论】:

  • 嗨。你能具体分享你的问题吗?另外,仅供参考,Airflow 有一个 DataflowOperator,可让您启动 Dataflow 作业模板。
  • 感谢@Pablo 的分享。我正在使用 DataflowTemplateOperator,它运行良好。我已经在 docker 内尝试了 DataflowPythonOperator 以及 Python V.3+ 并且接收 apache-beam 模块未发现问题,即使它安装在容器内。所以我想通过 BashOperator 尝试一下。
  • 很遗憾听到这个消息。这听起来像是 Airflow 执行环境的问题:/

标签: python google-cloud-dataflow airflow apache-beam


【解决方案1】:

我认为推荐的方法是使用Airflow's DataflowPythonOperator,它直接接收 Python 和 Dataflow 选项。

你会这样做:

test_DF_job = DataflowPythonOperator(
  py_file='/usr/local/airflow/dags/test_context.py',
  py_options=[...],
  dataflow_default_options={...},
  dag=dag
)

【讨论】:

  • 这对于作为模块的管道如何工作?我假设你只需要在你的气流环境中安装管道模块,然后指定 python 文件来运行管道?
猜你喜欢
  • 1970-01-01
  • 2018-08-05
  • 2021-12-08
  • 2020-11-03
  • 2019-06-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多