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

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'))

当前输出与期望输出

当前输出

当前输出

期望输出

期望输出

解决方案

核心思路

  1. 自动推断嵌套JSON的Schema:先解析顶层CDC JSON,提取嵌套JSON字符串字段,再通过样本数据推断其结构化Schema。
  2. 递归展平嵌套结构:编写通用函数自动展平所有嵌套的结构体字段,避免手动指定列名。
  3. 多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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 07:55:29