PySpark:如何将两个JSON列合并为格式正常的新列
PySpark合并JSON列生成无转义的正常JSON新列
问题场景
现有PySpark表包含两个JSON格式列Col1和Col2,结构如下:
Col1结构
{'table': [{'name': 'XXS', 'ranges': {'chestc': {'min': 87.88, 'max': 87.88}, 'waistc': {'min': 58.42, 'max': 58.42}}}, {'name': 'XS', 'ranges': {'chestc': {'min': 94.22, 'max': 94.22}, 'waistc': {'min': 66.04, 'max': 66.04}}}, {'name': 'S', 'ranges': {'chestc': {'min': 100.58, 'max': 100.58}, 'waistc': {'min': 73.66, 'max': 73.66}}}, {'name': 'M', 'ranges': {'chestc': {'min': 106.92, 'max': 106.92}, 'waistc': {'min': 81.28, 'max': 81.28}}}, {'name': 'L', 'ranges': {'chestc': {'min': 114.54, 'max': 114.54}, 'waistc': {'min': 93.98, 'max': 93.98}}}, {'name': 'XL', 'ranges': {'chestc': {'min': 122.16, 'max': 122.16}, 'waistc': {'min': 106.68, 'max': 106.68}}}, {'name': 'XXL', 'ranges': {'chestc': {'min': 131.06, 'max': 131.06}, 'waistc': {'min': 121.92, 'max': 121.92}}}], 'measurement_system': 'metric'}
Col2结构
{ "gender": "male", "measurement_system": "metric", "measurements": { "height": 178, "weight": 99 } }
需要生成新列Col3,将两列的JSON内容合并为一个无转义符的正常JSON对象。使用F.struct再转JSON会出现大量转义符,不符合需求,要求仅使用PySpark原生函数,禁止使用UDF、Python非PySpark函数或collect()。
解决方案
核心思路是先将JSON字符串解析为PySpark结构化数据(Struct/Array/Map类型),再合并结构,最后转换为无转义的JSON。
步骤1:定义JSON对应Schema
为Col1和Col2的JSON结构定义匹配的Schema,确保解析准确:
from pyspark.sql import functions as F from pyspark.sql.types import ( StructType, StructField, StringType, ArrayType, MapType, DoubleType, IntegerType ) # Col1的Schema size_range_schema = StructType([ StructField("min", DoubleType(), True), StructField("max", DoubleType(), True) ]) ranges_schema = StructType([ StructField("chestc", size_range_schema, True), StructField("waistc", size_range_schema, True) ]) table_item_schema = StructType([ StructField("name", StringType(), True), StructField("ranges", ranges_schema, True) ]) col1_schema = StructType([ StructField("table", ArrayType(table_item_schema), True), StructField("measurement_system", StringType(), True) ]) # Col2的Schema measurements_schema = StructType([ StructField("height", IntegerType(), True), StructField("weight", IntegerType(), True) ]) col2_schema = StructType([ StructField("gender", StringType(), True), StructField("measurement_system", StringType(), True), StructField("measurements", measurements_schema, True) ])
如果不确定Schema,可通过样本JSON自动推断:
# 自动推断Col1的Schema col1_sample = """{"table": [{"name": "XXS","ranges": {"chestc": {"min": 87.88, "max": 87.88},"waistc": {"min": 58.42, "max": 58.42}}}], "measurement_system": "metric"}""" col1_inferred_schema = F.json_schema(F.lit(col1_sample))
步骤2:解析JSON列为结构化数据
用F.from_json将JSON字符串解析为结构化列:
df_parsed = df.withColumn("col1_struct", F.from_json(F.col("Col1"), col1_schema)) \ .withColumn("col2_struct", F.from_json(F.col("Col2"), col2_schema))
步骤3:合并结构并转换为无转义JSON
使用F.struct合并两个结构化列的字段,再用F.to_json生成正常JSON。对于重复字段(如measurement_system),可选择保留其中一个或重命名:
# 合并字段示例:保留Col1的尺码表+Col2的所有非重复字段 df_final = df_parsed.withColumn( "Col3", F.to_json( F.struct( F.col("col1_struct.table").alias("size_table"), F.col("col2_struct.gender"), F.col("col2_struct.measurement_system"), F.col("col2_struct.measurements") ) ) ) # 若需保留所有字段(自动覆盖重复字段,以最后声明的为准) # df_final = df_parsed.withColumn( # "Col3", # F.to_json( # F.struct( # F.col("col1_struct.*"), # F.col("col2_struct.gender"), # F.col("col2_struct.measurements") # ) # ) # )
关键说明
直接对原始JSON字符串使用F.struct会导致嵌套JSON被转义,因为PySpark会将字符串类型的JSON当作普通文本处理。只有先解析为结构化数据,合并后再转JSON,才能生成无转义的正常嵌套JSON。
内容的提问来源于Stack Exchange,提问作者Galat
相关产品推荐
相关产品推荐

