Airflow中如何获取run_id、task_id、dag_id及尝试次数变量?
Airflow 失败回调获取日志所需参数方法
所有你需要的参数都可以直接从回调函数接收的context字典中获取,注意try_number的取值有个容易踩的坑:
dag_id:Airflow 2.0+版本可直接通过context['dag_id']获取,全版本兼容写法为context['dag'].dag_idtask_id:Airflow 2.0+版本可直接通过context['task_id']获取,全版本兼容写法为context['task'].task_idtry_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
相关产品推荐
相关产品推荐

