【问题标题】:LeftOuterJoin in Flink (JAVA API)Flink 中的 LeftOuterJoin (JAVA API)
【发布时间】: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


    【解决方案1】:

    这是左外连接的预期行为。

    鉴于您的程序,左外连接在两种情况下调用JoinFunction

    1. 如果两个输入 input1input2 具有具有相同连接键的记录,则为该键的笛卡尔积的每个元素调用 join()
    2. 如果左侧输入 input1 的记录具有右侧输入中不存在的键 (input2),则使用 input1null 的键为每个记录调用 join()正确的输入。

    您应该将in2 == null 的检查添加到您的JoinFunction

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2023-04-04
    • 1970-01-01
    • 1970-01-01
    • 2020-07-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-07-25
    相关资源
    最近更新 更多