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

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)

  • 已做配置:
    1. 创建 /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,
                },
            }
        },
    )
    
    1. 在 /opt/airflow/airflow.cfg 中添加:logging_config_class = log_conf.LOGGING_CONFIG
    2. 重启了 Airflow 的 scheduler、UI 及 PostgreSQL 所在 Kubernetes Pod,日志显示自定义配置已导入,但目标 INFO 日志仍未屏蔽。
  • 疑问:
    1. 为何 spark_submit.py 的 INFO 级日志仍会出现在 Airflow 日志中?
    2. 上述 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的,但存在两个可优化的点:

  1. Airflow 2.6.x中,官方推荐使用airflow.utils.log.logging_config.deep_update而非pydantic.utils.deep_update,后者虽能运行,但并非Airflow官方指定的工具函数。
  2. 针对目标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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 06:25:09