【发布时间】:2016-10-13 19:11:25
【问题描述】:
我正在尝试在 Flink 中进行 LeftOuterJoin。 我不会尝试自己实现 leftOuterJoin,因为它已经完成 在这里使用 CoGroupFunction:https://gist.github.com/mxm/c2e9c459a9d82c18d789
我正在尝试使用 FlatJoinFunction:
public static final class leftOuter implements FlatJoinFunction<Tuple3<String,String,String>, Tuple2<String,String>, Tuple2<String,String>>{
@Override
public void join(Tuple3<String, String, String> in1,
Tuple2<String, String> in2,
Collector<Tuple2<String, String>> out) throws Exception {
// TODO Auto-generated method stub
out.collect(new Tuple2<String,String>(in1.f0, in2.f1 == null ? "null" : in2.f1));
}
}
我把这个函数称为:
input1.leftOuterJoin(input2).where(0)
.equalTo(1)
.with(new leftOuter());
不幸的是,我在 out.collect 行中收到 NullPointerException。
提前感谢您的帮助!
【问题讨论】:
标签: java mapreduce apache-flink bigdata