【问题标题】:KeyBy is not creating different keyed streams for different keysKeyBy 不会为不同的键创建不同的键控流
【发布时间】:2019-09-13 16:28:56
【问题描述】:

我正在读取一个简单的 JSON 字符串作为输入,并根据两个字段 AB 键入流。但是 KeyBy 为 B 的不同值生成相同的键控流,但对于 AB 的特定组合。

输入:

{
    "A": "352580084349898",
    "B": "1546559127",
    "C": "A"
}

这是我的 Flink 代码的核心逻辑:

DataStream<GenericDataObject> genericDataObjectDataStream = inputStream
            .map(new MapFunction<String, GenericDataObject>() {
                @Override
                public GenericDataObject map(String s) throws Exception {
                    JSONObject jsonObject = new JSONObject(s);
                    GenericDataObject genericDataObject = new GenericDataObject();
                    genericDataObject.setA(jsonObject.getString("A"));
                    genericDataObject.setB(jsonObject.getString("B"));
                    genericDataObject.setC(jsonObject.getString("C"));
                    return genericDataObject;
                }
            });
DataStream<GenericDataObject> testStream = genericDataObjectDataStream
            .keyBy("A", "B")
            .map(new MapFunction<GenericDataObject, GenericDataObject>() {
                @Override
                public GenericDataObject map(GenericDataObject genericDataObject) throws Exception {
                    return genericDataObject;
                }
            });
testStream.print();

GenericDataObject 是一个 POJO,包含三个字段 ABC

这是B字段不同值的控制台输出。

5> GenericDataObject{A='352580084349898', B='1546559224', C='A'}
5> GenericDataObject{A='352580084349898', B='1546559127', C='A'}
4> GenericDataObject{A='352580084349898', B='1546559234', C='A'}
3> GenericDataObject{A='352580084349898', B='1546559254', C='A'}

请注意第 1 行和第 2 行。即使它们具有不同的 B 值,它们也被放入同一个键控流 (5)。我必须在这里做一些根本错误的事情,有人可以指出我正确的方向吗?

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    首先,你没有做错任何事。

    为什么他们在同一个子任务中?

    假设您有数千个键,Apache Flink 不可能为每个键创建数千个线程。因此,必须有另一种机制来确保一组键在一个线程中单独处理。

    因此,在 Apache Flink 中,每个子任务都有自己的键组,具有相同键组索引的不同键将在同一个子任务中处理。并且一个子任务通常会处理一些具有单独键控状态的键,以保持不同键的状态分开。

    keyBy 并不是说​​将不同的key分配给不同的子任务(或分区),而是将所有具有相同key的记录分配给同一个子任务。所以你只能通过编写一个 KeySelector 实例来决定不同的键是否在同一个组中。

    更多详情,可以查看Apache Flink官网的这篇文章。

    A Deep Dive into Rescalable State in Apache Flink

    【讨论】:

    • 感谢您的回复。如何控制操作员是否在同一个子任务中?我似乎了解了键和键组的概念,但似乎不太明白如何应用它们来解决我的问题。此外,奇怪的是,这个问题只出现在B 的一组特定值上,即15465592241546559127。对于所有其他值,我看到正在生成不同的键控流。
    • 视情况而定。由于 Flink 不知道它会收到多少个 key,所以它只是根据并行度设置了几个组。因此,您只能通过编写 Keyselector 实例来确定两个不同的键是否在同一个键组中。
    • 我明白了,当我增加并行度时,它似乎变得更好了。但是不应该根据键控流的定义将不同的键分配给不同的流吗?即使并行度为 1,有没有办法强制这种行为?
    • 不,keyBy 的合约不是不同的键分配给不同的子任务(或分区),而是所有具有相同键的记录都分配给相同的子任务。因此,子任务通常处理许多不同的键。 Flink 使用 keyed state 来保持不同 key 的状态分开。 @bupt_ljy,您想用 cmets 中的信息扩展您的答案吗?谢谢!
    • @FabianHueske 感谢您对 keyBy 合同的清晰解释!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-05-03
    • 2018-05-31
    • 2017-03-24
    • 1970-01-01
    • 1970-01-01
    • 2016-03-06
    • 1970-01-01
    相关资源
    最近更新 更多