【发布时间】:2021-12-13 12:52:55
【问题描述】:
来源:Kinesis 数据流
Sink:Elasticesearch
两者都使用 AWS 服务。
另外,在 AWS Kinesis 数据分析应用程序上运行我的 Flink 作业
我遇到了 flink 的窗口功能的问题。我的工作是这样的
DataStream<TrackingData> input = ...; // input from kinesis stream
input.keyBy(e -> e.getArea())
.window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
.reduce(new MyReduceFunction(), new MyProcessWindowFunction())
.addSink(<elasticsearch sink>);
private static class MyReduceFunction implements ReduceFunction<TrackingData> {
@Override
public TrackingData reduce(TrackingData trackingData, TrackingData t1) throws Exception {
trackingData.setVideoDuration(trackingData.getVideoDuration() + t1.getVideoDuration());
return trackingData;
}
}
private static class MyProcessWindowFunction extends ProcessWindowFunction<TrackingData, TrackingData, String, TimeWindow> {
public void process(String key,
Context context,
Iterable<TrackingData> in,
Collector<TrackingData> out) {
TrackingData trackingIn = in.iterator().next();
Long videoDuration =0l;
for (TrackingData t: in) {
videoDuration += t.getVideoDuration();
}
trackingIn.setVideoDuration(videoDuration);
out.collect(trackingIn);
}
}
示例事件:
{"area":"sessions","userId":4450,"date":"2021-12-03T11:00:00","videoDuration":5}
我在这里所做的是从 kinesis 流中获得大量这些事件我想为每 10 秒的窗口求和 videoDuration 然后我想将此 单个事件 存储到 elasticsearch .
在 Kinesis 中,每秒可以有 10,000 个事件。我不想在 elasticsearch 中存储所有 10,000 个事件,我只想每 10 秒存储一个事件。
问题是当我向这个作业发送一个事件时,它会快速处理这个事件并直接沉入弹性搜索,但我想实现:直到每 10 秒我希望事件 videoDuration 时间增加,10 秒后只有一个将事件存储在 elasticsearch 中。
我怎样才能做到这一点?
【问题讨论】:
标签: apache-flink flink-streaming