【问题标题】:Apache flink understanding of watermark idleness and relation to Bounded duration and window durationApache flink 对 watermark 空闲的理解以及与 Bounded duration 和 window duration 的关系
【发布时间】:2021-09-18 15:35:48
【问题描述】:

我有一个配置了 Kafka 连接器的 Flink 管道。

我已使用以下方法将水印生成频率设置为 2 秒:

env.getConfig().setAutoWatermarkInterval(2000);

现在我的滚动窗口是 60 秒的流窗口,我们在其中进行一些聚合,并且我们根据我们的数据字段之一的时间戳进行基于事件时间的处理。

我没有为我的水印策略或我的信息流配置 allowedLateness。

final ConnectorConfig topicConfig = config.forTopic("mytopic");
final FlinkKafkaConsumer<MyPojo> myEvents = new FlinkKafkaConsumer<>(
        topicConfig.name(),
        AvroDeserializationSchema.forSpecific(MyPojo.class),
        topicConfig.forConsumer()
);
myEvents.setStartFromLatest();



myEvents.assignTimestampsAndWatermarks(
    WatermarkStrategy
        .<MyPojo>forBoundedOutOfOrderness(
                Duration.ofSeconds(30))
        .withIdleness(Duration.ofSeconds(120))
        .withTimestampAssigner((evt, timestamp) -> evt.event_timestamp_field));

Q.1 从我正在阅读的内容来看,我的时间 0-60 的窗口将在 90 秒后计算,30-90 将在 120 秒后计算,依此类推。然而,由于我们正在做翻转窗口,即没有重叠,我的猜测是没有 30-90 窗口,0-60 之后的下一个窗口是 60-120,在 150 秒标记处触发,对吗?

Q.2 如果没有 allowedLateness,所有迟到的数据都将被丢弃,例如。 90 秒后到达的时间戳为 45 的事件被认为是乱序的,将超出第一个窗口,即 0-60。对于窗口 60-120,事件时间戳不匹配,因此将被丢弃且不包含在窗口在 150 秒处触发,对吗?

问题 3。对于源空闲持续时间,我选择 120 表示如果该主题的任何 Kakfa 分区没有数据,则在 2 分钟后将其标记为空闲,然后发送其他活动分区的水印。我的问题是关于这个数字的选择,即 2 分钟,以及它是否与窗口持续时间(60 秒)或无序(30 秒)有关。如果是这样,我应该在这里记住什么来进行恰当的选择,这样我就不会因为空闲分区导致的非高级水印而导致数据延迟?

或者 120 等待的时间太长,我可能会丢失数据,因此我应该将其设置为远小于 OutOfOrderness 持续时间以确保 0 数据丢失?

编辑:添加更多代码

【问题讨论】:

    标签: java apache-kafka apache-flink flink-streaming


    【解决方案1】:

    Q1:是的,没错。

    Q2:是的,这也是正确的。

    Q3:此处的详细信息取决于您是否让 Kafka 源应用 WatermarkStrategy,在这种情况下,它将执行每个分区的水印,或者 WatermarkStrategy 是否在之后某处部署为单独的操作符(通常在之后立即链接)源运算符。

    在第一种情况下(由FlinkKafkaConsumer 完成每个分区的水印),您将执行以下操作:

    FlinkKafkaConsumer<MyType> kafkaSource = new FlinkKafkaConsumer<>(...);
    
    kafkaSource.assignTimestampsAndWatermarks(WatermarkStrategy ...);
    
    DataStream<MyType> stream = env.addSource(kafkaSource);
    

    而在源代码之后单独进行水印,看起来像这样:

    DataStream<MyType> events = env.addSource(...);
    
    DataStream<MyType> timestampedEvents = events
      .assignTimestampsAndWatermarks(
          WatermarkStrategy
            .<MyType>forBoundedOutOfOrderness(Duration ...)
            .withTimestampAssigner((event, timestamp) -> event.timestamp));
    

    当基于每个分区完成水印时,单个空闲分区将为处理该分区的消费者/源实例保留水印 - 直到空闲超时开始(在您的示例中为 120 秒)。相比之下,如果水印是在链接到源的单独运算符中完成的,那么只有分配给该源实例的所有分区(具有空闲分区的分区)都空闲时,才会保留水印(同样,对于 120秒)。

    但不管这些细节如何,希望不会有数据丢失。会有一段时间不触发窗口(因为水印没有前进),但会继续处理事件并分配给它们相应的窗口。水印恢复后,这些窗口将关闭并提供结果。

    会发生数据丢失的情况是分区空闲,因为上游的某些故障导致了中断,最终产生了一堆迟到的事件。空闲超时到期后,水印将前进,如果源空闲是因为上游的某些东西被破坏(而不是因为根本没有事件),最终到达的那些事件将迟到(除非你有界-out-of-有序延迟足够大以容纳它们)。如果您选择忽略迟到的事件,那么这些事件将会丢失。

    【讨论】:

    • 太棒了!万分感谢。我更新了一些代码。只是为了确认最后一件事,我相信代码正在使用 FlinkKafkaConsumer,因此这属于第一个水印选项的类别,即每个分区水印?所以这就是在这里适用的,对吗? “当基于每个分区完成水印时,单个空闲分区将为处理该分区的消费者/源实例保留水印 - 直到空闲超时开始(在您的示例中为 120 秒)。”
    • 我已经扩展了我的答案以更明确地说明这一点。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-08-06
    • 1970-01-01
    • 2022-06-10
    • 2015-03-23
    • 2016-07-20
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多