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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 18:00:05