【问题标题】:Apache flink - Early firing window implementation issue - duplicate elements receivedApache flink - 早期触发窗口实现问题 - 收到重复的元素
【发布时间】: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


    【解决方案1】:

    每个任务都将 currentWatermark 初始化为 Long.MIN_VALUE,并且这仍然是 currentWatermark 的本地值,直到更大的水印从该任务的所有输入流中到达。希望知道这将帮助您更好地了解正在发生的事情。

    就其价值而言,使用 ProcessFunction 实现这种逻辑通常比使用 Window API 更直接。

    【讨论】:

    • 嗨,David,实际上 process 函数对我来说似乎是一个低级别的处理操作(即,在不泄漏状态的情况下维护窗口边界似乎更难)。
    • 我后来从“使用 Apache Flink 进行流处理”一书中发现触发器非常适合早期触发。 (这来自本书的第 6 章:“自定义触发器也可用于在水印到达窗口的结束时间戳之前从事件时间窗口计算和发出早期结果”)。至于您的评论,水印值 > Long.MIN_VALUE 已经到达窗口,我无法弄清楚为什么要重新评估它们。我不知道它是否相关,但 watermark = Long.MIN_VALUE 是我在发生这种重复时观察到的行为。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-05-25
    • 1970-01-01
    • 2020-08-30
    • 2021-04-09
    • 2023-04-10
    • 1970-01-01
    相关资源
    最近更新 更多