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
相关产品推荐
相关产品推荐

