PySpark:基于CSV列生成指定键的JSON列的实现问题
问题:根据指定字段列表生成JSON列
输入数据集
| id | fields | f1 | f2 | f3 | f4 |
|---|---|---|---|---|---|
| 1 | f1, f2, f3 | 3 | 2 | 0 | 2 |
| 2 | f2, f4 | 2 | 4 | 2 | 5 |
| 3 | f1 | 7 | 6 | 4 | 6 |
期望输出
| id | fields | json_field |
|---|---|---|
| 1 | f1, f2, f3 | {"f1": 3, "f2": 2, "f3": 0} |
| 2 | f2, f4 | {"f2": 4, "f4": 5} |
| 3 | f1 | {"f1": 7} |
原错误代码
input_df.selet( col("id"), col("fields"), to_json(struct(split(col("fields"), ","))).alias("json_field") )
问题分析
- 拼写错误:
selet应为select - 逻辑错误:
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
相关产品推荐
相关产品推荐

