使用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中改为LongTypeLegacyStoreId:JSON中是数值类型,Schema中改为LongTypeFirstDeliveryTime: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
相关产品推荐
相关产品推荐

