如何使用PySpark DataFrame处理指定格式键值对文本得到目标结构化结果
PySpark 实现键值对行转列处理方案
核心思路
- 读取原始文本后先做基础清洗,移除每行首尾的单引号
- 为每行分配唯一行标识,方便后续分组还原每行数据
- 将每行的逗号分隔键值对拆分为多行<键,值>结构
- 用pivot方法完成行转列,手动指定a-j全量列名,保证列顺序正确且缺失列自动补全
- 空值替换为空字符串后,所有列用|拼接为目标输出格式
完整代码
from pyspark.sql import SparkSession from pyspark.sql.functions import ( split, explode, regexp_replace, col, concat_ws, monotonically_increasing_id ) # 初始化SparkSession spark = SparkSession.builder.appName("kv_transform").getOrCreate() # 1. 读取原始文本文件 raw_df = spark.read.text("your_file_path.txt") # 2. 清洗数据+分配行ID clean_df = raw_df.select( regexp_replace(col("value"), "^'|'$", "").alias("line") ).withColumn("row_id", monotonically_increasing_id()) # 3. 拆分键值对并展开为多行 kv_df = clean_df.select( "row_id", explode(split(col("line"), ",")).alias("kv") ).select( "row_id", split(col("kv"), ":")[0].alias("key"), split(col("kv"), ":")[1].alias("value") ) # 4. 行转列补全所有列 target_cols = ["a", "b", "c", "d", "e", "f", "g", "h", "i", "j"] pivot_df = kv_df.groupBy("row_id")\ .pivot("key", target_cols)\ .agg({"value": "first"})\ .na.fill("")\ .orderBy("row_id") # 5. 拼接为|分隔的输出格式 result_df = pivot_df.select(concat_ws("|", *target_cols).alias("output")) # 打印结果(如果要写出到文件直接用result_df.write.text("output_path")即可) print("|".join(target_cols)) for row in result_df.collect(): print(row["output"])
关键说明
- pivot方法传入第二个参数
target_cols可以强制指定输出列的顺序,同时自动补充原始数据中不存在的g列,避免pivot自动按key出现顺序排列的问题 na.fill("")将不存在的key对应值替换为空字符串,和要求的输出格式对齐- 行ID的作用是区分不同行的键值对,避免pivot时所有数据合并为一行
内容的提问来源于stack exchange,提问作者Ismael
相关产品推荐
相关产品推荐

