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

Airflow中如何在下游任务中检查上游任务状态并获取错误信息以执行决策

刚好之前处理过类似的Airflow审计场景,给你分享几个最实用的实现方案,都是生产环境验证过的:

方案1:通过TaskInstance API直接查询上游状态与错误信息

这个方案最省心,不需要对上游任务做太多修改,直接在下游审计任务里通过Airflow的TaskInstance对象获取上游任务的状态和异常信息。

下游审计任务建议用PythonOperator(或者Airflow 2.x+的@task装饰器),代码示例如下:

from airflow.decorators import task

@task(trigger_rule='all_done')
def audit_task(**context):
    # 获取当前任务实例
    ti = context['ti']
    
    # 替换成你的上游任务ID
    upstream_task_id = "upstream_python_or_bq_task"
    # 获取上游任务的实例对象
    upstream_ti = ti.get_task_instance(task_id=upstream_task_id)
    
    # 获取上游任务的最终状态(success/failed/skipped等)
    upstream_state = upstream_ti.state
    
    # 获取错误信息(仅当状态为failed时有效)
    error_msg = None
    if upstream_state == "failed":
        # 直接获取异常对象的字符串表示
        error_msg = str(upstream_ti.exception)
        # 如果需要更详细的日志,可以通过日志API提取,不过一般exception足够审计用
    
    # 这里执行你的审计逻辑,比如写入审计表、发送告警等
    print(f"上游任务[{upstream_task_id}]状态: {upstream_state}")
    if error_msg:
        print(f"错误详情: {error_msg}")

这个方案适用于所有Operator,不管上游是PythonOperator还是BigQueryInsertJobOperator,只要任务失败时抛出了异常,就能通过upstream_ti.exception拿到错误信息。

方案2:结合XCom手动传递自定义状态与错误详情

如果需要更定制化的错误信息(比如除了异常栈,还要传递任务执行的参数、中间结果等),可以让上游任务主动把状态和错误推送到XCom,下游直接拉取。

针对PythonOperator的示例

在Python任务里捕获异常,主动推送状态和错误到XCom:

from airflow.decorators import task
from airflow.exceptions import AirflowException

@task
def upstream_python_task(**context):
    try:
        # 你的业务逻辑
        data = fetch_some_data()
        process_data(data)
        
        # 推送成功状态到XCom(默认key是'return_value')
        return {"status": "success", "processed_rows": len(data)}
    except Exception as e:
        # 推送失败状态和错误详情到XCom
        context['ti'].xcom_push(
            key="task_result",
            value={"status": "failed", "error": str(e), "failed_at": context['execution_date'].isoformat()}
        )
        # 必须抛出异常,让Airflow标记任务为失败
        raise AirflowException(f"任务执行失败: {str(e)}") from e

下游审计任务拉取XCom:

@task(trigger_rule='all_done')
def audit_task(**context):
    ti = context['ti']
    upstream_task_id = "upstream_python_task"
    
    # 拉取上游推送的自定义结果
    task_result = ti.xcom_pull(task_ids=upstream_task_id, key="task_result")
    # 如果上游成功,return_value会是成功的结果,失败的话用task_result
    if not task_result:
        task_result = ti.xcom_pull(task_ids=upstream_task_id, key="return_value")
    
    # 执行审计逻辑
    print(f"上游任务审计结果: {task_result}")

针对BigQueryInsertJobOperator的示例

BigQueryInsertJobOperator失败时会抛出BigQueryJobFailedException,我们可以通过on_failure_callback推送自定义错误信息:

from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator
from airflow.utils.context import Context

def bq_task_failure_callback(context: Context):
    ti = context['ti']
    exception = context['exception']
    # 解析BigQuery的错误详情,比如作业ID、错误原因
    error_details = {
        "status": "failed",
        "job_id": exception.job_id,
        "error_msg": str(exception),
        "execution_date": context['execution_date'].isoformat()
    }
    ti.xcom_push(key="bq_task_result", value=error_details)

# 定义BigQuery上游任务
upstream_bq_task = BigQueryInsertJobOperator(
    task_id="upstream_bq_task",
    configuration={
        "query": {
            "query": "SELECT * FROM your_dataset.your_table",
            "useLegacySql": False
        }
    },
    on_failure_callback=bq_task_failure_callback,
    dag=dag
)

# 下游审计任务
@task(trigger_rule='all_done')
def audit_bq_task(**context):
    ti = context['ti']
    task_result = ti.xcom_pull(task_ids="upstream_bq_task", key="bq_task_result")
    
    # 如果任务成功,task_result会是None,此时可以查询任务状态
    if not task_result:
        upstream_ti = ti.get_task_instance(task_id="upstream_bq_task")
        task_result = {"status": upstream_ti.state}
    
    print(f"BigQuery任务审计结果: {task_result}")
最佳实践总结
  • 优先使用方案1,因为它不需要修改上游任务,代码简洁,适合大多数审计场景。
  • 如果需要自定义错误信息或者传递额外上下文,再用方案2。
  • 无论用哪种方案,下游任务必须设置trigger_rule='all_done',确保上游任务无论成功、失败还是被跳过,审计任务都会执行。
  • 对于敏感的错误信息,注意Airflow的XCom和日志是否有权限控制,避免泄露敏感数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 08:52:35