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

PySpark:基于CSV列生成指定键的JSON列的实现问题

问题:根据指定字段列表生成JSON列

输入数据集

idfieldsf1f2f3f4
1f1, f2, f33202
2f2, f42425
3f17646

期望输出

idfieldsjson_field
1f1, f2, f3{"f1": 3, "f2": 2, "f3": 0}
2f2, f4{"f2": 4, "f4": 5}
3f1{"f1": 7}

原错误代码

input_df.selet(
    col("id"),
    col("fields"),
    to_json(struct(split(col("fields"), ","))).alias("json_field")
)

问题分析

  1. 拼写错误:selet 应为 select
  2. 逻辑错误:split(col("fields"), ",")返回字符串数组,而struct()需要传入具体列对象,无法直接用数组生成包含对应字段值的结构体。

解决方案

方法1:SQL表达式动态构建结构体(性能更优)

from pyspark.sql import functions as F

input_df = input_df.withColumn(
    "json_field",
    F.to_json(
        F.expr("struct(" + ", ".join([f"`{f.strip()}`" for f in F.split(F.col("fields"), ",")]) + ")")
    )
).select("id", "fields", "json_field")

方法2:map_from_arrays生成键值对后转JSON

from pyspark.sql import functions as F

input_df = input_df.withColumn(
    "field_names",
    F.transform(F.split(F.col("fields"), ", "), lambda x: F.trim(x))
).withColumn(
    "field_values",
    F.transform(F.col("field_names"), lambda x: F.col(x))
).withColumn(
    "json_field",
    F.to_json(F.map_from_arrays(F.col("field_names"), F.col("field_values")))
).select("id", "fields", "json_field")

代码说明

  • 两种方法都会先清理fields中的空格,确保字段名能匹配DataFrame列名。
  • 方法1通过拼接SQL表达式直接生成结构体,适合常规字段名场景;方法2用键值对映射更灵活,可兼容特殊字段名。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 12:50:03