Spark Structured Streaming任务报错:Initial job未接收任何资源求助
运行命令
spark-submit --master spark://{SparkMasterIP}:7077 --deploy-mode cluster --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2, com.datastax.spark:spark-cassandra-connector_2.12:3.2.0, com.github.jnr:jnr-posix:3.1.15 --conf spark.dynamicAllocation.enabled=false --conf com.datastax.spark:spark.cassandra.connectiohost={SparkMasterIP==CassandraIP}, spark.sql.extensions=com.datastax.spark.connector.CassandraSparkExtensions test.py
源代码
from pyspark.sql import SparkSession from pyspark.sql.functions import * from pyspark.sql.types import * from pyspark.sql import SQLContext # Spark Bridge local to spark_master == Connect master spark = SparkSession.builder \ .master("spark://{SparkMasterIP}:7077") \ .appName("Spark_Streaming+kafka+cassandra") \ .config('spark.cassandra.connection.host', '{SparkMasterIP==CassandraIP}') \ .config('spark.cassandra.connection.port', '9042') \ .getOrCreate() # Parse Schema of json schema = StructType() \ .add("col1", StringType()) \ .add("col2", StringType()) \ .add("col3", StringType()) \ .add("col4", StringType()) \ .add("col5", StringType()) \ .add("col6", StringType()) \ .add("col7", StringType()) # Read Stream From {TOPIC} at BootStrap df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "{KAFKAIP}:9092") \ .option('startingOffsets','earliest') \ .option("subscribe", "{TOPIC}") \ .load() \ .select(from_json(col("value").cast("String"), schema).alias("parsed_value")) \ .select("parsed_value.*") df.printSchema() # write Stream at cassandra ds = df.writeStream \ .trigger(processingTime='15 seconds') \ .format("org.apache.spark.sql.cassandra") \ .option("checkpointLocation","./checkPoint") \ .options(table='{TABLE}',keyspace="{KEY}") \ .outputMode('append') \ .start() ds.awaitTermination()
错误信息
Initial job has not accepted any resources; check your cluster UI to ensure that workers are registered and have sufficient resources
部署架构与问题描述
- 架构:Kafka(DBIP) → 本地(DriverIP,无Spark环境,使用Python虚拟环境运行PySpark) → 写入Spark&Kafka&Cassandra(MasterIP),DBIP、DriverIP、MasterIP为不同IP
- 已检查Spark UI,显示Worker节点注册正常,资源状态无异常(截图显示Master节点运行正常,Worker已注册且有可用资源)
解决方案
1. 修正spark-submit命令的配置错误
- 拼写错误:
com.datastax.spark:spark.cassandra.connectiohost应为spark.cassandra.connection.host,多余前缀和拼写错误会导致Cassandra连接配置失效 - --conf参数格式错误:多个配置项需单独用
--conf指定,不能用逗号分隔,正确写法:spark-submit --master spark://{SparkMasterIP}:7077 \ --deploy-mode cluster \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2,com.datastax.spark:spark-cassandra-connector_2.12:3.2.0,com.github.jnr:jnr-posix:3.1.15 \ --conf spark.dynamicAllocation.enabled=false \ --conf spark.cassandra.connection.host={SparkMasterIP==CassandraIP} \ --conf spark.sql.extensions=com.datastax.spark.connector.CassandraSparkExtensions \ test.py - Checkpoint路径问题:cluster模式下Driver运行在集群节点,
./checkPoint是本地路径仅当前节点可访问,需改为集群共享存储路径(如HDFS路径hdfs://{MasterIP}:9000/checkpoint、NFS挂载路径等)
2. 调整部署模式匹配环境
本地无Spark环境时,cluster模式会将Driver调度到集群Worker节点,本地虚拟环境无法生效:
- 若要使用本地虚拟环境运行Driver,改为
--deploy-mode client,但需确保本地能访问Spark Master(7077端口)、Worker节点、Kafka(9092)和Cassandra(9042)的网络 - 若坚持使用
cluster模式,需确保集群所有Worker节点安装了与提交环境版本一致的Python,且Spark配置中指定正确的Python路径(spark.pyspark.python)
3. 资源与网络检查
- 显式指定资源:提交命令时添加资源参数,避免默认资源不足,例如:
--executor-cores 2 --executor-memory 4g --driver-memory 2g - 网络连通性:确认Spark Master与所有Worker节点双向网络连通(ping通+端口开放);Worker节点能访问Kafka(DBIP:9092)和Cassandra(MasterIP:9042)
4. 代码优化
- 移除代码中
SparkSession.builder的.master("spark://{SparkMasterIP}:7077")配置,命令行指定的master优先级更高,避免冲突 - 确认代码中
kafka.bootstrap.servers和spark.cassandra.connection.host的IP地址能被Executor所在Worker节点访问
内容的提问来源于stack exchange,提问作者hi-inbeom
相关产品推荐
相关产品推荐

