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

如何将Airflow与Spark集群连接?遇MetricsSystem运行异常报错

问题描述

我尝试通过Airflow任务执行本地Spark作业,使用的代码如下:

spark = (SparkSession
            .builder
            .master("spark://172.22.102.229:7077")
            .appName("Test")
            .getOrCreate())

运行时出现如下错误:

py4j.protocol.Py4JJavaError: An error occurred while calling None.org.apache.spark.api.java.JavaSparkContext.
: java.lang.IllegalArgumentException: requirement failed: Can only call getServletHandlers on a running MetricsSystem

将.master("spark://172.22.102.229:7077")替换为.master("local")时,任务可正常运行。我的Spark集群已部署完成,且能通过http://172.22.102.229:4040/访问其Web UI。

排查与解决方案

1. 验证网络连通性

  • 确认Airflow Worker节点能ping通Spark Master节点(172.22.102.229),且7077端口可正常访问,用telnet 172.22.102.229 7077或nc -zv 172.22.102.229 7077测试。
  • 若Airflow基于Docker部署,需确保容器网络能接入Spark集群所在网络,必要时切换为host模式或配置端口映射。

2. 调整SparkSession配置

添加MetricsSystem相关配置,强制初始化该系统,修改后的代码如下:

spark = (SparkSession
            .builder
            .master("spark://172.22.102.229:7077")
            .appName("Test")
            .config("spark.metrics.conf.driver.source.jvm.class", "org.apache.spark.metrics.source.JvmSource")
            .config("spark.metrics.conf.executor.source.jvm.class", "org.apache.spark.metrics.source.JvmSource")
            .getOrCreate())

3. 检查Spark集群MetricsSystem状态

  • 登录Spark Master节点,查看$SPARK_HOME/logs下的日志,确认MetricsSystem是否正常启动,有无相关报错。
  • 检查spark-defaults.conf,确保没有spark.metrics.enabled=false这类禁用配置,若有则删除或改为true。

4. 直接测试Spark集群可用性

在Airflow Worker节点上用spark-submit提交测试作业:

spark-submit --master spark://172.22.102.229:7077 --class org.apache.spark.examples.SparkPi $SPARK_HOME/examples/jars/spark-examples_2.12-3.3.0.jar 10

如果该命令执行成功,说明集群本身无问题,需排查Airflow任务的配置或环境;若失败,优先修复Spark集群的配置问题。

5. 统一Airflow与Spark的环境配置

  • 确保Airflow Worker的PYSPARK_PYTHON和PYSPARK_DRIVER_PYTHON环境变量指向与Spark集群兼容的Python版本,避免版本不匹配。
  • 若使用Airflow的SparkOperator,明确指定spark_home参数,确保任务能正确加载Spark依赖。

内容的提问来源于stack exchange,提问作者Марсель Абдуллин

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 08:55:22