【问题标题】:How can I run specific task/s from the Airflow dag如何从 Airflow dag 运行特定任务
【发布时间】:2021-06-19 13:13:30
【问题描述】:

当前气流状态:

ml_processors = [a, b, c, d, e]

abc_task >> ml_processors (all ml models from a to e run in parallel after abc task is successfully completed)
ml_processors >> xyz_task (once a to e all are successful xyz task runs)

问题陈述:在某些情况下,其中一种机器学习模型(气流中的任务)以更高的准确性进入新版本,并且我们想要重新处理我们的数据。现在让我们说 c_processor 获得了新版本,并且需要重新处理才能重新处理该处理器的数据。在这种情况下,我只想运行 c_processor >> xyz_task。

我知道/尝试过的事情

  1. 我知道我可以在成功的 dag 运行中返回并在一段时间内清除任务以仅运行特定任务。但是,当我说要重新运行 c_processor、d_classifier 时,这种方式可能不是很有效。我最终会在这里执行两个步骤:

  2. c_processor >> xyz_task

  3. d_processor >> xyz_task 我想避免

  4. 我读过“气流回填”,但看起来更像是整个 dag,而不是 dag 中的特定/选定任务

环境/设置

  1. 使用 google composer 环境。
  2. 在 GCP 存储中上传文件时触发 Dag。

我很想知道是否有任何其他方法可以仅从气流 dag 重新运行特定任务。

【问题讨论】:

  • 遗憾的是,这不是气流的好用例。你可以做的是你清除 c_processor 和 d_processor,当它们都完成时,xyz_task 将只运行一次。

标签: google-cloud-platform dockerfile airflow google-cloud-composer


【解决方案1】:

"clear"1 还允许您使用 --task-regex 标志清除 DAG 中的某些特定任务。在这种情况下,您可以运行 airflow tasks clear --task-regex "[c|d]_processor" --downstream -s 2021-03-22 -e 2021-03-23 <dag_id>,这会清除 c 和 d 处理器及其下游的状态。

但需要注意的是,这也会清理原始任务运行的状态。

【讨论】:

  • 感谢斌的回复。这听起来像我上面提到的#1,但使用这个命令会很有用。
  • 是的,您至少需要知道要清理或回填哪些任务。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-01-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多