【问题标题】:Flink and Kinesis stream app for non continous data用于非连续数据的 Flink 和 Kinesis 流应用
【发布时间】:2022-01-10 10:53:26
【问题描述】:

我们已经构建了一个 Flink 应用程序来处理来自 Kinesis 流的数据。该应用程序的执行流程包含基于注册类型过滤数据、基于事件时间戳分配水印、映射、处理和聚合函数应用于5分钟数据窗口的基本操作,如下所示:

    final SingleOutputStreamOperator<Object> inputStream = env.addSource(consumer)
            .setParallelism(..)
            .filter(..)
            .assignTimestampsAndWatermarks(..);

    // Processing flow
    inputStream
            .map(..)
            .keyBy(..)
            .window(..)
            .sideOutputLateData(outputTag)
            .aggregate(aggregateFunction, processWindowFunction);

    // store processed data to external storage
    AsyncDataStream.unorderedWait(...);

我的水印分配器的参考代码:

    @Override
public void onEvent(@NonNull final MetricSegment metricSegment,
                    final long eventTimestamp,
                    @NonNull final WatermarkOutput watermarkOutput) {
    if (eventTimestamp > eventMaxTimestamp) {
        currentMaxTimestamp = Instant.now().toEpochMilli();
    }
    eventMaxTimestamp = Math.max(eventMaxTimestamp, eventTimestamp);
}

@Override
public void onPeriodicEmit(@NonNull final WatermarkOutput watermarkOutput) {
    final Instant maxEventTimestamp = Instant.ofEpochMilli(eventMaxTimestamp);
    final Duration timeElaspsed = Duration.between(Instant.ofEpochMilli(lastCurrentTimestamp), Instant.now());
    if (timeElaspsed.getSeconds() >= emitWatermarkIntervalSec) {
        final long watermarkTimestamp = maxEventTimestamp.plus(1, ChronoUnit.MINUTES).toEpochMilli();
        watermarkOutput.emitWatermark(new Watermark(watermarkTimestamp));
    }
}

现在,这个应用程序在某个时候以良好的性能运行(就延迟而言,大约为几秒)。但是,最近上游系统帖子发生了变化,Kinesis 流中的数据以突发方式发布到流中(每天仅 2-3 小时)。发布此更改后,我们看到我们的应用程序的延迟出现了巨大的峰值(使用 flink gauge 方法测量,方法是在第一个过滤方法中记录开始时间,然后在 Async 方法中通过计算时间戳中的差异来发出度量标准开始时间图)。想知道在使用带有 Kinesis 流的 Flink 应用程序处理突发流量/非连续数据流时是否存在任何问题?

【问题讨论】:

    标签: apache-flink flink-streaming amazon-kinesis


    【解决方案1】:

    由于输入流现在长时间处于空闲状态,这可能会造成水印被搁置的情况。如果是这种情况,那么我预计延迟会出现很大差异,因为它(可能)只是每个突发的最终窗口,其结果会延迟到下一个突发到达。

    【讨论】:

    • 啊,我明白了,但在我的应用程序中可能不是这种情况,因为我在 onPeriodicEmit 函数中发出水印,基于从最大事件时间的第一个事件开始经过 20 秒。所以即使对于最后的窗口,它会在第一个事件的 20 秒后发出水印,并仅在下一次爆发出现之前保持该水印。如果那里看起来有问题,请在我的问题中添加代码参考以供参考。
    • 检测突发并推进水印以清除结果是有道理的 - 但我无法确定您共享的代码是否会正常工作,因为对未初始化的变量的引用太多或未在此 sn-p 中使用。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-10-01
    • 2021-06-20
    • 2021-11-01
    • 1970-01-01
    • 2020-06-25
    • 1970-01-01
    相关资源
    最近更新 更多