如何用PySpark将3000万条销售数据按Item聚合为嵌套列表?
使用PySpark实现销售数据聚合转换
针对你3000万条销售数据的聚合需求,可以通过两次分组聚合实现目标格式,具体步骤如下:
步骤1:按item+type分组,收集单类型字段列表
首先将每个item下的同type数据聚合,把days_diff、placed_orders、cancelled_orders分别收集为单类型对应的列表:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, collect_list # 初始化SparkSession spark = SparkSession.builder.appName("SalesDataAggregation").getOrCreate() # 假设你的原始数据已经加载为DataFrame df # 示例数据加载(可替换为你的实际数据读取逻辑) data = [ ("console", "ps5", -10, 8, 1), ("console", "xbox", -8, 6, 0), ("console", "ps5", -5, 4, 4), ("console", "xbox", -1, 10, 7), ("console", "xbox", 0, 2, 3), ("games", "ps5", -11, 48, 9), ("games", "ps5", -3, 2, 4), ("games", "xbox", 5, 10, 2) ] df = spark.createDataFrame(data, ["item", "type", "days_diff", "placed_orders", "cancelled_orders"]) # 第一次分组:按item和type聚合,收集每个type的字段列表 grouped_by_type = df.groupBy("item", "type") \ .agg( collect_list(col("days_diff")).alias("days_list"), collect_list(col("placed_orders")).alias("placed_list"), collect_list(col("cancelled_orders")).alias("cancelled_list") )
步骤2:按item分组,生成嵌套列表和类型列表
将第一步的结果按item再次聚合,收集所有type到types列表,同时将每个类型对应的字段列表收集为嵌套列表:
# 第二次分组:按item聚合,生成最终嵌套格式 # 先按type排序,保证types顺序与嵌套列表顺序一致 final_df = grouped_by_type.orderBy("type") \ .groupBy("item") \ .agg( collect_list(col("type")).alias("types"), collect_list(col("days_list")).alias("days_diff"), collect_list(col("placed_list")).alias("placed_orders"), collect_list(col("cancelled_list")).alias("cancelled_orders") ) # 查看结果 final_df.show(truncate=False)
性能优化建议
针对3000万条数据的规模,可做以下优化:
- 调整Shuffle分区数:设置
spark.sql.shuffle.partitions为合适值(如根据CPU核心数调整,默认200,可设为1000左右),避免Shuffle过程中分区过多/过少影响性能。 - 数据分区预处理:如果原始数据有分区字段,可利用分区减少后续分组的Shuffle量。
- 内存配置:确保Spark executor有足够内存处理聚合过程中的数据缓存。
内容的提问来源于stack exchange,提问作者ohh
相关产品推荐
相关产品推荐

