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

使用PySpark解析Kafka流JSON时Stores数组返回空值的求助

修复PySpark解析Kafka流JSON中Stores ArrayType字段为空的问题

问题描述

在使用PySpark处理Kafka流数据时,无法正确解析JSON中的Stores ArrayType字段,解析后该字段返回空值。目标JSON数据结构如下:

{"CustomerId":606811,"Latitude":35.834896,"Longitude":50.019657,"Response":{"Stores":[{"Id":771,"LegacyStoreId":5497,"LegacyStoreTypeId":1,"PartnerId":3,"StoreName":"test","StoreDisplayName":"test","PartnerName":"test","ServiceRadius":3,"Longitude":56.009797,"Latitude":35.829067,"StatusCode":1,"CityId":200,"CityName":null,"Description":"test","IsOk24":false,"RouteDistanceInMeter":0,"IsRouteDistanceValid":false,"IsOutOfOrders":false,"AirDistanceInMeter":1100,"IsAirDistanceInValid":true,"IsDeliveryCoverage":true,"IsNonCoverageArea":false,"Rate":3.9,"Reviews":560,"IsHighPriorityStore":false,"StoreScore":0,"PartnerRank":1,"DeliveryCost":"80000","FirstDeliveryTime":"test","Labels":[],"MinCartTotalPrice":500000},{"Id":463,"LegacyStoreId":4765,"LegacyStoreTypeId":1,"PartnerId":3,"StoreName":"test","StoreDisplayName":"test","PartnerName":"test","ServiceRadius":3,"Longitude":56.995281,"Latitude":35.82251,"StatusCode":1,"CityId":200,"CityName":null,"Description":"test","IsOk24":false,"RouteDistanceInMeter":0,"IsRouteDistanceValid":false,"IsOutOfOrders":false,"AirDistanceInMeter":2593,"IsAirDistanceInValid":true,"IsDeliveryCoverage":true,"IsNonCoverageArea":false,"Rate":3.8,"Reviews":532,"IsHighPriorityStore":false,"StoreScore":0,"PartnerRank":1,"DeliveryCost":"80000","FirstDeliveryTime":"test","Labels":[],"MinCartTotalPrice":500000}]},"Id":"f2655da4-c236-4f86-9ca0-8063a4c77da8","CreationAt":"2023-06-18T14:49:11.8545562+03:30"}

用户的PySpark代码如下:

spark= SparkSession \
    .builder \
    .appName("striming") \
    .config("spark.jars.packages",
            "*****************") \
    .config('spark.driver.extraClassPath', '/usr/local/spark/resources/jars/sqljdbc42.jar') \
    .config('spark.executor.extraClassPath', '/usr/local/spark/resources/jars/sqljdbc42.jar') \
    .config("spark.cores.max", "1") \
    .config("spark.executor.memory", "1g") \
    .config("spark.executor.cores", "1") \
    .config("spark.dynamicAllocation.initialExecutors", "1") \
    .master("local[1]") \
    .getOrCreate()

sc = spark.sparkContext

sqlContext = SQLContext(sc)

kafka_df = spark.readStream.format("kafka").option("kafka.bootstrap.servers",
                                                   "My Kafka Servers").option(
    "subscribe",
    "Event").option("startingOffsets", "earliest").option("failOnDataLoss", "false").option('multiline', "True").load()


schema = StructType([StructField("CreationAt", StringType(), True),
             StructField("CustomerId", LongType(), True),
             StructField("Id", StringType(), True),
             StructField("Latitude", DoubleType(), True),
             StructField("Longitude", DoubleType(), True),
             StructField("Response",
                 StructType(
                     [StructField("Stores",
                        ArrayType(StructType(
                             [
                                 StructField("Id", StringType(), True),
                                 StructField("LegacyStoreId", StringType(), True),
                                 StructField("PartnerName", StringType(), True),
                                 StructField("FirstDeliveryTime", LongType(), True),
                                 StructField("StatusCode", LongType(), True),
                                 StructField("StoreName", StringType(), True)]), True),
                                                                         True)]), True)
                             ]
            )


