Databricks中基于PySpark/SQL实现按Id分组排序聚合字符串
需求说明
我有一个包含Id、Value和Timestamp列的数据集,其中Id和Value为字符串类型。需要将相同Id对应的Value按Timestamp排序后用分号拼接。
示例输入数据
| Id | Value | Timestamp |
|---|---|---|
| Id1 | 100 | 1658919600 |
| Id1 | 200 | 1658919602 |
| Id1 | 300 | 1658919601 |
| Id2 | 433 | 1658919677 |
期望输出示例(以Id1为例)
| Id | Values |
|---|---|
| Id1 | 100;300;200 |
Databricks SQL 实现
直接使用STRING_AGG函数并指定排序规则即可,修正伪代码中的拼写错误(WITHING改为WITHIN):
SELECT Id, STRING_AGG(Value, ';' ORDER BY Timestamp) AS Values FROM your_table_name GROUP BY Id
PySpark 实现
提供两种实现方式,按需选择:
方法1:分组排序拼接
from pyspark.sql import functions as F df = spark.table("your_table_name") # 分组后收集(Timestamp, Value)结构体,排序后提取Value再拼接 result_df = df.groupBy("Id") \ .agg( F.concat_ws( ";", F.sort_array(F.collect_list(F.struct("Timestamp", "Value"))).Value ).alias("Values") ) result_df.show()
方法2:窗口函数先排序再分组
from pyspark.sql import functions as F from pyspark.sql.window import Window df = spark.table("your_table_name") # 窗口内按Timestamp排序并收集Value window_spec = Window.partitionBy("Id").orderBy("Timestamp") df_ordered = df.withColumn("ordered_values", F.collect_list("Value").over(window_spec)) # 分组取每组的完整排序列表并拼接 result_df = df_ordered.groupBy("Id") \ .agg(F.max("ordered_values").alias("ordered_values")) \ .withColumn("Values", F.concat_ws(";", "ordered_values")) \ .drop("ordered_values") result_df.show()
内容的提问来源于stack exchange,提问作者Alcibiades
相关产品推荐
相关产品推荐

