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

Spark Dataset分组聚合时限制列表长度并拆分的实现求助

实现方案

你可以直接使用Scala原生集合的grouped()方法搭配Spark的flatMapGroups算子完成需求,修改后的核心代码如下:

// 自定义每个booking列表的最大长度
val maxListLength = 2

val resultDS: Dataset[AggrBook] = bookDS.groupByKey(r => (r.city, r.state))
  .flatMapGroups((key, value) => {
    // 先收集当前分组的所有bookingId
    val allBookIds = value.map(_.bookingId).toList
    // 按固定长度拆分列表,每个拆分后的子列表对应一条输出记录
    allBookIds.grouped(maxListLength).map(subList => AggrBook(key._2, key._1, subList))
  })

resultDS.show()

运行上述代码后输出结果和你预期的完全一致:

+----+---------+------------+
|city|    state|       books|
+----+---------+------------+
|  CA|   Eureka|      [1, 2]|
|  CA|   Eureka|      [3, 5]|
|  AK|Kettlemen|         [4]|
|  CA|   Irvine|         [6]|
+----+---------+------------+

注意:如果单分组数据量极大,建议先通过窗口函数给每个分组内的元素打上行号,再按行号分片处理,避免单task内存溢出,常规场景下上述方案可直接使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 21:57:00