Spark Streaming读取Kafka时startingOffsets的earliest与latest区别及问题咨询
Spark Streaming Kafka
startingOffsets=latest 无数据读取问题排查 先确认Kafka分区的实时偏移量与消息写入情况
当startingOffsets设为latest时,Spark会从每个Kafka分区的当前最大偏移量的下一个位置开始读取。如果任务启动时,目标topic的所有分区已经没有未消费的新消息(偏移量已经到了末尾),那自然读不到数据。
可以用Kafka自带工具查看分区最新偏移量:kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list <你的broker地址> --topic <目标topic> --time -1同时要确认任务启动后,有没有新的消息写入该topic——如果没有新消息流入,就算配置正确也不会有数据输出。
检查消费组的已提交偏移量是否干扰配置
如果你的Spark任务指定了group.id,且这个消费组之前已经消费过该topic,Spark会优先使用消费组已提交的偏移量,而非startingOffsets(这个配置仅在消费组首次启动、无提交偏移量时生效)。
查看消费组当前的偏移量:kafka-consumer-groups.sh --bootstrap-server <你的broker地址> --describe --group <你的group.id>如果发现消费组的偏移量已经追到了分区末尾,那只有新消息写入时才会读到数据,和
latest配置的预期一致,但如果没有新消息,就会出现无数据的情况。验证
startingOffsets配置是否真正生效
确保配置没有被覆盖。比如在Structured Streaming中,正确的配置方式应该是:spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "<broker地址>") .option("subscribe", "<目标topic>") .option("startingOffsets", "latest") .load()检查代码中是否重复设置了偏移量参数,或者提交任务时的命令行参数是否覆盖了代码中的配置。
确认Kafka Topic分区状态正常
排查目标topic的分区是否存在且可用。如果分区被删除、处于离线状态,或者Spark无法连接到Kafka broker,也会出现无法读取数据的情况。
内容的提问来源于stack exchange,提问作者user3185460
相关产品推荐
相关产品推荐

