Spark中按order_id聚合行对象生成嵌套订单结构的实现方法
解决方案
核心误区澄清
你担心不同order_id的数据会在merge阶段合并完全是多余的:UDAF的merge操作仅会处理同一个group by分组内的中间聚合结果,Spark执行group by时已经将相同order_id的所有数据路由到同一个处理单元,不同order_id的计算链路完全隔离,不会出现跨分组合并的情况。
方案1:直接使用内置函数实现(推荐,无需自定义UDAF)
你的需求完全可以通过Spark原生函数实现,无需额外开发UDAF,示例代码如下:
Spark SQL 示例
SELECT order_id, STRUCT( order_id, collect_list(STRUCT(item_id, price)) AS items ) AS order FROM ( -- 先拆解item结构体的字段 SELECT item.order_id AS order_id, item.item_id AS item_id, item.price AS price FROM source_table ) t GROUP BY order_id
Scala 示例
import org.apache.spark.sql.functions.{collect_list, struct} sourceTable .select( $"item.order_id".alias("order_id"), $"item.item_id".alias("item_id"), $"item.price".alias("price") ) .groupBy("order_id") .agg( collect_list(struct("item_id", "price")).alias("items") ) .select( $"order_id", struct("order_id", "items").alias("order") )
方案2:自定义Aggregator类型UDAF的实现方式
如果有特殊场景需要自定义UDAF,可以参考如下逻辑实现,重点明确merge方法的作用:
- 定义三个泛型:
- 输入类型IN:对应item字段的结构体类型
- 缓冲类型BUF:样例类
AggBuffer(var orderId: Long, var items: ArrayBuffer[ItemStruct]),用于存储当前分组的order_id和已收集的商品列表 - 输出类型OUT:对应最终的order结构体类型
- 核心方法实现:
zero: 初始化缓冲,orderId设为0,items设为空ArrayBufferreduce: 如果缓冲orderId为0则赋值为当前输入的order_id,将当前输入的item_id和price组成的结构体加入itemsmerge: 直接将两个同分组的缓冲的items合并,返回新的缓冲即可,示例逻辑:buffer1.items ++= buffer2.items; buffer1finish: 将缓冲的orderId和items组装为最终的order结构体返回
内容的提问来源于stack exchange,提问作者user1888955
相关产品推荐
相关产品推荐

