Scala环境下Spark Consumer无法读取Kafka Producer消息求助
Hey there! Let's troubleshoot why your Spark Consumer isn't picking up data from topic1 even though your Kafka Producer is running fine. Since you're using the spotify/kafka Docker image and Scala, here are the most common issues and fixes to check:
1. Verify Kafka Broker Address Configuration
When running Kafka in Docker, the broker's advertised address is critical—Spark needs to reach the right endpoint.
- First, check your docker-compose Kafka config: ensure
ADVERTISED_LISTENERSincludes an address your Spark app can access. For example:
(Useenvironment: ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092,PLAINTEXT://kafka:9092localhost:9092if Spark runs outside Docker, orkafka:9092if it's in the same Docker network.) - In your Spark code, make sure
bootstrap.serversmatches this address. Here's a sample config snippet:val kafkaParams = Map[String, Object]( "bootstrap.servers" -> "localhost:9092", "key.deserializer" -> classOf[StringDeserializer], "value.deserializer" -> classOf[StringDeserializer], "group.id" -> "spark-consumer-group-1", "auto.offset.reset" -> "earliest" // Don't skip this! )
2. Fix the Auto Offset Reset Strategy
If your Consumer uses a new group ID, the default auto.offset.reset is latest—meaning it only reads messages sent after the Consumer starts. If your Producer already sent data before launching Spark, you'll see nothing.
- Force the Consumer to read from the start of the topic by setting
"auto.offset.reset" -> "earliest"in your Kafka params.
3. Validate the Topic and Messages Exist
Rule out Kafka-side issues first:
- Enter the Kafka container and list topics to confirm
topic1exists:docker exec -it <kafka-container-name> kafka-topics.sh --list --zookeeper zookeeper:2181 - Test with Kafka's built-in console consumer to check if messages are actually in
topic1:kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic topic1 --from-beginning
If this reads messages, the problem is definitely in your Spark code.
4. Fix Incomplete Spark Streaming Code
Your snippet cuts off, so here's a complete, working example for both old Spark Streaming and Structured Streaming:
Old Spark Streaming API
import org.apache.spark.streaming._ import org.apache.spark.streaming.kafka010._ import org.apache.kafka.common.serialization.StringDeserializer object SparkConsumer { def main(args: Array[String]) { val spark = org.apache.spark.sql.SparkSession .builder() .appName("KafkaSparkStreaming") .master("local[*]") .getOrCreate() val ssc = new StreamingContext(spark.sparkContext, Seconds(3)) val topic1 = "topic1" val kafkaParams = Map[String, Object]( "bootstrap.servers" -> "localhost:9092", "key.deserializer" -> classOf[StringDeserializer], "value.deserializer" -> classOf[StringDeserializer], "group.id" -> "spark-consumer-group", "auto.offset.reset" -> "earliest" ) val topics = Array(topic1) val stream = KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) ) // Print received messages to verify stream.foreachRDD { rdd => println(s"Received ${rdd.count()} new messages") rdd.foreach(record => println(s"Key: ${record.key()}, Value: ${record.value()}")) } ssc.start() ssc.awaitTermination() } }
Structured Streaming API (Recommended for Spark 2.0+)
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ object SparkStructuredConsumer { def main(args: Array[String]) { val spark = SparkSession .builder() .appName("KafkaStructuredStreaming") .master("local[*]") .getOrCreate() import spark.implicits._ val df = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "localhost:9092") .option("subscribe", "topic1") .option("startingOffsets", "earliest") .load() // Parse and print messages val query = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") .writeStream .outputMode("append") .format("console") .start() query.awaitTermination() } }
Make sure you're not forgetting to call ssc.start() / query.start() and awaitTermination()—these are easy to miss!
5. Check Dependency Version Compatibility
Mismatched Spark and Kafka versions can cause silent failures. For example:
- Spark 3.3.x works best with Kafka 2.4+
- Add the correct dependencies in your
build.sbt:// For old Streaming API libraryDependencies ++= Seq( "org.apache.spark" %% "spark-streaming" % "3.3.0" % Provided, "org.apache.spark" %% "spark-streaming-kafka-0-10" % "3.3.0" ) // For Structured Streaming libraryDependencies += "org.apache.spark" %% "spark-sql-kafka-0-10" % "3.3.0"
Start with the simplest checks (offset reset, broker address) first—those are the most common culprits!
内容的提问来源于stack exchange,提问作者Cassie

