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
相关产品推荐
相关产品推荐

