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

PySpark 3.4.1 Kafka流式DataFrame的JSON数组展开与转换问题

PySpark 3.4.1 流式作业:扁平化嵌套JSON数组结构

问题核心

你遇到的explode无法展开SH数组的问题,大概率是因为通过json_get_object提取的SH字段仍是字符串类型,而非Spark可识别的ArrayType[StructType]结构。以下是完整的解决方案,可将原始DataFrame转换为目标扁平化结构:

完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, col, lit, from_json
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DoubleType, ArrayType

# 1. 定义嵌套JSON的Schema
# SH数组元素的结构体Schema
sh_item_schema = StructType([
    StructField("AJ", IntegerType(), True),
    StructField("BW", IntegerType(), True),
    StructField("CA", IntegerType(), True),
    StructField("CF", IntegerType(), True),
    StructField("DT", IntegerType(), True),
    StructField("EP", IntegerType(), True),
    StructField("ML", DoubleType(), True),
    StructField("IL", DoubleType(), True),  # 可选字段,允许为null
    StructField("FP", StructType([
        StructField("AD", IntegerType(), True),
        StructField("DD", IntegerType(), True),
        StructField("NA", IntegerType(), True),
        StructField("NW", IntegerType(), True)
    ]), True),
    StructField("NJ", StructType([
        StructField("FD", IntegerType(), True),
        StructField("PL", IntegerType(), True),
        StructField("TH", DoubleType(), True)
    ]), True)
])

# 根数据结构Schema
root_schema = StructType([
    StructField("ID", StringType(), True),
    StructField("SH", ArrayType(sh_item_schema), True)
])

# 2. 处理流式DataFrame
# 假设从Kafka消费的原始DataFrame为kafka_df,其中"value"是JSON字符串
parsed_df = kafka_df.select(col("value").cast(StringType()).alias("json_str")) \
    .select(from_json(col("json_str"), root_schema).alias("data")) \
    .select("data.*")

# 若你已通过json_get_object提取了ID和SH字符串,替换为以下代码:
# parsed_df = original_df.select(
#     col("ID"),
#     from_json(col("SH"), ArrayType(sh_item_schema)).alias("SH")
# )

# 3. 展开SH数组
exploded_df = parsed_df.select(col("ID"), explode(col("SH")).alias("sh_item"))

# 4. 扁平化嵌套字段并匹配目标结构
flattened_df = exploded_df.select(
    lit(None).alias("EN"),
    col("ID"),
    lit(None).alias("EP"),
    col("sh_item.AJ").alias("SH_AJ"),
    col("sh_item.BW").alias("SH_BW"),
    col("sh_item.CA").alias("SH_CA"),
    col("sh_item.CF").alias("SH_CF"),
    col("sh_item.DT").alias("SH_DT"),
    col("sh_item.EP").alias("SH_EP"),
    col("sh_item.FP.AD").alias("SH_FP_AD"),
    col("sh_item.FP.DD").alias("SH_FP_DD"),
    col("sh_item.FP.NA").alias("SH_FP_NA"),
    col("sh_item.FP.NW").alias("SH_FP_NW"),
    col("sh_item.NJ.FD").alias("SH_NJ_FD"),
    col("sh_item.NJ.PL").alias("SH_NJ_PL"),
    col("sh_item.NJ.TH").alias("SH_NJ_TH")
)

# 查看结果
flattened_df.show(truncate=False)

关键步骤说明

  1. Schema定义:严格匹配Kafka原始JSON的嵌套结构,确保字符串类型的SH被正确解析为数组+结构体类型,这是explode生效的前提。
  2. JSON解析:通过from_json将字符串类型的JSON转换为Spark可操作的结构化数据。
  3. 数组展开:使用explode将SH数组的每个元素拆分为单独行,保留对应的ID。
  4. 字段扁平化:提取嵌套结构体的字段并按目标格式重命名,同时添加值为null的EN、EP列,列顺序完全匹配目标结构。

注意事项

  • 确保Schema与JSON字段完全对应,可选字段需设置nullable=True(第三个参数)。
  • 该方案完全支持PySpark 3.4.1的流式处理场景,无需额外配置。

内容的提问来源于stack exchange,提问作者Virendar Kumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 19:50:56