如何自定义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
相关产品推荐
相关产品推荐

