You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Apache Airflow中配置通过API更新状态的被动任务?

在Apache Airflow中配置被动任务(通过API更新状态)

要实现这类仅通过外部API更新状态的被动任务,核心是让Airflow任务不执行实际业务逻辑,只等待外部系统触发状态变更,以下是两种常用实现方案:

方案1:自定义等待型Operator(推荐)

这种方案让任务启动后一直处于运行状态,持续检查自身状态,直到外部API修改它的状态为成功/失败/跳过。

自定义Operator代码

import time
from airflow.models.baseoperator import BaseOperator
from airflow.utils.state import State
from airflow.exceptions import AirflowSkipException, AirflowFailException

class PassiveTaskOperator(BaseOperator):
    def execute(self, context):
        task_instance = context['task_instance']
        # 持续检查任务状态,直到外部触发变更
        while True:
            task_instance.refresh_from_db()
            current_state = task_instance.state
            if current_state == State.SUCCESS:
                return
            elif current_state == State.FAILED:
                raise AirflowFailException("被动任务被外部标记为失败")
            elif current_state == State.SKIPPED:
                raise AirflowSkipException("被动任务被外部标记为跳过")
            # 每隔30秒刷新一次状态,避免Airflow判定任务无响应
            time.sleep(30)

在DAG中使用

from airflow import DAG
from datetime import datetime
from your_module import PassiveTaskOperator

with DAG(
    dag_id='passive_approval_workflow',
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,  # 手动触发或由其他DAG触发
    catchup=False
) as dag:
    manual_approval = PassiveTaskOperator(
        task_id='await_manual_approval',
        execution_timeout=None,  # 取消超时限制,或设置足够长的超时时间
        retries=0  # 不需要重试,状态由外部控制
    )

方案2:使用DummyOperator配合API强制更新状态

如果不需要任务保持运行中状态,可以用空的DummyOperator,直接通过API修改任务实例的状态。

DAG定义

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from datetime import datetime

with DAG(
    dag_id='external_system_triggered_dag',
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    external_task = DummyOperator(
        task_id='external_process_task'
    )

API调用示例

需要先获取目标任务实例的dag_id、task_id和execution_date(ISO 8601格式,如2024-05-20T00:00:00+00:00),然后通过Airflow REST API修改状态:

  • 标记任务为成功:
curl -X PATCH \
  http://<airflow-webserver-url>/api/v1/dags/<dag_id>/dagRuns/<execution_date>/taskInstances/<task_id> \
  -H "Authorization: Bearer <your-auth-token>" \
  -H "Content-Type: application/json" \
  -d '{"state": "success"}'
  • 标记任务为失败:
curl -X PATCH \
  http://<airflow-webserver-url>/api/v1/dags/<dag_id>/dagRuns/<execution_date>/taskInstances/<task_id> \
  -H "Authorization: Bearer <your-auth-token>" \
  -H "Content-Type: application/json" \
  -d '{"state": "failed"}'

关键注意事项

  • 权限控制:确保调用API的账号拥有edit级别的权限,可通过Airflow RBAC角色配置。
  • 超时设置:方案1中需将execution_timeout设为None或足够长的时间,防止Airflow自动将无响应的任务标记为失败。
  • 状态一致性:外部系统必须准确传递execution_date,否则无法匹配到对应的任务实例。
  • 重复操作防护:外部系统要保证只对同一个任务实例触发一次状态变更,避免状态混乱。

内容的提问来源于stack exchange,提问作者tony

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 07:01:08