【问题标题】:Execute flink sink after tumbling window翻滚窗口后执行 flink sink
【发布时间】: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


    【解决方案1】:

    我认为你误诊了问题。

    您编写的代码将在每个 10 秒长的窗口中为每个在窗口期间发生事件的不同键生成一个事件。 MyProcessWindowFunction 没有任何效果:由于窗口结果已经预先聚合,每个 Iterable 将只包含一个事件。

    我相信你想这样做:

    input.keyBy(e -> e.getArea())
                    .window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
                    .reduce(new MyReduceFunction())
                    .windowAll(TumblingProcessingTimeWindows.of(Time.seconds(10)))
                    .reduce(new MyReduceFunction())
                    .addSink(<elasticsearch sink>);
    

    你也可以这样做

    input.windowAll(TumblingProcessingTimeWindows.of(Time.seconds(10)))
                    .reduce(new MyReduceFunction())
                    .addSink(<elasticsearch sink>);
    

    但第一个版本会更快,因为它能够在计算 windowAll 中的全局总和之前并行计算每个键的窗口结果。

    FWIW,Table/SQL API 通常更适合这种类型的应用程序,并且应该产生比其中任何一个都更优化的管道。

    【讨论】:

    • 谢谢!大卫,第一个解决方案奏效了。但我不明白为什么我们需要 2 个 reduce 和 2 个窗口函数,你能详细说明一下吗?
    • 第一层窗口是单独计算每个键的总和(并行计算),而第二层缩减窗口是计算所有键的总和(不能并行完成) .
    • 谢谢你,大卫。关于这个问题,我还有一个问题:stackoverflow.com/questions/70411797/… 你介意回答这个问题吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-26
    • 1970-01-01
    相关资源
    最近更新 更多