PySpark嵌套JSON展平失败求助:无法获取预期输出
解决PySpark展平嵌套JSON时数组对应展开的问题
问题分析
你的通用flatten函数会对event_properties和entities两个数组分别执行explode操作,这会产生笛卡尔积(比如start事件会生成9行而非3行),无法满足你需要的数组元素按索引一一对应的展开需求。
正确实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, arrays_zip, col # 初始化SparkSession spark = SparkSession.builder \ .master("local[1]") \ .appName("PySpark Read JSON") \ .getOrCreate() # 读取JSON文件 df = spark.read.option("multiline","true").json(r"C:\Users\Lajo\Downloads\spark_ex1_input.json") # 步骤1:展开最外层的events数组 df_events = df.select(explode(col("events")).alias("event")) # 步骤2:将event结构体展开为单独字段 df_expanded = df_events.select( col("event.event_name"), col("event.event_properties"), col("event.entities"), col("event.event_timestamp") ) # 步骤3:用arrays_zip将两个数组按位置配对,再展开配对后的数组 df_final = df_expanded.select( col("event_name"), col("event_timestamp"), explode(arrays_zip(col("event_properties"), col("entities"))).alias("zipped") ).select( col("event_name"), col("zipped.event_properties").alias("event_properties"), col("zipped.entities").alias("entities"), col("event_timestamp") ) # 查看结果 df_final.show()
代码解释
- 展开events数组:先把最外层的
events数组展开,得到每个独立的event行。 - 展开event结构体:将event结构体中的字段提取为单独列,方便后续处理。
- 数组配对与展开:
arrays_zip函数会把event_properties和entities两个数组中同索引的元素打包成结构体(比如(property1, entityI)),生成一个新的结构体数组。- 对这个结构体数组执行
explode,就能得到一一对应的行。
- 提取最终字段:从结构体中提取
event_properties和entities,保留其他字段,得到你需要的输出格式。
预期输出
+----------+----------------+----------+-------------------+ |event_name|event_properties| entities| event_timestamp| +----------+----------------+----------+-------------------+ | start| property1| entityI|2022-05-01 00:00:00| | start| property2| entityII|2022-05-01 00:00:00| | start| property3|entityIII|2022-05-01 00:00:00| | stop| propertyA| entityW|2022-05-01 01:00:00| | stop| propertyB| entityX|2022-05-01 01:00:00| | stop| propertyC| entityY|2022-05-01 01:00:00| | stop| propertyD| entityZ|2022-05-01 01:00:00| +----------+----------------+----------+-------------------+
内容的提问来源于stack exchange,提问作者Priyanka
相关产品推荐
相关产品推荐

