Spark Scala:依据另一列指定的列名提取对应列值
动态提取指定列值并拼接成新列的Spark实现
实现思路
利用Spark的create_map、split、transform和array_join函数组合,动态根据best_col中的列名提取对应值并拼接:
- 将所有列名与对应值构建成Map结构,方便按列名取值
- 拆分
best_col为列名字符串数组 - 遍历数组从Map中提取对应列值,生成值数组
- 将值数组用
#连接成最终的result列
Python代码示例
from pyspark.sql import functions as F # 示例数据 data = [ ("A#B", 14, 26, 32), ("C", 13, 17, 96), ("B#C", 23, 19, 42) ] df = spark.createDataFrame(data, ["best_col", "A", "B", "C"]) # 构建列名-列值的Map df = df.withColumn("col_map", F.create_map(*[F.lit(c), F.col(c)] for c in df.columns)) # 拆分best_col为列名数组,提取对应值并拼接 df = df.withColumn("col_names", F.split(F.col("best_col"), "#")) \ .withColumn("result_values", F.transform(F.col("col_names"), lambda x: F.col("col_map").getItem(x))) \ .withColumn("result", F.array_join(F.col("result_values"), "#")) \ .drop("col_map", "col_names", "result_values") df.show()
Scala代码示例
import org.apache.spark.sql.functions._ // 示例数据 val data = Seq( ("A#B", 14, 26, 32), ("C", 13, 17, 96), ("B#C", 23, 19, 42) ) val df = spark.createDataFrame(data).toDF("best_col", "A", "B", "C") // 构建列名-列值的Map val colMap = createMap(df.columns.flatMap(c => Seq(lit(c), col(c))): _*) // 生成result列 val resultDF = df .withColumn("colNames", split(col("best_col"), "#")) .withColumn("resultValues", transform(col("colNames"), name => colMap.getItem(name))) .withColumn("result", array_join(col("resultValues"), "#")) .drop("colMap", "colNames", "resultValues") resultDF.show()
为什么col(col("best_col").toString())不生效
col()函数接受的是字符串字面量(比如col("A")),而col("best_col").toString()只是把Column对象转换成字符串(结果是"best_col"),并不能获取每行实际的列名值,所以无法动态引用对应列。
内容的提问来源于stack exchange,提问作者techie
相关产品推荐
相关产品推荐

