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

PySpark中实现多层groupBy嵌套分组的技术咨询

实现RDD双层分组(ID内按Duration再分组)

首先先确认你的原始数据:

iddurationaction1action2
110AD
110BE
225AE
17AG

你已经完成了按ID的基础分组,现在要在每个ID组内再按duration分组并整理成指定的嵌套结构,咱们可以通过两次分组来实现:

解决方案代码

# 第一步:先按(id, duration)组合键分组,提取每组的action对列表
grouped_by_id_duration = df.rdd.groupBy(lambda x: (x['id'], x['duration'])) \
    .mapValues(lambda rows: [(row['action1'], row['action2']) for row in rows])

# 第二步:按id重新分组,整理成目标嵌套格式
final_result = grouped_by_id_duration.groupBy(lambda x: x[0][0]) \
    .mapValues(lambda duration_groups: tuple((dg[0][1], dg[1]) for dg in duration_groups)) \
    .map(lambda x: (f"id={x[0]}", x[1])) \
    .collect()

代码解释

  • 第一步:把每个元素的分组键设为(id, duration),这样同一ID且相同时长的记录会被聚合到一组,再用mapValues把组内的每条数据转换成(action1, action2)的元组,得到((id, duration), [action_tuple_list])格式的RDD。
  • 第二步:基于第一步的结果,再次按id(也就是组合键的第一个元素)分组,然后把每个时长组转换成(duration值, action列表)的元组,打包成元组集合,最后把外层的ID格式化成id=X的形式,就得到你想要的结构了。

运行结果

执行后你会得到如下格式的结果:

[('id=1', ((10, [('A', 'D'), ('B', 'E')]), (7, [('A', 'G')]))), ('id=2', ((25, [('A', 'E')]),))]

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:49:58