【发布时间】:2018-06-02 19:02:23
【问题描述】:
这是一个关于连接键控流的非常基本的问题。
如果我有两个具有相关事件的流共享相同的逻辑键,并且这些流正在连接(使用键逻辑连接)并且这一切都以并行度 > 1 运行,那么 Flink 如何保证来自不同的两个事件具有相同逻辑键的流最终在同一个并行运算符实例中?
这是一个关于医院患者流的虚构示例 - 温度流和心跳流。我们希望使用ConnectedStream 和CoFlatMapFunction,通过患者ID 加入这两个流。
DataStream<PatientTemperature> temperatureStream = ..
DataStream<HeartbeatStream> heartbeatStream = ..
temperatureStream
.keyBy(pt -> pt.getPatientId())
.connect (heartBeatStream.keyBy(hbt -> hbt.getPatientId() )
.flatMap (new RichCoFlatMapFunction() {
ValueState<PatientTemperatureAndHeartBeat> state = ...
public void flatMap1(PatientTemperature value, Collector<PatientTemperatureAndHeartBeat> out) {
state.value().setTemperature(value);
}
public void flatMap2(PatentHeartbeat value, Collector<PatientTemperatureAndHeartBeat> out) {
PatientTemperatureAndHeartBeat temperatureAndHeartBeat = state.value()
temperatureAndHeartBeat.setHeartBeat(value)
out.collect(temperatureAndHeartBeat);
}
});
假设它以并行度 = 3 运行,操作员任务 A、B、C,并且它们都在不同的物理机器上运行。
Flink 将保证患者“JohnDoe”的所有Temperature 事件将最终在同一个并行运算符实例中结束。假设它最终出现在 Operator B 中。
但是,当 Flink 接收到“JohnDoe”的 HeartBeat 事件时,它如何知道将它们发送给操作员 B,患者的 Temperature 事件被发送到那里。除非Temperature 和HeartBeat 事件都发送到同一个并行实例运算符,否则连接将不起作用。
两个流使用相同的逻辑键(即患者的 ID)这一事实是特定于应用程序的,Flink 不知道。这两个连接的流可能使用它们自己的彼此无关的密钥。
【问题讨论】:
标签: apache-flink flink-streaming