PySpark中实现多层groupBy嵌套分组的技术咨询
实现RDD双层分组(ID内按Duration再分组)
首先先确认你的原始数据:
| id | duration | action1 | action2 |
|---|---|---|---|
| 1 | 10 | A | D |
| 1 | 10 | B | E |
| 2 | 25 | A | E |
| 1 | 7 | A | G |
你已经完成了按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
相关产品推荐
相关产品推荐

