【发布时间】: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