Java操作Apache Spark如何传入列名数组实现Dataset列拼接
Java API操作Spark传入列名数组实现多列拼接方案
问题说明
现有存储待拼接列名的字符串数组定义如下:
String[] col_arr = new String[] { "colname_1", "colname_2"};
预期直接传入数组调用concat函数生成新列,预期写法如下:
Dataset<Row> new_abc = dataset_abc.withColumn("new_concat_Column", concat(col_arr));
目前已知显式逐个传入col()转换后的列对象的写法可正常运行:
Dataset<Row> new_abc = dataset_abc.withColumn("new_concat_Column", concat(col("colname_1"), col("colname_2")));
需求为无需显式逐个传入列名,直接传入列名数组完成Dataset多列拼接。
实现原理
Spark SQL内置的concat函数方法签名为public static Column concat(Column... exprs),仅接收Column类型的可变参数,不支持直接传入字符串数组,因此只需要提前将字符串格式的列名批量转换为Column类型数组,即可直接传入方法完成调用。
具体实现代码
- 引入依赖类
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.functions; import org.apache.spark.sql.Column; import java.util.Arrays;
- 批量转换列名类型并调用拼接方法
// 将字符串数组批量转换为Column类型数组 Column[] concatColumns = Arrays.stream(col_arr) .map(functions::col) .toArray(Column[]::new); // 直接传入转换后的Column数组完成拼接,可变参数会自动适配数组输入 Dataset<Row> new_abc = dataset_abc.withColumn("new_concat_Column", functions.concat(concatColumns));
扩展:如果需要指定拼接分隔符,可替换为
concat_ws方法,第一个参数传入分隔符即可,后续列参数逻辑完全一致,示例代码如下:// 以下划线为分隔符拼接指定列 Dataset<Row> new_abc = dataset_abc.withColumn("new_concat_Column", functions.concat_ws("_", concatColumns));
内容的提问来源于stack exchange,提问作者Meta Sapien
相关产品推荐
相关产品推荐

