【问题标题】:Flink streaming: Do the events get distributed to each task slots separately according to their keys?Flink 流式处理:事件是否根据它们的键分别分配到每个任务槽?
【发布时间】:2021-12-16 07:37:11
【问题描述】:

因此,例如,如果我有按键 A 的事件和按键 B 的事件以及 2 的并行度。所有带有键 A 的事件都进入一个任务槽,而键 B 的事件进入另一个任务槽?

如果我只使用密钥 A 按顺序获取事件会发生什么。它们是否也被分配到两个任务槽。这是否意味着我失去了它们来的顺序?

【问题讨论】:

    标签: parallel-processing streaming apache-flink data-stream


    【解决方案1】:

    不,这不是它的工作原理。

    发生的情况是,每个键都映射到一个键组,其中键组的总数由集群的最大并行度(配置设置)决定。然后将键组映射到任务槽。如果有两个键和两个槽,则完全有可能将两个键分配到同一个槽。

    key的key组是:

    MathUtils.murmurHash(key.hashCode()) % maxParallelism
    

    一个密钥组的槽是:

    keyGroup * actualParallelism / maxParallelism
    

    关于维护排序,见https://stackoverflow.com/a/69094404/2000823https://stackoverflow.com/a/69757412/2000823

    【讨论】:

    • 嘿大卫!我希望可以进一步迭代这个问题。有两个不同的流包含相同类型的事件,每个流都单独键控。这意味着,即使两个流都被键控在某个相似的键 A 上,流 1 中的 A 将被映射到与流 2 中的 A 不同的键组,因为它们都是单独键控的,对吧?在我使用的情况下,这是有问题的,因为如果我尝试使用 joinFunction 加入两个流,它们不会加入,因为它们可能存在于单独的任务槽中。还有其他解决方法吗?
    • “来自流 1 的 A 将被映射到与来自流 2 的 A 不同的密钥组......”不,这是不正确的。集群中的所有参与者都将给定的密钥映射到同一个插槽。 Flink 的 join 依赖于此。
    • 这很奇怪,在我的实际工作中这是目前正在发生的事情。加入两个正在获得键控的流。在测试中,我基本上分别在每个流中发送具有相同键的相同事件,但是除非并行度 = 1,否则它们永远不会加入。
    • 打印后似乎两个事件都在同一个实例上。他们没有加入会不会是水印问题?
    • 请创建一个新问题,并提供足够的详细信息来重现该问题。另外,仅供参考,github.com/apache/flink-training/tree/master/rides-and-fares 中有一个示例,您可能会觉得有帮助。
    猜你喜欢
    • 1970-01-01
    • 2012-10-27
    • 1970-01-01
    • 1970-01-01
    • 2018-11-08
    • 2017-01-28
    • 2014-04-07
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多