You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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逐条读取:

  1. {"uuid":10150,"event_time":"2023-08-07T08:00:00.004071876Z"}:正常通过,控制台显示
  2. {"uuid":10151,"event_time":"2023-08-07T09:00:00.004071876Z"}:正常通过,控制台显示
  3. {"uuid":10152,"event_time":"2023-08-07T10:00:00.004071876Z"}:正常通过,控制台显示
  4. {"uuid":10150,"event_time":"2023-08-07T11:00:00.004071876Z"}:被判定为重复丢弃(此时Watermark阈值应为2023-08-07T09:59:50,第一条同uuid记录的event_time早于该阈值,理论上状态已清理)
  5. {"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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 03:35:55