求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;
语句解释
- 内层查询:将字符串类型的
timeStamp通过TO_DATE转换为日期类型(避免字符串排序的潜在问题),和column2打包成struct类型字段status_struct——结构体中的日期字段用于后续排序,column2是要合并的状态值。 - 分组收集:按
column1分组,调用COLLECT_LIST收集每组内所有的status_struct,得到一个结构体数组。 - 排序拼接:
- 用
sort_array对结构体数组按日期字段ts升序排序,保证状态顺序和时间线一致; - 通过
TRANSFORM遍历排序后的数组,提取每个结构体中的column2值,得到纯状态值的数组; - 最后用
CONCAT_WS将状态数组用-拼接成完整的字符串。
- 用
如果你的timeStamp字段本身就是日期类型,直接使用timeStamp作为结构体的排序字段即可,无需TO_DATE转换。
内容的提问来源于stack exchange,提问作者hellopa1
相关产品推荐
相关产品推荐

