PySpark将嵌套JSON字符串动态解析为DataFrame列的问题
动态解析CDC嵌套JSON至Databricks Delta表的解决方案
问题背景
处理从Kafka接收的CDC数据并加载到Databricks Delta表时,除嵌套JSON字符串(如input_data)外,其余字段均可正常解析。使用from_json、spark.read.json无法自动解析这类嵌套JSON,且spark.read.json(df.rdd.map(lambda row: row.value)).schema会将INPUT_DATA识别为字符串类型。需要实现多Topic动态处理,无需预先存储Schema,自动适配Schema变更。
示例嵌套JSON
{ "after": { "id_transaction": "121", "product_id": 25, "transaction_dt": 1662076800000000, "creation_date": 1662112153959000, "product_account": "40012", "input_data": "{\"amount\":[{\"type\":\"CASH\",\"amount\":1000.00}],\"currency\":\"USD\",\"coreData\":{\"CustId\":11021,\"Cust_Currency\":\"USD\",\"custCategory\":\"Premium\"},\"context\":{\"authRequired\":false,\"waitForConfirmation\":false,\"productAccount\":\"CA12001\"},\"brandId\":\"TOYO-2201\",\"dealerId\":\"1\",\"operationInfo\":{\"trans_Id\":\"3ED23-89DKS-001AA-2321\",\"transactionDate\":1613420860087},\"ip_address\":null,\"last_executed_step\":\"PURCHASE_ORDER_CREATED\",\"last_result\":\"OK\",\"output_dataholder\":\"{\\\"DISCOUNT_AMOUNT\\\":\\\"0\\\",\\\"BONUS_AMOUNT_APPLIED\\\":\\\"10000\\\"}\"", "dealer_id": 1, "dealer_currency": "USD", "Cust_id": 11021, "process_status": "IN_PROGRESS", "tot_amount": 10000, "validation_result_code": "OK_SAVE_AND_PROCESS", "operation": "Create", "timestamp_ms": 1675673484042 }, "op": "c" }
尝试过的代码
提取JSON结构列的代码
import json json_keys = {} child_members = [] table_column_schema = {} column_schema = [] dbname = "mydb" tbl_name = "tbl_name" def get_table_keys(dbname): table_values_extracted = "select value from {mydb}.{tbl_name} limit 1" cmd_key_pair_data = spark.sql(table_values_extracted) jsonkeys=cmd_key_pair_data.collect()[0][0] json_keys = json.loads(jsonkeys) column_names_as_keys = json_keys["after"].keys() value_column_data = json_keys["after"].values() column_schema = list(column_names_as_keys) for i in value_column_data: if ("{" in str(i) and "}" in str(i)): a = json.loads(i) for i2 in a.values(): if (str(i2).startswith("{") and str(i2).endswith('}')): column_schema = column_schema + list(i2.keys()) table_column_schema['temp_table1'] = column_schema return 0 get_table_keys("dbname")
处理JSON生成DataFrame的代码
from pyspark.sql.functions import from_json, to_json, col from pyspark.sql.types import StructType, StructField, StringType, LongType, MapType import time dbname = "mydb" tbl_name = "tbl_name" start = time.time() df = spark.sql(f'select value from {mydb}.{tbl_name} limit 2') tbl_columns = table_column_schema[tbl_name] data = [] for i in tbl_columns: if i == 'input_data': data.append(StructField(f'{i}', MapType(StringType(),StringType()), True)) else: data.append(StructField(f'{i}', StringType(), True)) schema2 = spark.read.json(df.rdd.map(lambda row: row.value)).schema print(type(schema2)) df2 = df.withColumn("value", from_json("value", schema2)).select(col('value.after.*'), col('value.op'))
当前输出与期望输出
当前输出

期望输出

解决方案
核心思路
- 自动推断嵌套JSON的Schema:先解析顶层CDC JSON,提取嵌套JSON字符串字段,再通过样本数据推断其结构化Schema。
- 递归展平嵌套结构:编写通用函数自动展平所有嵌套的结构体字段,避免手动指定列名。
- 多Topic动态适配:封装处理逻辑为通用函数,自动识别所有JSON字符串字段,适配不同Topic的Schema差异。
完整实现代码
1. 定义通用展平函数
from pyspark.sql.functions import from_json, col, rlike from pyspark.sql.types import StructType def flatten_struct(df, prefix=""): """递归展平DataFrame中的结构体字段""" fields = [] for field in df.schema.fields: col_alias = f"{prefix}{field.name}" if prefix else field.name if isinstance(field.dataType, StructType): # 递归处理嵌套结构体 nested_df = flatten_struct(df.select(field.name), prefix=f"{col_alias}_") fields.extend(nested_df.columns) else: fields.append(col(field.name).alias(col_alias)) return df.select(fields)
2. 编写动态CDC处理函数
def process_cdc_topic(dbname, tbl_name, sample_size=10): """ 动态处理指定Topic的CDC数据,自动解析嵌套JSON并展平结构 :param dbname: 数据库名 :param tbl_name: 青铜层表名 :param sample_size: 用于推断Schema的样本数量 """ # 读取青铜层原始数据 df_raw = spark.sql(f'select value from {dbname}.{tbl_name}') # 推断顶层CDC数据的Schema并解析 top_schema = spark.read.json(df_raw.rdd.map(lambda row: row.value)).schema df_parsed = df_raw.withColumn("value", from_json("value", top_schema)) \ .select("value.after.*", "value.op") # 自动识别所有可能的JSON字符串字段(包含{}的非空字符串) json_string_cols = [] for col_name, dtype in df_parsed.dtypes: if dtype == "string": has_json = df_parsed.filter(col(col_name).isNotNull() & col(col_name).rlike(r"\{.*\}")).count() > 0 if has_json: json_string_cols.append(col_name) # 逐个解析嵌套JSON字段 for col_name in json_string_cols: # 提取样本数据推断嵌套Schema samples = df_parsed.select(col_name) \ .filter(col(col_name).isNotNull()) \ .limit(sample_size) \ .rdd.map(lambda x: x[0]) \ .collect() if samples: # 推断嵌套JSON的Schema nested_schema = spark.read.json(spark.sparkContext.parallelize(samples)).schema # 将JSON字符串解析为结构体 df_parsed = df_parsed.withColumn(f"{col_name}_struct", from_json(col(col_name), nested_schema)) # 展平结构体并移除原字段和中间结构体字段 df_parsed = flatten_struct(df_parsed.select("*", f"{col_name}_struct.*")) \ .drop(col_name, f"{col_name}_struct") # 写入黄金层Delta表(自动适配Schema变更) df_parsed.write.mode("append") \ .format("delta") \ .option("mergeSchema", "true") \ .saveAsTable(f"{dbname}.{tbl_name}_gold") return df_parsed
3. 调用处理函数
# 处理指定Topic的数据 processed_df = process_cdc_topic("mydb", "tbl_name") # 查看结果 processed_df.show(truncate=False)
关键说明
- Schema自动适配:Spark的
read.json会自动合并样本中的Schema,支持新增字段的动态识别。 - 性能优化:仅使用少量样本(默认10条)推断Schema,避免全表扫描带来的性能损耗。
- 深层嵌套支持:递归展平函数可处理任意层级的嵌套结构体,无需手动配置。
- 多Topic兼容:通用函数可直接用于不同Topic的处理,无需修改核心逻辑。
内容的提问来源于stack exchange,提问作者Yuva
相关产品推荐
相关产品推荐