value_df = kafka_df.select(from_json(col("value").cast("string"), schema).alias("StoreSelection"))

value_df.printSchema()

sites_flat = value_df.selectExpr("StoreSelection")\
                       .select("StoreSelection.CustomerId", "StoreSelection.Latitude", "StoreSelection.Longitude", "StoreSelection.CreationAt", explode_outer("StoreSelection.Response.Stores").alias("Stores"))\
                       .select( "CustomerId","Latitude","Longitude", "CreationAt","Stores.LegacyStoreId", "Stores.StoreName", "Stores.PartnerName")\
                       .select( "CustomerId","Latitude","Longitude", "CreationAt","LegacyStoreId", "StoreName", "PartnerName")

sites_flat.printSchema()

def foreach_batch_function(df, epoch_id):

   df.write \
   df.show()

invoiceWriterQuery = sites_flat.writeStream.foreachBatch(foreach_batch_function).outputMode("update") \
   .option("checkpointLocation", "/usr/local/airflow/dags/otime/log_file").trigger(
    processingTime="1 minute").start().awaitTermination()

invoiceWriterQuery.awaitTermination()

解决方案

1. 修正Schema字段类型不匹配问题

Stores数组解析为空的核心原因是Schema定义的字段类型与JSON实际值类型不匹配,PySpark在遇到类型不匹配时会将整个复杂字段(如数组、结构体)置为空。需要修正以下字段的类型:

  • Id:JSON中是数值类型,Schema中改为LongType
  • LegacyStoreId:JSON中是数值类型,Schema中改为LongType
  • FirstDeliveryTime:JSON中是字符串类型(如"test"),Schema中改为StringType

修正后的Schema如下:

from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType, ArrayType

schema = StructType([
    StructField("CreationAt", StringType(), True),
    StructField("CustomerId", LongType(), True),
    StructField("Id", StringType(), True),
    StructField("Latitude", DoubleType(), True),
    StructField("Longitude", DoubleType(), True),
    StructField("Response", StructType([
        StructField("Stores", ArrayType(StructType([
            StructField("Id", LongType(), True),
            StructField("LegacyStoreId", LongType(), True),
            StructField("PartnerName", StringType(), True),
            StructField("FirstDeliveryTime", StringType(), True),
            StructField("StatusCode", LongType(), True),
            StructField("StoreName", StringType(), True)
        ])), True)
    ]), True)
])

2. 移除无效配置

Kafka数据源没有multiline配置项,该配置仅适用于JSON文件读取,需移除:

kafka_df = spark.readStream.format("kafka")\
    .option("kafka.bootstrap.servers", "My Kafka Servers")\
    .option("subscribe", "Event")\
    .option("startingOffsets", "earliest")\
    .option("failOnDataLoss", "false")\
    .load()

3. 修复foreach_batch_function语法错误

原代码中foreach_batch_function存在语法错误,df.write \未完成写入逻辑且与df.show()的写法冲突,调整为:

def foreach_batch_function(df, epoch_id):
    # 打印解析后的数据
    df.show(truncate=False)
    # 若需写入存储,补充完整的write逻辑,例如:
    # df.write.mode("append").format("jdbc").options(url="xxx", dbtable="xxx", user="xxx", password="xxx").save()

invoiceWriterQuery = sites_flat.writeStream.foreachBatch(foreach_batch_function)\
    .outputMode("update")\
    .option("checkpointLocation", "/usr/local/airflow/dags/otime/log_file")\
    .trigger(processingTime="1 minute")\
    .start()

invoiceWriterQuery.awaitTermination()

验证

修正后重新运行代码,Stores数组将被正常解析,sites_flat中的LegacyStoreId、StoreName等字段会正确显示JSON中的对应值。


内容的提问来源于stack exchange,提问作者Ali Moayed

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 06:37:09