JoinWindows.of(Duration.ofMinutes(5))如何正确配置WindowBytesStoreSupplier?
问题原因
JoinWindows.of(Duration.ofMinutes(5))默认会将窗口的前后匹配偏移都设为5分钟,因此实际join窗口的总大小 = 前偏移量 + 后偏移量 = 5min + 5min = 10分钟,你之前传入的windowSize参数为5分钟,与要求不匹配。- 窗口存储的保留周期必须和JoinWindows的总保留周期完全一致,JoinWindows总保留周期 = 窗口总大小 + 优雅关闭时间(graceMs)。从异常日志可知当前JoinWindows总保留为
300000 + 300000 + 85800000 = 86400000毫秒(即24小时),你之前传入的86700000毫秒数值不符合要求。
正确配置代码
如果你没有修改JoinWindows的默认参数,直接使用如下配置即可:
fun mangelJoinStoreSupplier(): WindowBytesStoreSupplier = Stores.inMemoryWindowStore( MangelJoinStore.NAME, // 总保留周期和JoinWindows要求对齐:24小时 Duration.ofHours(24), // join窗口实际总大小:10分钟 Duration.ofMinutes(10), true )
如果你希望窗口总大小保持为5分钟,可以在定义JoinWindows时手动调整偏移量:
// 手动指定仅单侧偏移为5分钟,总窗口大小为5分钟 val joinWindows = JoinWindows.ofTimeDifferenceAndGrace(Duration.ofMinutes(5), Duration.ofMillis(85800000)) .after(Duration.ZERO)
此时StoreSupplier的windowSize参数就可以传入5分钟,保留周期仍保持24小时即可。
内容的提问来源于stack exchange,提问作者Andras Hatvani
相关产品推荐
相关产品推荐

