Spark Structured Streaming withWatermark去重行为异常咨询
Spark Structured Streaming流去重异常问题
场景与代码
从Kafka读取数据(startingOffsets设为earliest),写入控制台做流去重测试:
- 设置10秒Watermark,基于
event_time字段 - 使用
dropDuplicates按uuid去重
核心代码:
result_df = parsed_df\ .withWatermark("event_time", "10 seconds")\ .dropDuplicates(["uuid"]) result_df\ .writeStream\ .option("checkpointLocation", "gs://test/checkpoint_ts/")\ .format("console")\ .start()
测试步骤与异常现象
逐条向Kafka发送记录,Spark逐条读取:
{"uuid":10150,"event_time":"2023-08-07T08:00:00.004071876Z"}:正常通过,控制台显示{"uuid":10151,"event_time":"2023-08-07T09:00:00.004071876Z"}:正常通过,控制台显示{"uuid":10152,"event_time":"2023-08-07T10:00:00.004071876Z"}:正常通过,控制台显示{"uuid":10150,"event_time":"2023-08-07T11:00:00.004071876Z"}:被判定为重复丢弃(此时Watermark阈值应为2023-08-07T09:59:50,第一条同uuid记录的event_time早于该阈值,理论上状态已清理){"uuid":10153,"event_time":"2023-08-07T06:00:00.004071876Z"}:被丢弃,符合预期(event_time早于Watermark阈值)
原因分析与解决方案
1. Checkpoint历史状态残留
如果之前使用过相同的checkpointLocation运行任务,未清理checkpoint目录,Spark会加载历史状态数据。uuid=10150的去重标记仍保存在状态中,导致新记录被误判为重复。
解决:停止流任务,删除gs://test/checkpoint_ts/目录下所有内容,重新启动任务。
2. 状态清理的微批次延迟
Spark的状态清理操作在微批次处理的末尾执行,若第4条记录与第3条记录被分配到同一个微批次,此时状态清理还未完成,uuid=10150的状态仍存在,导致新记录被拦截。
验证与解决:发送第4条记录前等待3-5秒,让Spark完成前一个微批次的状态清理后再发送。
3. 流去重状态逻辑的细节
Spark对于withWatermark + dropDuplicates([uuid])的状态管理逻辑是:
- 为每个uuid维护状态,记录该uuid所有事件的时间
- 仅当全局Watermark阈值超过该uuid的所有事件中的最大event_time时,才会清理该uuid的状态
- 第一条uuid=10150的event_time为08:00:00,处理第3条记录时Watermark阈值为09:59:50,确实满足清理条件,但上述两种情况会导致状态未及时清理。
4. 额外验证建议
- 确认Spark版本(建议使用3.0+,流去重状态逻辑更稳定)
- 通过Spark UI的Streaming页面查看状态大小变化,确认uuid=10150的状态是否被清理
内容的提问来源于stack exchange,提问作者steve
相关产品推荐
相关产品推荐

