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)
关键步骤说明
- Schema定义:严格匹配Kafka原始JSON的嵌套结构,确保字符串类型的SH被正确解析为数组+结构体类型,这是
explode生效的前提。 - JSON解析:通过
from_json将字符串类型的JSON转换为Spark可操作的结构化数据。 - 数组展开:使用
explode将SH数组的每个元素拆分为单独行,保留对应的ID。 - 字段扁平化:提取嵌套结构体的字段并按目标格式重命名,同时添加值为
null的EN、EP列,列顺序完全匹配目标结构。
注意事项
- 确保Schema与JSON字段完全对应,可选字段需设置
nullable=True(第三个参数)。 - 该方案完全支持PySpark 3.4.1的流式处理场景,无需额外配置。
内容的提问来源于stack exchange,提问作者Virendar Kumar
相关产品推荐
相关产品推荐

