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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 15:45:33