Spark中dropDuplicatesWithinWatermark工作原理及性能疑问
关于Spark Structured Streaming中dropDuplicatesWithinWatermark的问题解析
一、这个算子的核心工作逻辑
dropDuplicatesWithinWatermark本质是基于事件时间+水印的有状态去重,和普通的dropDuplicates()完全不是一回事:
- 它会在集群中维护一份状态存储,记录每个去重键对应的最新事件时间
- 水印的作用是定期“过期”旧状态:当水印推进到某个时间点后,所有事件时间早于该点的状态项会被清理,避免状态无限膨胀
- 关键:这个算子的触发不只是看有没有新数据,水印推进到清理阈值时,Spark会自动触发单独的微批来做状态清理
二、为什么3批数据跑出6个微批?
你遇到的情况是每一批数据对应两个微批:
- 第一个微批:处理刚进来的新数据,更新状态存储,同时计算当前的水印值
- 第二个微批:当新数据的事件时间足够让水印推进到能清理旧状态的时刻,Spark会触发一个专门的微批来执行状态清理——这个微批不会处理新数据,只做状态的过期清理
- 3批数据就对应3次数据处理微批 + 3次状态清理微批,加起来就是6个
举个实际场景:假设你设的水印延迟是5分钟,每批数据的事件时间是0分、6分、12分。处理完0分的批后,水印可能还没到清理线;但处理完6分的批后,水印推进到1分,此时可以清理事件时间早于1分的状态,触发清理微批;处理12分的批后,水印推进到7分,又触发一次清理微批——如果你的MemoryStream数据是连续生成的,事件时间间隔刚好够触发水印清理,就会每批都带一个清理微批。
三、为什么性能会暴跌?
这个算子的性能开销主要来自三个方面:
- 状态存储的读写开销:不管是处理数据还是清理状态,每一个微批都要读写分布式状态存储,要是状态量比较大,磁盘/网络IO会拖慢整个作业
- 额外微批的调度开销:多出来的清理微批会增加集群的调度成本,尤其是数据量不大的场景,这种调度开销占比会非常高
- 状态清理的扫描开销:清理过期状态时需要扫描状态存储的内容,如果状态没有合理分区,扫描全量状态的过程会很耗时
四、可以怎么优化?
- 调大水印延迟:不要设太小的延迟,比如业务能接受1小时的数据延迟,就把水印延迟设为1小时,减少状态清理的频率
- 控制去重键的基数:尽量让去重键不要太细碎,避免状态存储无限膨胀
- 给状态分区:如果用的是RocksDB这类支持分区的状态存储,按去重键或事件时间分区,缩小状态扫描的范围
- 换无状态去重:如果你的场景不需要严格基于事件时间的去重,直接用普通的
dropDuplicates(),或者改成定期批处理去重,性能会好很多
内容的提问来源于stack exchange,提问作者Pavel Orekhov
相关产品推荐
相关产品推荐

