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

求Spark-SQL查询语句:按唯一column1及时间升序合并多行数据

Spark-SQL 分组按时间合并状态实现方案

需求说明

现有enrollment表数据如下:

column1 column2 timeStamp
abc enrolled 2022/09/01
abc changed 2022/09/02
abc registered 2022/09/04
abc blocked 2022/09/05
abc left 2022/09/06
def enrolled 2022/09/20
def changed 2022/09/21
def changed 2022/09/21
def changed 2022/09/24
def left 2022/09/25
ghi registered 2022/09/01
ghi changed 2022/09/02
ghi left 2022/09/03
ghi returned 2022/10/03

需要编写Spark-SQL语句,按column1分组,将每组内的column2按timeStamp升序用-拼接成一行,输出格式如下:

out_column1 out_column2
abc enrolled-changed-registered-blocked-left
def enrolled-changed-changed-left
ghi registered-changed-left-returned

注:PySpark中无group_concat函数。

实现语句

通过struct、collect_list、sort_array和concat_ws组合实现需求,SQL语句如下:

SELECT
    column1 AS out_column1,
    CONCAT_WS('-', TRANSFORM(sort_array(COLLECT_LIST(status_struct)), x -> x.column2)) AS out_column2
FROM (
    SELECT
        column1,
        STRUCT(TO_DATE(timeStamp, 'yyyy/MM/dd') AS ts, column2) AS status_struct
    FROM enrollment
) t
GROUP BY column1;

语句解释

  1. 内层查询:将字符串类型的timeStamp通过TO_DATE转换为日期类型(避免字符串排序的潜在问题),和column2打包成struct类型字段status_struct——结构体中的日期字段用于后续排序,column2是要合并的状态值。
  2. 分组收集:按column1分组,调用COLLECT_LIST收集每组内所有的status_struct,得到一个结构体数组。
  3. 排序拼接:
    • 用sort_array对结构体数组按日期字段ts升序排序,保证状态顺序和时间线一致;
    • 通过TRANSFORM遍历排序后的数组,提取每个结构体中的column2值,得到纯状态值的数组;
    • 最后用CONCAT_WS将状态数组用-拼接成完整的字符串。

如果你的timeStamp字段本身就是日期类型,直接使用timeStamp作为结构体的排序字段即可,无需TO_DATE转换。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 17:01:10