Airflow 2.6.1设置spark_submit模块日志级别为WARN不生效问题
问题
- 环境:Airflow 2.6.1 调度 Spark 3.4.1 任务,使用
SparkSubmitOperator,Spark 运行在 Cluster Mode,只能通过spark_submit.py轮询 Spark Driver Pod 状态。 - 日志问题:Airflow 日志中充斥大量 INFO 级条目,示例:
[2024-07-29, 06:47:33 UTC] {spark_submit.py:523} INFO - 24/07/29 08:47:33 INFO LoggingPodStatusWatcherImpl: Application status for spark-dc8c170895df4383be2c6933606ee764 (phase: Running)
- 已做配置:
- 创建
/opt/airflow/config/log_conf.py,内容:
from copy import deepcopy from pydantic.utils import deep_update from airflow.config_templates.airflow_local_settings import DEFAULT_LOGGING_CONFIG LOGGING_CONFIG = deep_update( deepcopy(DEFAULT_LOGGING_CONFIG), { "loggers": { "airflow.providers.apache.spark.operators.spark_submit": { "handlers": ["task"], "level": "WARNING", "propagate": True, }, } }, )- 在
/opt/airflow/airflow.cfg中添加:logging_config_class = log_conf.LOGGING_CONFIG - 重启了 Airflow 的 scheduler、UI 及 PostgreSQL 所在 Kubernetes Pod,日志显示自定义配置已导入,但目标 INFO 日志仍未屏蔽。
- 创建
- 疑问:
- 为何
spark_submit.py的 INFO 级日志仍会出现在 Airflow 日志中? - 上述
log_conf.py配置是否兼容 Airflow 2.6.1?
- 为何
解答
疑问1:为何INFO日志未被屏蔽
你看到的这些日志并非 spark_submit.py 模块自身输出的日志,而是该模块直接转发了Spark Driver Pod的日志内容。spark_submit.py 里的 LoggingPodStatusWatcherImpl 会把从K8s Pod获取到的Spark日志原样打印到Airflow任务日志中,这部分日志不受Airflow的Python日志级别控制——因为它不是Airflow模块生成的日志,只是被Airflow任务进程读取后输出的外部日志。
另外,你配置里的propagate: True会让日志继续向上传递到父logger,可能也会导致部分日志通过其他handler输出。
疑问2:配置是否兼容Airflow 2.6.1
你的log_conf.py配置在结构上是兼容Airflow 2.6.1的,但存在两个可优化的点:
- Airflow 2.6.x中,官方推荐使用
airflow.utils.log.logging_config.deep_update而非pydantic.utils.deep_update,后者虽能运行,但并非Airflow官方指定的工具函数。 - 针对目标logger,建议将
propagate设为False,避免日志向上传递到更上层的logger被输出。修改后的配置如下:
from copy import deepcopy from airflow.utils.log.logging_config import deep_update from airflow.config_templates.airflow_local_settings import DEFAULT_LOGGING_CONFIG LOGGING_CONFIG = deep_update( deepcopy(DEFAULT_LOGGING_CONFIG), { "loggers": { "airflow.providers.apache.spark.operators.spark_submit": { "handlers": ["task"], "level": "WARNING", "propagate": False, }, } }, )
不过要注意,这个优化只能屏蔽spark_submit.py模块自身生成的INFO日志,无法解决你当前遇到的Spark Pod日志转发问题。
如果要屏蔽Spark Pod的轮询状态日志,你需要修改SparkSubmitOperator的相关参数:在Operator初始化时设置get_logs=False,或者调整Spark Driver的日志级别(在Spark配置里设置log4j.logger.org.apache.spark.deploy.k8s.LoggingPodStatusWatcherImpl=WARN)。
内容的提问来源于stack exchange,提问作者Phil

