【问题标题】:Apache Flink: How are events partitioned for a keyed CoFlatMapFunction?Apache Flink:如何为键控 CoFlatMapFunction 划分事件?
【发布时间】:2018-06-02 19:02:23
【问题描述】:

这是一个关于连接键控流的非常基本的问题。

如果我有两个具有相关事件的流共享相同的逻辑键,并且这些流正在连接(使用键逻辑连接)并且这一切都以并行度 > 1 运行,那么 Flink 如何保证来自不同的两个事件具有相同逻辑键的流最终在同一个并行运算符实例中?

这是一个关于医院患者流的虚构示例 - 温度流和心跳流。我们希望使用ConnectedStreamCoFlatMapFunction,通过患者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 事件被发送到那里。除非TemperatureHeartBeat 事件都发送到同一个并行实例运算符,否则连接将不起作用。

两个流使用相同的逻辑键(即患者的 ID)这一事实是特定于应用程序的,Flink 不知道。这两个连接的流可能使用它们自己的彼此无关的密钥。

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    当然,键的选择是特定于应用程序的。但是,Flink 知道如何访问键,因为您提供了键选择器功能(pt -&gt; pt.getPatientId()hbt -&gt; hbt.getPatientId())。 Flink 确保两个流的键具有相同的类型,并在两个流上应用相同的哈希函数来确定将记录发送到哪里。

    因此,两个流的相同值被传送到同一个运算符实例。

    【讨论】:

    • 感谢您的回复。如果它们具有不同的键,它是否仍然有效 - 例如,如果一个流以患者 ID(整数)为键,而另一个流以患者的全名(字符串)为键。 Flink 是否能够将同一患者的记录从两个流发送到同一个操作员实例? temperatureStream .keyBy(pt -> pt.getPatientId()) .connect (heartBeatStream.keyBy(hbt -> hbt.getPatientFullName() ) .flatMap (new RichCoFlatMapFunction()
    • Flink 检查键的类型是否匹配。因此,一侧的整数键和另一侧的字符串键将被拒绝。除此之外,Flink 使用您指定的键。
    • @FabianHueske 所说的“两个流的相同值被传送到同一个操作员实例”你也指的是同一台物理机器吧?
    • 是的,一个算子实例在单台机器上运行。但是一台机器可以运行多个算子实例。
    猜你喜欢
    • 2018-11-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多