如何用PySpark将两个JSON字符串列合并为JSON列表字符串列?
问题
我有两个StringType类型的JSON字符串列,示例数据如下:
| JSON 1 | JSON 2 |
|---|---|
| {"key1":"value1"} | {"key2":"value2"} |
| {"key3":"value3"} | {"key4":"value4"} |
需要将这两列合并为第三列(同样为字符串类型),该列的值是包含这两个JSON的列表,示例结果如下:
| JSON 1 | JSON 2 | JSON 3 |
|---|---|---|
| {"key1":"value1"} | {"key2":"value2"} | [{"key1":"value1"},{"key2":"value2"}] |
| {"key3":"value3"} | {"key4":"value4"} | [{"key3":"value3"},{"key4":"value4"}] |
我已经用以下PySpark代码实现了需求,但觉得不是最优方案,想找更好的实现思路:
from pyspark.sql import DataFrame from pyspark.sql.functions import col, concat, concat_ws, lit def concat_jsons(df: DataFrame, columns: list): df = df.withColumn( 'JSON 3', concat_ws(',', *columns) ) return df.withColumn( 'JSON 3', concat(lit('['), col('JSON 3'), lit(']')) )
优化方案
思路1:利用Spark数组与JSON转换函数(推荐)
原方法存在隐患:如果某列的JSON字符串本身包含逗号(比如{"key":"val,ue"}),直接用concat_ws(',', ...)会把JSON内部的逗号当成分隔符,导致最终的JSON列表格式错误。
更可靠的方式是先将字符串列解析为JSON对象,组合成数组后再转回JSON字符串,从根源避免格式破坏:
from pyspark.sql import DataFrame from pyspark.sql.functions import col, array, to_json, from_json, schema_of_json def concat_jsons_optimized(df: DataFrame, columns: list): # 提取第一个JSON列的Schema,用于统一解析所有目标列 sample_schema = schema_of_json(df.select(col(columns[0])).first()[0]) # 将每个字符串列解析为JSON对象,组合成数组后转为JSON字符串 json_array = array(*[from_json(col(c), sample_schema) for c in columns]) return df.withColumn('JSON 3', to_json(json_array))
如果各JSON列结构不一致,可改用通用Schema兼容,比如from_json(col(c), "map<string, string>")或from_json(col(c), "struct<*>")。
思路2:简化字符串拼接(仅适用于无内部逗号的场景)
若能确保所有JSON字符串内部绝对没有逗号,可简化原代码为一次拼接操作,减少中间列创建:
from pyspark.sql import DataFrame from pyspark.sql.functions import col, concat, lit def concat_jsons_simplified(df: DataFrame, columns: list): return df.withColumn( 'JSON 3', concat(lit('['), col(columns[0]), lit(','), col(columns[1]), lit(']')) )
但该方法仍存在格式风险,不建议在生产环境使用。
内容的提问来源于stack exchange,提问作者joao pedro imamura
相关产品推荐
相关产品推荐

