【问题标题】:Apache Flink: NullPointerException in DataSet API Outer JoinApache Flink:DataSet API 外部联接中的 NullPointerException
【发布时间】:2018-04-12 01:50:44
【问题描述】:

我正在尝试在 Flink 的 Dataset API 中实现以下简单查询。

select 
    t1_value1 
from  
    table1 
where  
    t1_suppkey not in ( 
        select  
            t2_suppkey
        from  
            table2
     )

所以我的想法是执行左外连接 (table1.leftOuterJoin(table2)...),然后删除我获得 t1_suppkey 和 t2_suppkey 值的所有行。

所以我这样尝试:

     output = table1
    .leftOuterJoin(table2).where("t1_suppkey").equalTo("t2_suppkey")
    .with((Table1 t1, Table2 t2) -> new Tuple2<>(t1.ps_suppkey, t2.s_suppkey))
    .returns(new TypeHint <Tuple2<Integer, Integer>>() {});

但是,如果我这样做,它总是会因“java.lang.NullPointerException”而失败,我不知道为什么。如果我使用普通联接而不是左外部联接,则代码可以工作,但这不是我想要的。

我需要以不同的方式实现 Left Join,还是有更简单的方法来重写 Dataset API 中的“not in”语句?

【问题讨论】:

    标签: java sql dataset apache-flink


    【解决方案1】:
    output = table1
    .leftOuterJoin(table2)
    .where("t1_suppkey").equalTo("t2_suppkey") 
    .with((Table1 t1, Table2 t2, Collector<Tuple2<Integer, Integer>> c) -> { 
    if(t2 == null) {
        c.collect(new Tuple2<>(t1.t1_suppkey, t1.t1_value1)); 
    } 
    else { 
        //Do nothing. 
    }})
    

    【讨论】:

      【解决方案2】:

      DataSet API 的外连接调用JoinFunction 也用于在内侧找不到连接记录的外记录。在这种情况下,the JoinFunction.join() method is called with null

      由于您使用的是 LEFT OUTER JOIN,第二个参数 Table2 t2 可以是 nullNullPointerException 是由 t2.s_suppkey 引起的。您需要检查t2 == null,如果t2 不为空,则仅访问它。

      您还可以使用具有Collector 参数的FlatJoinFunction 实现NOT IN 连接,并且仅在t2 == null 时发出t1

      另一种选择是使用 Flink 的批处理 SQL support,它支持您示例中的查询。

      【讨论】:

      • 谢谢,像this 一样解决了它,它可以工作。
      猜你喜欢
      • 2012-10-12
      • 1970-01-01
      • 1970-01-01
      • 2011-12-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-08-04
      • 1970-01-01
      相关资源
      最近更新 更多