Har*_*lar 0 apache-flink flink-streaming
我在读一个简单的JSON字符串作为输入和键控基于两个领域流A和B。但KeyBy产生了对不同值的相同键合流B,但为的特定组合A和B。
输入:
{
"A": "352580084349898",
"B": "1546559127",
"C": "A"
}
Run Code Online (Sandbox Code Playgroud)
这是我的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();
Run Code Online (Sandbox Code Playgroud)
GenericDataObject是一个POJO,具有三个字段A,B并且C。
这是不同field值的控制台输出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'}
Run Code Online (Sandbox Code Playgroud)
注意第1行和第2行。即使它们的B值不同,也将它们放在相同的键控流(5)中。我肯定在这里做错了什么,有人可以指出正确的方向吗?
首先,您没有做错任何事情。
为什么它们在同一个子任务中?
假设您有数千个密钥,并且Apache Flink无法为每个密钥创建数千个线程。因此,必须有另一种机制来确保一组密钥在一个线程中被单独处理。
因此,在Apache Flink中,每个子任务都有其自己的密钥组,具有相同密钥组索引的不同密钥将在同一子任务中处理。子任务通常使用个别键控状态处理几个键,以保持不同键的状态分开。
keyBy并不意味着将不同的键分配给不同的子任务(或分区),但是具有相同键的所有记录都将分配给相同的子任务。因此,只能通过对KeySelector实例进行编程来确定不同的密钥是否在同一组中。
有关更多详细信息,您可以在Apache Flink的官方网站上查看此文章。
| 归档时间: |
|
| 查看次数: |
329 次 |
| 最近记录: |