Spark按col1分组、依col2条件新增列,新增排序列后的实现需求
Spark DataFrame分组聚合及排序需求处理
初始需求
定义初始DataFrame
df = spark.createDataFrame( [ ['A', '1', '3'], ['A', '2', '7'], ['A', '3', '1'], ['A', '1', '5'], ['A', '3', '4'], ['A', '5', '2'], ['B', '1', '8'], ['B', '2', '4'], ['B', '4', '2'], ['B', '6', '8'] ], ['col1', 'col2', 'col3'] ) df.show()
需求
按col1分组,将每组内的连续行批次(A组前3行、后3行各为一批,B组4行为一批)的col2和col3分别聚合成数组,最终得到如下输出:
期望输出
+----+------------+------------+ |col1| col2| col3| +----+------------+------------+ | A| [1, 2, 3]| [3, 7, 1]| | A| [1, 3, 5]| [5, 4, 2]| | B|[1, 2, 4, 6]|[8, 4, 2, 8]| +----+------------+------------+
修改后需求
定义修改后的DataFrame
新增col4列用于划分行批次,col4值相同的行视为同一批次;若col4存在重复值,对应行的顺序无需严格保证:
df = spark.createDataFrame( [ ['A', '1', '3','2'], ['A', '2', '7','2'], ['A', '3', '1','2'], ['A', '1', '5','3'], ['A', '3', '4','3'], ['A', '5', '2','4'], ['B', '1', '8','4'], ['B', '2', '4','5'], ['B', '4', '2','6'], ['B', '6', '8','7'] ], ['col1', 'col2', 'col3', 'col4'] ) df.show()
需求
按col1分组,再结合col4的批次划分规则,将同一col1下、同一col4批次内的col2和col3分别聚合成数组。
内容的提问来源于stack exchange,提问作者Ivan Lee
相关产品推荐
相关产品推荐

