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

Airflow中如何获取run_id、task_id、dag_id及尝试次数变量?

Airflow 失败回调获取日志所需参数方法

所有你需要的参数都可以直接从回调函数接收的context字典中获取,注意try_number的取值有个容易踩的坑:

  • dag_id:Airflow 2.0+版本可直接通过context['dag_id']获取,全版本兼容写法为context['dag'].dag_id
  • task_id:Airflow 2.0+版本可直接通过context['task_id']获取,全版本兼容写法为context['task'].task_id
  • try_number:不要直接取context['try_number']——任务失败触发回调时,Airflow已经提前将尝试次数+1用于判断是否需要重试,此时直接拿到的是下一次重试的序号,当前失败的实际尝试次数需要取context['ti'].try_number - 1,这个写法不管任务有没有配置重试都能正确拿到对应attempt值。

完整的日志路径拼接示例代码如下:

def send_failure_email_with_log(context):
    # 提取日志路径所需参数
    dag_id = context['dag_id']
    run_id = context['run_id']
    task_id = context['task_id']
    current_attempt = context['ti'].try_number - 1

    # 拼接日志文件路径
    log_file_path = f"dag_id={dag_id}/run_id={run_id}/task_id={task_id}/attempt={current_attempt}.log"

    # 后续读取日志、添加为邮件附件、发送通知的逻辑自行补充即可
    with open(log_file_path, 'r', encoding='utf-8') as log_f:
        log_content = log_f.read()
    # 邮件组装发送逻辑省略

如果你还在使用Airflow 1.x版本,context没有顶层的dag_id、task_id字段,统一用context['dag'].dag_id、context['task'].task_id的写法即可,ti是TaskInstance实例的上下文键,所有Airflow版本都通用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 15:12:32