Spark Streaming读取Kafka主题时已消费数据重复返回是什么原因?
问题根因
- checkpointLocation配置位置错误:Spark Structured Streaming中,checkpointLocation是**流写入(writeStream)**的专属配置,仅用于在流计算执行过程中持久化偏移量、状态等信息,你在readStream阶段配置的checkpointLocation不会被识别,也不会生效,无法帮你记录已经消费的Kafka偏移量。
- 偏移量从未被提交:你当前只是创建了流DataFrame并调用display展示,没有触发偏移量提交的逻辑。如果是原生Spark的普通展示、或者未配置持久化checkpoint的Databricks display,每次运行代码都会启动一个全新的流作业,加上你配置了
startingOffsets = 'earliest',作业每次启动都会从Kafka主题的最早偏移量开始读取,自然每次都返回全量的4条数据。
解决方案
- 移除readStream中的无效配置:删除
spark.readStream里的.option('checkpointLocation', checkpoint_location)配置,该参数在读阶段无效,避免混淆。 - 正确配置写入阶段的checkpoint:如果需要持久化消费偏移量,必须在writeStream阶段配置checkpointLocation,示例如下:
# 先修改extract_kafka_data方法,移除readStream里的checkpointLocation配置 stream_df = extract_kafka_data(kafka_config, topic_name, column_schema) # 写入时配置checkpoint,偏移量会自动持久化到指定路径 query = stream_df.writeStream \ .format("console") # 输出到控制台,也可以换成其他输出格式 .option("checkpointLocation", "/持久化存储的checkpoint路径") # 路径要保证任务有权限读写,且不要随意删除 .start() query.awaitTermination()
- 如果使用Databricks环境的display方法展示流数据,需要给display方法传入checkpointLocation参数才能持久化偏移量,写法如下:
df = extract_kafka_data(kafka_config, topic_name, column_schema) display(df, checkpointLocation="/持久化存储的checkpoint路径")
补充说明
startingOffsets = 'earliest'仅在对应checkpoint路径不存在、首次启动流作业时生效,只要checkpoint存在,后续启动都会优先从checkpoint中记录的偏移量继续消费,不会重复读历史数据。
另外你提供的join_kafka_streams_po_denorm方法存在笔误,调用kafka_ingest时第一个参数错误传了方法本身kafka_ingest,应该传入转换后的final_df,需要修正避免写入报错。
内容的提问来源于stack exchange,提问作者Metadata
相关产品推荐
相关产品推荐

