Spark Streaming消费Kafka报ConcurrentModificationException问题求助
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()ands1.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=trueare unnecessary here (you only create one context pair), andspark.streaming.concurrentJobs=3might 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.concurrentJobsto1(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

