【问题标题】:Persist Apache Flink window持久化 Apache Flink 窗口
【发布时间】:2021-11-03 21:05:04
【问题描述】:

我正在尝试使用 Flink 在流式传输中使用来自消息队列的有界数据。数据将采用以下格式:

{"id":-1,"name":"Start"}
{"id":1,"name":"Foo 1"}
{"id":2,"name":"Foo 2"}
{"id":3,"name":"Foo 3"}
{"id":4,"name":"Foo 4"}
{"id":5,"name":"Foo 5"}
...
{"id":-2,"name":"End"}

可以使用事件 id 确定消息的开始和结束。我想接收此类批次并将最新的(通过覆盖)批次存储在磁盘或内存中。我可以编写一个自定义窗口触发器来使用开始和结束标志提取事件,如下所示:

DataStream<Foo> fooDataStream = ...
AllWindowedStream<Foo, GlobalWindow> fooWindow = fooDataStream.windowAll(GlobalWindows.create())
.trigger(new CustomTrigger<>())
.evictor(new Evictor<Foo, GlobalWindow>() {
    @Override
    public void evictBefore(Iterable<TimestampedValue<Foo>> elements, int size, GlobalWindow window, EvictorContext evictorContext) {
        for (Iterator<TimestampedValue<Foo>> iterator = elements.iterator();
             iterator.hasNext(); ) {
            TimestampedValue<Foo> foo = iterator.next();
            if (foo.getValue().getId() < 0) {
                iterator.remove();
            }
        }
    }

    @Override
    public void evictAfter(Iterable<TimestampedValue<Foo>> elements, int size, GlobalWindow window, EvictorContext evictorContext) {

    }
});

但是我怎样才能持久化最新窗口的输出。一种方法是使用ProcessAllWindowFunction 接收所有事件并将它们手动写入磁盘,但感觉就像是黑客攻击。我也在研究带有 Flink CEP 模式的 Table API(例如 question),但找不到在每批之后清除表以丢弃前一批事件的方法。

【问题讨论】:

    标签: apache-flink flink-streaming flink-cep


    【解决方案1】:

    有几件事妨碍了你想要的东西:

    (1) Flink 的窗口操作符产生追加流,而不是更新流。它们并非旨在更新先前发出的结果。 CEP 也不产生更新流。

    (2) Flink 的文件系统抽象不支持覆盖文件。这是因为对象存储(如 S3)不能很好地支持此操作。

    我认为您的选择是:

    (1) 重做您的工作,使其生成更新(更改日志)流。您可以使用toChangelogStream 或使用创建更新流的表/SQL 操作来执行此操作,例如GROUP BY(在没有时间窗口的情况下使用时)。除此之外,您还需要选择一个支持撤回/更新的接收器,例如数据库。

    (2) 坚持生成追加流并使用FileSink 之类的东西将结果写入一系列滚动文件。然后在 Flink 之外编写一些脚本来得到你想要的。

    【讨论】:

    • 谢谢。除了在一个运算符中组合所有事件的性能瓶颈之外,您是否发现使用 ProcessAllWindowFunction 手动写入数据有任何问题?
    • 容错和恢复——这就是这种方法的问题所在。 Flink 能够提供恰好一次的保证,因为它的接收器以精心设计的方式参与检查点。你会放弃的。
    • 谢谢。这是有道理的。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多