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