【发布时间】:2019-07-27 22:58:40
【问题描述】:
我很难理解 flink 窗口原理,如果您能指出正确的方向,我会非常高兴。
我的目的是计算一个时间间隔内重复事件的数量,如果重复事件的数量大于阈值,则生成警报事件。
据我了解,窗口化非常适合这种情况。
如果窗口中的重复事件计数为 2,则附加要求是生成早期警报(即应在不等待窗口结束的情况下生成警报)。
我认为警报事件生成过程窗口函数可用于聚合窗口事件,自定义触发器可用于根据重复事件计数从窗口发出早期结果(在水印到达窗口的结束时间戳之前) .
我正在使用事件时间语义并且对自定义触发器有问题/疑问。
您可以在 gist 中找到实际的实现:https://gist.github.com/simpleusr/7c56d4384f6fc9f0a61860a680bb5f36
我正在使用键控状态来跟踪窗口 encounteredElementsCountState 中的元素计数
收到第一个元素后,我将EventTimeTimer 注册到窗口末端。这应该会触发 FIRE_AND_PURGE 以关闭窗口并按预期工作。
如果计数超过阈值,我会尝试触发早期触发。这似乎也成功了,processwindow 函数在此触发后立即调用。
问题是,我不得不在不了解原因的情况下在代码中插入下面的检查。因为之前收集的元素再次提供给onElement方法:
if (ctx.getCurrentWatermark() < 0) {
logger.debug(String.format("onElement processing skipped for eventId : %s for watermark: %s ", element.getEventId(), ctx.getCurrentWatermark()));
return TriggerResult.CONTINUE;
}
我无法弄清楚原因。我看到的是,当这种情况发生时,水印值为(ctx.getCurrentWatermark()) Long.MIN_VALUE(导致上述检查)。怎么会这样?
此检查似乎避免了重复的早期事件生成,但我不知道为什么会发生这种情况以及这种解决方法是否合适。
您能否告知为什么相同的元素在窗口中被处理两次?
另一个问题是关于键控状态的使用。这个实现在处理窗口后是否会泄漏任何状态?我正在尝试以触发器的 clear 方法清除所有已使用的状态,但这是否足够?
问候。
【问题讨论】:
标签: apache-flink flink-streaming