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

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_LISTENERS includes an address your Spark app can access. For example:
    environment:
      ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092,PLAINTEXT://kafka:9092
    
    (Use localhost:9092 if Spark runs outside Docker, or kafka:9092 if it's in the same Docker network.)
  • In your Spark code, make sure bootstrap.servers matches 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 topic1 exists:
    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 08:26:16