【问题标题】:How Airflow can create a dataflow job from a python operator?Airflow 如何从 python 运算符创建数据流作业?
【发布时间】:2019-06-21 09:32:29
【问题描述】:

当我通过命令行运行我的 Beam 管道时,使用直接运行器或数据流运行器,它工作正常......

例子:

$ python my_pipeline.py --key /path/to/gcp/service/key.json --project gcp_project_name

但是当我尝试使用气流时,我有两个选择,bash 运算符或 python 运算符。

使用 bash 运算符会成功,但会限制我使用气流功能的能力。

但我想做的是将它作为 python 运算符运行。所以我将模块导入到airflow dg文件中,然后作为python操作符运行。

如果我使用本地运行器也可以正常工作,但是当我将其更改为数据流运行器时,在 GCP 数据流上创建作业后它会失败并出现此错误

ImportError: No module named airflow.bin.cli

Airflowpython 运算符创建 dataflow 作业我缺少什么?

【问题讨论】:

  • 你有没有试过通过 python 操作符运行其他 python 任务?你成功了吗?我曾经玩过 Airflow,这个错误看起来像 Airflow 设置问题。
  • 是的,我有多个 python 运算符,其中一个只创建一个数据流作业。在气流创建数据流作业之后出现此错误......

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


【解决方案1】:

好的, 这不是完美的解决方案,但您可以使用

DataFlowPythonOperator()

它将运行我们之前提到的完全相同的 bash 命令。 这是一种解决方法,不等于PythonOperator,但更像是运行BashOperator... 在当前情况下仍然无法使用 Airflow 功能的优势(如 xcom)... Docs

【讨论】:

  • 我们如何在 DataFlowPythonOperator() 中传递输入文件位置和输出位置?
猜你喜欢
  • 1970-01-01
  • 2020-01-25
  • 2022-10-13
  • 2023-03-20
  • 1970-01-01
  • 2022-11-17
  • 1970-01-01
  • 1970-01-01
  • 2019-06-04
相关资源
最近更新 更多