【问题标题】:Spark-Java : How to add an array column in spark DataframeSpark-Java:如何在 spark Dataframe 中添加数组列
【发布时间】:2021-06-11 18:51:10
【问题描述】:

我正在尝试向我的 Spark 数据框添加一个新列。添加的新列的大小将基于一个变量(比如salt)帖子,我将使用该列来分解数据集以用于加盐连接。

目前,我在array 函数中使用连续的lit,但存在无法参数化的问题,并且在编码实践中看起来最糟糕。我当前的实现如下所示。

int salt =3;

Dataset<Row> Reference_with_Salt_Col = Reference.withColumn("salt_array", array(lit(0), lit(1), lit(2)));

我已经参考并查看了各种方法,但似乎没有一个可以解决 Java 中的问题。 functions.typedlit 方法虽然适用于 Python/Scala,但似乎不适用于 Java。进一步传递数组或列表也无助于 spark 给出相同的错误。

我使用的是 Spark 2.2.0 和 Java 1.8 版本

【问题讨论】:

  • 在调用 typedlit 之前尝试将 java 列表转换为 scala seq。见this
  • 嗨,mck,试过这个但没能实现。可能是由于某些 scala 版本不匹配。能够使用地图实现。尽管来自@blackbishop 的下面一个是直接的实现,似乎确实可以解决这个问题。

标签: java arrays list apache-spark apache-spark-sql


【解决方案1】:

您可以使用array 函数,但首先将列表的每个元素转换为lit。在流中使用map 函数的示例:

import org.apache.spark.sql.*;

import java.util.Arrays;
import java.util.List;

// example of input dataframe
Reference.show();

//+-----+
//|label|
//+-----+
//|    a|
//|    b|
//|    c|
//+-----+

List<Integer> salt_array = Arrays.asList(0, 1, 2);

Reference.withColumn(
        "salt_array",
        functions.array(salt_array.stream().map(functions::lit).toArray(Column[]::new))
).show();

//+-----+----------+
//|label|salt_array|
//+-----+----------+
//|    a| [0, 1, 2]|
//|    b| [0, 1, 2]|
//|    c| [0, 1, 2]|
//+-----+----------+

要生成一个包含从 0 到 salt - 1 的范围的数组,您可以像这样使用 IntStream.rangeClosed:

import java.util.stream.IntStream;

int salt = 3;

Dataset<Row> Reference_with_Salt_Col = Reference.withColumn(
        "salt_array",
        functions.array(IntStream.rangeClosed(0, salt - 1).mapToObj(functions::lit).toArray(Column[]::new))
);

【讨论】:

  • 谢谢。正是我想要的。能够使用 map 实现一个混乱的实现,但你的方式更干净。
猜你喜欢
  • 2015-12-23
  • 2016-08-27
  • 2018-11-13
  • 1970-01-01
  • 1970-01-01
  • 2021-09-21
相关资源
最近更新 更多