PySpark Streaming对接Kafka仅读新增数据 如何读取存量历史记录
解决方案
问题根因
你当前使用的Spark Streaming Kafka Direct Stream配置中,auto.offset.reset参数设置为largest,该配置的规则为:当消费组没有已提交的有效偏移量时,直接从Kafka分区的最新消息位置开始消费,因此只会读取程序启动后新产生的增量记录,存量历史记录会被直接跳过。
调整步骤
- 修改Kafka消费参数
将代码中Kafka配置项的auto.offset.reset值从largest改为smallest:
kStream = KafkaUtils.createDirectStream(ssc, [topic],{"metadata.broker.list": brokers, 'group.id':'ozy-group', 'fetch.message.max.bytes':'15728640', 'auto.offset.reset':'smallest'})
参数说明:
smallest对应高版本Kafka消费者API的earliest:无有效偏移量时从分区最早的消息位置开始消费,可拉取到Topic中所有存量历史记录
- 清理已有消费组偏移量(必做)
如果你的消费组ozy-group之前已经运行过,Kafka中已经存储了该消费组的已提交偏移量,此时修改auto.offset.reset不会生效(该参数仅在无有效偏移量时才会触发),可选择以下任意一种方式处理:
- 直接修改
group.id为一个从未使用过的新值,例如ozy-group-v2,首次运行时无历史偏移量,会自动触发从最开头消费的逻辑 - 使用Kafka自带的命令行工具删除
ozy-group消费组的对应偏移量
验证效果
修改完成后启动程序,会先批量处理Topic中所有存量历史记录,处理完成后会继续监听并处理后续新增的增量记录,和你使用./kafka-console-consumer.sh加--from-beginning参数的消费效果一致。
内容的提问来源于stack exchange,提问作者Gandharv Suri
相关产品推荐
相关产品推荐

