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

Spark Streaming消费Kafka报ConcurrentModificationException问题求助

Fixing KafkaConsumer ConcurrentModificationException in Spark Streaming

Hey there, let's dig into this issue you're facing. That ConcurrentModificationException is a classic sign that your KafkaConsumer is being accessed by multiple threads at the same time—and KafkaConsumer is explicitly not designed to handle concurrent thread access.

Root Cause Breakdown

Looking at your setup and code, here are the key factors triggering this error:

  • Single Kafka Partition: Your topic has num.partition=1, which means all incoming data is stuck in one partition. Spark Streaming's DirectStream creates a KafkaRDD with only one partition to match, so all processing tasks will compete for the same cached KafkaConsumer instance.
  • Dual Output Operations: You're calling both s1.print() and s1.saveAsTextFiles(...) on the same DStream. Each of these triggers a separate Spark job, and with only one Kafka partition, both jobs' tasks will try to use the same consumer concurrently.
  • Misconfigured Spark Settings: Flags like spark.driver.allowMultipleContexts=true are unnecessary here (you only create one context pair), and spark.streaming.concurrentJobs=3 might be allowing more concurrent job execution than your single-partition setup can handle safely.

Step-by-Step Fixes

1. Combine Output Operations to Avoid Concurrent Jobs

Instead of running separate jobs for printing and saving, merge these actions into a single foreachRDD call. This way, you process the RDD once and perform both actions in a single job, eliminating the consumer thread conflict.

2. Increase Kafka Topic Partition Count

A single partition is a critical bottleneck here. Boost your topic's partition count to match or exceed your Spark parallelism (since you're using local[4], start with at least 4 partitions). This distributes data across multiple partitions, each with its own KafkaConsumer instance, preventing thread clashes.

Run this command on your Kafka server (0.0.0.178) to adjust partitions:

kafka-topics.sh --alter --topic topics1 --bootstrap-server 0.0.0.178:9092 --partitions 4

3. Clean Up Spark Configuration

Remove unnecessary flags and align settings with best practices:

  • Delete spark.driver.allowMultipleContexts=true (you don't need multiple contexts here)
  • Set spark.streaming.concurrentJobs to 1 (the default) unless you truly need multiple independent streaming jobs running.

Modified Code Example

Here's your updated code incorporating all fixes:

// Create the context with a 3 second batch size
val sparkConf = new SparkConf()
  .setAppName("SparkScript")
  .set("spark.streaming.concurrentJobs", "1")
  .setMaster("local[4]")

val sc = new SparkContext(sparkConf)
val ssc = new StreamingContext(sc, Seconds(3))

// Case class definitions remain unchanged
case class Thema(name: String, metadata: String)
case class Tempo(unit: String, count: Int, metadata: String)
case class Spatio(unit: String, metadata: String)
case class Stt(spatial: Spatio, temporal: Tempo, thematic: Thema)
case class Location(latitude: Double, longitude: Double, name: String)
case class Datas1(location : Location, timestamp : String, windspeed : Double, direction: String, strenght : String)
case class Sensors1(sensor_name: String, start_date: String, end_date: String, data1: Datas1, stt: Stt)

val kafkaParams = Map[String, Object](
  "bootstrap.servers" -> "0.0.0.178:9092",
  "key.deserializer" -> classOf[StringDeserializer].getCanonicalName,
  "value.deserializer" -> classOf[StringDeserializer].getCanonicalName,
  "group.id" -> "test_luca",
  "auto.offset.reset" -> "earliest",
  "enable.auto.commit" -> (false: java.lang.Boolean)
)

val topics1 = Array("topics1")
val s1 = KafkaUtils.createDirectStream[String, String](
  ssc, 
  PreferConsistent, 
  Subscribe[String, String](topics1, kafkaParams)
).map(record => {
  implicit val formats = DefaultFormats
  parse(record.value).extract[Sensors1]
})

// Combine print and save operations into one job
s1.foreachRDD { rdd =>
  if (!rdd.isEmpty()) {
    // Print the data
    rdd.foreach(println)
    // Save to text files
    rdd.saveAsTextFiles("results/", "")
  }
}

ssc.start()
ssc.awaitTermination()

内容的提问来源于stack exchange,提问作者lu_ferra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:16:47