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

Spark DataFrame按指定列名列表条件拆分字段问题求解

问题修正方案

你的代码问题在于循环中每次都会覆盖new_name和new_value列,最终只有最后一个列名(col4)的规则生效。我们可以通过动态正则匹配的方式一次性完成所有拆分逻辑,无需循环。

修正后的完整代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, explode, regexp_extract
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, ArrayType

# 定义结构体数组 schema
value_schema = ArrayType(
    StructType([
        StructField("name", StringType(), True),
        StructField("location", StringType(), True)
    ])
)

# 测试数据
data = [
    (1, [
        {"name": "col1_US", "location": "usa"},
        {"name": "col2_name_plex", "location": "usa"},
        {"name": "col4_false_val", "location": "usa"},
        {"name": "col3_name_is_fantasy", "location": "usa"}
    ])
]

# 初始化SparkSession
spark = SparkSession.builder.appName("NameSplit").getOrCreate()

# 创建DataFrame并展开数组
df = spark.createDataFrame(data, ["id", "values"])
df = df.withColumn("values", explode(col("values")))
df = df.select(col("id"), col("values.name").alias("name"))

# 待匹配的列名列表
col_names = ["col1","col2_name","col3_name_is","col4"]

# 构建动态正则表达式:匹配以指定列名开头并带下划线的格式
pattern = f"^({'|'.join(col_names)})_(.*)$"

# 提取匹配的列名和后续内容
df = df.withColumn("new_name", regexp_extract(col("name"), pattern, 1)) \
       .withColumn("new_value", regexp_extract(col("name"), pattern, 2))

df.display()

代码说明

  1. 动态正则构建:将col_names中的元素用|拼接成正则分组,匹配以指定列名开头并紧跟下划线的字符串。
  2. 正则提取:
    • regexp_extract(..., 1):提取匹配到的列名(第一个分组)作为new_name
    • regexp_extract(..., 2):提取下划线后的所有内容(第二个分组)作为new_value
  3. 避免循环覆盖:一次性完成所有行的匹配和拆分,无需逐一遍历列名。

输出结果

idnamenew_namenew_value
1col1_UScol1US
1col2_name_plexcol2_nameplex
1col4_false_valcol4false_val
1col3_name_is_fantasycol3_name_isfantasy

内容的提问来源于stack exchange,提问作者steve

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:02:04