如何将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,提问作者Марсель Абдуллин
相关产品推荐
相关产品推荐

