增大Join窗口大小与设置grace period的Join输出差异
Kafka Streams 两种JoinWindows配置的Join输出结果差异对比
以下分析默认基于Kafka Streams流对流Join的事件时间处理场景:
配置本身的逻辑差异
- 第一种配置
JoinWindows.of(Duration.ofMillis(a)).grace(Duration.ofMillis(b))- 窗口的事件匹配阈值为
a:仅当两条待Join的事件时间戳差值的绝对值≤a时,才会被纳入匹配范围 grace为窗口关闭容忍时长:窗口匹配逻辑的时间边界到达后,还会额外保留b毫秒接收迟到数据,这段时间内到达的迟到事件只要满足时间差≤a的要求,依然可以完成Join,直到b毫秒后窗口才会彻底关闭,丢弃所有后续到达的相关迟到数据
- 窗口的事件匹配阈值为
- 第二种配置
JoinWindows.of(Duration.ofMillis(a + b))- 未显式配置
grace时默认值为0,窗口的事件匹配阈值直接为a+b,窗口到达时间边界后立刻关闭,不会接收任何迟到数据
- 未显式配置
输出结果的核心差异
- 匹配的事件范围不同:第二种配置允许两条事件的最大时间差为
a+b,明显宽于第一种的a。举个例子,当a=100、b=50时,时间差为120毫秒的事件对,仅第二种配置会输出Join结果,第一种永远不会匹配该类事件对 - 迟到数据的处理结果不同:第一种配置可以接收窗口结束后
b毫秒内的迟到事件,只要满足时间差要求就能生成Join结果;第二种配置窗口到点立刻关闭,哪怕事件只晚到1毫秒,也会被直接丢弃,不会生成对应的Join输出 - 输出时机不同:第一种配置的部分Join结果会因为等待迟到数据延迟输出,第二种配置的所有Join结果都会在窗口结束时立刻输出,无额外延迟
场景示例
假设a=100、b=50,水印已经推进到时间点t0+120:
- 现有右流事件到达,事件时间为
t0+80,和已存的左流事件t0的时间差为80毫秒≤100毫秒 - 第一种配置因为设置了50毫秒的grace期,允许事件时间≥
t0+120-50 = t0+70的事件参与计算,所以该事件会正常匹配,输出Join结果 - 第二种配置grace默认为0,仅允许事件时间≥
t0+120的事件参与计算,该事件会被直接丢弃,无对应输出
内容的提问来源于stack exchange,提问作者Evgeniy Berezovsky
相关产品推荐
相关产品推荐

