【发布时间】:2020-03-26 13:10:47
【问题描述】:
我正在寻找合并多个 (>20) 个 Flink 流的最佳方法,这些流代表我们系统中不同的事件来源,它们都具有相同的类型。
List<DataStream<Event>> dataStreams = ...
每个对象都是一个 POJO(显然是抽象表示)
public class Event implements Serializable {
public String userId;
public long eventTimestamp;
public String eventData;
}
我最终希望以单个流结束
DataStream<Event> merged;
有不同的方法来管理它:join、coGroup、map/flatMap(使用CoGroup)和union。
我不确定它们中的哪一个会给我从原始流到合并流的最快事件吞吐量。
此外,是否有一个运算符可以同时用于所有流,还是我应该一次只调用每 2 个流?
我希望得到一个流,然后将是 keyedBy userId 字段,这有什么不同吗?
附带说明一下,下一步是按eventTimestamp 对每个userId 的事件(在每个window 中)进行“排序”,以获得每个userId 事件的时间顺序。
【问题讨论】:
标签: apache-flink flink-streaming stream-processing