You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Scala:依据另一列指定的列名提取对应列值

动态提取指定列值并拼接成新列的Spark实现

实现思路

利用Spark的create_map、split、transform和array_join函数组合,动态根据best_col中的列名提取对应值并拼接:

  1. 将所有列名与对应值构建成Map结构,方便按列名取值
  2. 拆分best_col为列名字符串数组
  3. 遍历数组从Map中提取对应列值,生成值数组
  4. 将值数组用#连接成最终的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 20:42:03