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()
代码说明
- 动态正则构建:将
col_names中的元素用|拼接成正则分组,匹配以指定列名开头并紧跟下划线的字符串。 - 正则提取:
regexp_extract(..., 1):提取匹配到的列名(第一个分组)作为new_nameregexp_extract(..., 2):提取下划线后的所有内容(第二个分组)作为new_value
- 避免循环覆盖:一次性完成所有行的匹配和拆分,无需逐一遍历列名。
输出结果
| id | name | new_name | new_value |
|---|---|---|---|
| 1 | col1_US | col1 | US |
| 1 | col2_name_plex | col2_name | plex |
| 1 | col4_false_val | col4 | false_val |
| 1 | col3_name_is_fantasy | col3_name_is | fantasy |
内容的提问来源于stack exchange,提问作者steve
相关产品推荐
相关产品推荐

