Spark 3.4.0 Structured Streaming读取Kafka已创建Topic报错求助
问题:Spark 3.4.0读取Kafka 3.4.0创建的Topic报UnknownTopicOrPartitionException
我正在将应用从Spark 2.2升级到3.4.0,同步升级了Kafka客户端。由于Kafka 3.4.0无需依赖ZooKeeper,改用Admin API创建Topic,代码如下:
import org.apache.kafka.clients.admin.{AdminClientConfig, NewTopic, Admin} val properties = new Properties properties.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9094") val admin = Admin.create(properties) if (!admin.listTopics().names().get().contains(topicName)) { val newTopic = new NewTopic(topicName, 1, 1.toShort) val result = admin.createTopics(Collections.singleton(newTopic)) val future = result.values.get(topicName) future.get() }
Topic创建成功,但使用Spark Structured Streaming读取该Topic时抛出异常,读取代码如下:
val kafkaStream = spark .readStream .format("kafka") .option("kafka.bootstrap.servers", "localhost:9094") .option("kafka.max.partition.fetch.bytes", settings.kafka.maxRequestSize) .option("startingOffsets", settings.kafka.startingOffsets) .option("maxOffsetsPerTrigger", settings.kafka.maxOffsetsPerTrigger.getOrElse(1000000L)) .option("failOnDataLoss", "false") .option("subscribe", topicName) .load()
异常信息如下:
java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.UnknownTopicOrPartitionException: This server does not host this topic-partition. at java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357) at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1908) at org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:165) at org.apache.spark.sql.kafka010.ConsumerStrategy.retrieveAllPartitions(ConsumerStrategy.scala:66) at org.apache.spark.sql.kafka010.ConsumerStrategy.retrieveAllPartitions$(ConsumerStrategy.scala:65) at org.apache.spark.sql.kafka010.SubscribeStrategy.retrieveAllPartitions(ConsumerStrategy.scala:102) at org.apache.spark.sql.kafka010.SubscribeStrategy.assignedTopicPartitions(ConsumerStrategy.scala:113) at org.apache.spark.sql.kafka010.KafkaOffsetReaderAdmin.$anonfun$partitionsAssignedToAdmin$1(KafkaOffsetReaderAdmin.scala:499) at org.apache.spark.sql.kafka010.KafkaOffsetReaderAdmin.withRetries(KafkaOffsetReaderAdmin.scala:518) at org.apache.spark.sql.kafka010.KafkaOffsetReaderAdmin.partitionsAssignedToAdmin(KafkaOffsetReaderAdmin.scala:498) at org.apache.spark.sql.kafka010.KafkaOffsetReaderAdmin.fetchLatestOffsets(KafkaOffsetReaderAdmin.scala:297) at org.apache.spark.sql.kafka010.KafkaMicroBatchStream.$anonfun$getOrCreateInitialPartitionOffsets$1(KafkaMicroBatchStream.scala:251) at scala.Option.getOrElse(Option.scala:121) at org.apache.spark.sql.kafka010.KafkaMicroBatchStream.getOrCreateInitialPartitionOffsets(KafkaMicroBatchStream.scala:246) at org.apache.spark.sql.kafka010.KafkaMicroBatchStream.initialOffset(KafkaMicroBatchStream.scala:98) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$getStartOffset$2(MicroBatchExecution.scala:455) at scala.Option.getOrElse(Option.scala:121) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.getStartOffset(MicroBatchExecution.scala:455) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$constructNextBatch$4(MicroBatchExecution.scala:489) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:411) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:409) at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:67) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$constructNextBatch$2(MicroBatchExecution.scala:488) at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:234) at scala.collection.AbstractIterator.foreach(Iterator.scala:932) at scala.collection.AbstractIterable.foreach(Iterable.scala:54) at scala.collection.TraversableLike.map$(TraversableLike.scala:234) at scala.collection.AbstractTraversable.map(Traversable.scala:104) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$constructNextBatch$1(MicroBatchExecution.scala:477) at scala.runtime.java8.JFunction0$mcZ$sp.apply(JFunction0$mcZ$sp.java:12) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.withProgressLocked(MicroBatchExecution.scala:802) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.constructNextBatch(MicroBatchExecution.scala:473) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$2(MicroBatchExecution.scala:266) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:12) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:411) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:409) at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:67) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$1(MicroBatchExecution.scala:247) at org.apache.spark.sql.execution.streaming.ProcessingTimeExecutor.execute(TriggerExecutor.scala:67) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:237) at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$1(StreamExecution.scala:306) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:12) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:827) at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:284) at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:207)
解决方案
- 等待元数据同步:刚创建的Topic需要时间在Kafka集群节点间同步元数据,可在创建Topic后添加短暂延迟(比如1-2秒),或在读取前通过AdminClient确认Topic分区信息已存在。
- 对齐Kafka客户端版本:Spark 3.4.0官方依赖的kafka-clients版本为3.3.1,确保项目中kafka-clients版本与Spark依赖版本一致,避免版本不兼容导致元数据读取异常。
- 缩短元数据刷新间隔:在Spark读取配置中添加
kafka.metadata.max.age.ms参数,设置为较小值(如3000毫秒),让Spark更快获取最新的Topic元数据:.option("kafka.metadata.max.age.ms", "3000") - 检查Topic状态:用Kafka命令行工具确认Topic的分区和副本分配正常:
kafka-topics.sh --describe --topic <topicName> --bootstrap-server localhost:9094 - 添加重试机制:如果是同一应用内创建后立即读取,可在读取逻辑中添加重试逻辑,捕获
UnknownTopicOrPartitionException后重试几次,等待元数据同步完成。
内容的提问来源于stack exchange,提问作者Chandan Gawri
相关产品推荐
相关产品推荐

