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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 15:00:47