【发布时间】:2019-09-13 16:28:56
【问题描述】:
我正在读取一个简单的 JSON 字符串作为输入,并根据两个字段 A 和 B 键入流。但是 KeyBy 为 B 的不同值生成相同的键控流,但对于 A 和 B 的特定组合。
输入:
{
"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,包含三个字段 A、B 和 C。
这是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