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

如何自定义Airflow日志格式 将RUN_ID和alert_name加入日志输出

Airflow自定义带RUN_ID和alert_name的日志实现方案

要实现按RUN_ID、alert_name快速检索的自定义日志格式,按以下3步配置即可:

1. 编写自定义日志过滤器

Airflow任务运行时的RUN_ID存在任务上下文中,可通过日志过滤器自动提取注入所有日志记录,无需手动传参。
在$AIRFLOW_HOME/config目录下新建log_config.py,写入以下内容:

import os
import logging
from airflow.models import TaskInstance

class TaskContextFilter(logging.Filter):
    """自动从当前任务上下文提取run_id注入日志"""
    def filter(self, record):
        record.run_id = ""
        try:
            ti = TaskInstance._get_current_context()
            if ti:
                record.run_id = ti.run_id
        except Exception:
            # 非任务运行场景(如scheduler、webserver进程日志)run_id留空即可
            pass
        return True

class AlertContextFilter(logging.Filter):
    """初始化alert_name字段,避免无该字段时日志格式化报错"""
    def filter(self, record):
        record.alert_name = getattr(record, "alert_name", "")
        return True

2. 覆写Airflow日志配置

在同一个log_config.py中追加配置,自定义日志格式、处理器,将RUN_ID、alert_name加入日志输出模板,同时适配本地日志和S3日志存储:

LOG_LEVEL = logging.INFO
# 日志格式和预期结构对齐
LOG_FORMAT = (
    "[%(asctime)s] "
    "[%(run_id)s] "
    "%(alert_name)s"
    "{%(filename)s:%(lineno)d} %(levelname)s - %(message)s"
)
LOG_DATEFORMAT = "%Y-%m-%d, %H:%M:%S UTC"

DEFAULT_LOGGING_CONFIG = {
    "version": 1,
    "disable_existing_loggers": False,
    "filters": {
        "task_context": {"()": TaskContextFilter},
        "alert_context": {"()": AlertContextFilter},
    },
    "formatters": {
        "airflow_custom": {
            "format": LOG_FORMAT,
            "datefmt": LOG_DATEFORMAT,
        }
    },
    "handlers": {
        "console": {
            "class": "logging.StreamHandler",
            "formatter": "airflow_custom",
            "filters": ["task_context", "alert_context"],
            "stream": "ext://sys.stdout",
        },
        "task_file": {
            "class": "airflow.utils.log.file_task_handler.FileTaskHandler",
            "formatter": "airflow_custom",
            "filters": ["task_context", "alert_context"],
            "base_log_folder": os.path.expandvars("$AIRFLOW_HOME/logs"),
            "filename_template": "{{ ti.dag_id }}/{{ ti.task_id }}/{{ ts }}/{{ try_number }}.log",
        },
        "s3_task": {
            "class": "airflow.providers.amazon.aws.log.s3_task_handler.S3TaskHandler",
            "formatter": "airflow_custom",
            "filters": ["task_context", "alert_context"],
            "base_log_folder": os.path.expandvars("$AIRFLOW_HOME/logs"),
            # 替换为实际的S3日志存储路径
            "s3_log_folder": "s3://your-airflow-log-bucket/path",
            "filename_template": "{{ ti.dag_id }}/{{ ti.task_id }}/{{ ts }}/{{ try_number }}.log",
        }
    },
    "loggers": {
        "airflow.task": {
            "handlers": ["console", "task_file", "s3_task"],
            "level": LOG_LEVEL,
            "propagate": False,
        }
    }
}

修改Airflow主配置文件airflow.cfg的[logging]段,指定使用自定义日志配置:

[logging]
logging_config_class = log_config.DEFAULT_LOGGING_CONFIG

3. 业务代码中动态注入alert_name

alert_name是任务运行时循环处理告警的动态变量,无法全局自动注入,使用LoggerAdapter给对应代码块的日志注入该字段即可,修改DAG代码如下:

from airflow import DAG
from datetime import datetime, timedelta
from airflow.operators.python import PythonOperator
import logging

log = logging.getLogger(__name__)

default_args = {
    'owner': 'SRE',
    'execution_timeout': timedelta(minutes=150)
}

dag = DAG(
    dag_id = 'new_dag',
    default_args = default_args,
    start_date = datetime(2021, 11, 22),
    schedule_interval = timedelta(days=1),
    catchup = False,
    max_active_runs = 3,
)

def implement_alert_logic(alert_name, alert_log):
    alert_log.info(f'In the implementation for {alert_name}')

def myfunc(**wargs):
    for alert in ['alert_1', 'alert_2', 'alert_3']:
        # 给当前告警对应的日志注入alert_name,输出时自动带[alert_x]前缀
        alert_log = logging.LoggerAdapter(log, {"alert_name": f"[{alert}]"})
        alert_log.info(f'Executing logic for {alert}')
        implement_alert_logic(alert, alert_log)


t1 = PythonOperator(
    task_id='testing_this',
    python_callable = myfunc,
    dag=dag)

t2 = PythonOperator(
    task_id='testing_this2',
    python_callable = myfunc,
    dag=dag)

t1 >> t2

配置生效验证

重启Airflow scheduler、worker组件后运行DAG,输出的日志格式完全符合预期,推送到S3的日志也会保留该格式,直接通过RUN_ID或alert_name关键词即可快速检索:

[2022-06-13, 08:16:54 UTC] [manual__2022-06-13T08:16:54.103265+00:00] {myenv.py:32} INFO - Executing logic for alert_1
[2022-06-13, 08:16:54 UTC] [manual__2022-06-13T08:16:54.103265+00:00] [alert_1]{myenv.py:27} INFO - In the implementation for alert_1
[2022-06-13, 08:16:54 UTC] [manual__2022-06-13T08:16:54.103265+00:00] {myenv.py:32} INFO - Executing logic for alert_2
[2022-06-13, 08:16:54 UTC] [manual__2022-06-13T08:16:54.103265+00:00] [alert_2]{myenv.py:27} INFO - In the implementation for alert_2

注意:Airflow 2.x版本的PythonOperator无需手动设置provide_context=True,上下文会自动注入。如果使用CeleryExecutor,需要保证所有worker节点上都同步了log_config.py配置文件,否则worker上的日志格式不生效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 22:42:23