【问题标题】:How to write dynamic join condition in spark Java API如何在 Spark Java API 中编写动态连接条件
【发布时间】:2019-04-23 23:56:18
【问题描述】:

我想使用 spark Java API 对数据集执行左外连接。如何编写动态条件以匹配连接条件中的多个列。

我有两个数据集对象。它们都有 2 列或更多列。我无法定义条件

将一列与另一列匹配的示例

dataSet = resultData.as("resultData").join(distinctData.as("distinctData"), resultData.col("A").equalTo(distinctData.col("B")), "leftouter").selectExpr(select.toString());

现在由于有多个列,我无法使用 Java API 定义动态表达式来匹配多个列。

【问题讨论】:

  • 您可能投了反对票,因为您没有提供任何关于您的数据是什么样子的信息,或者到目前为止您已经尝试过什么。如果您能提供这些信息,我很乐意为您提供帮助。
  • 编辑问题
  • 您是否收到错误消息?运行上面的代码会发生什么?
  • 例如问题中提到的我没有收到任何错误。问题是我想指定匹配多个列的条件,但我找不到任何定义相同的参考。

标签: java apache-spark


【解决方案1】:

未经测试的代码 - 但这会从列名列表动态生成连接条件

public Column makeJoinConditional(Dataset<Row> df1, Dataset<Row> df2, List<String> columnNames, Column c)  {

        if (c==null) {
            String  top = columnNames.get(0);
            columnNames.remove(0);
            Column first = df1.col(top).equalTo(df2.col(top));

            return makeJoinConditional(df1,df2, columnNames,first);

        } else {

            if (columnNames.size()==0) {
                return c;
            } else {
                String  top = columnNames.get(0);
                columnNames.remove(0);
                Column next = c.and( df1.col(top).equalTo(df2.col(top)) );
                return makeJoinConditional(df1,df2, columnNames,next);
            }
        }
    }

    public Dataset<Row> joinDataFrames(Dataset<Row> df1, Dataset<Row> df2, List<String> columns) {
        Column joinCols = makeJoinConditional(df1,df2,columns,null);
        return df1.join(df2,joinCols);
    }

【讨论】:

  • 是的,这将在列数固定时起作用。但是列号因每种情况而异,我不能每次都为此更改代码:)
  • 那么,你需要一个函数,给定一个列列表,可以生成条件语句吗?
  • 是的,和你说的有点像。但是 dataste.join() 接受 column、columnExpr 和 seq。我正在尝试查找将返回条件列/语句的 columnExpr 是什么。
  • 好的,更新答案以根据列名列表动态生成条件
猜你喜欢
  • 2020-10-02
  • 1970-01-01
  • 1970-01-01
  • 2010-12-23
  • 2020-07-25
  • 2017-05-07
  • 1970-01-01
  • 1970-01-01
  • 2013-12-29
相关资源
最近更新 更多