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

如何在Azure Databricks的PySpark中每6行拼接列值?

问题:Azure Databricks + PySpark 3.4.1 实现组内按窗口拼接列值

我目前使用Azure Databricks及PySpark 3.4.1,需要实现按分组内每6行(不足6行则取全部)以空格分隔拼接col4的值,最终让组内所有行都对应拼接后的完整结果(如示例中第5列所示)。但之前的尝试没有成功。

数据集示例

data = [
    ("1000","999", "1", "123", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "2", "abc", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "3", "456", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "4", "def", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "5", "789", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "6", "ghi", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "7", "101112", "101112 jkl 131415"),
    ("1000","999", "8", "jkl", "101112 jkl 131415"),
    ("1000","999", "9", "131415", "101112 jkl 131415"),
    ("1001","999", "1", "aaa", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "2", "bbb", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "3", "ccc", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "4", "ddd", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "5", "eee", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "6", "fff", "aaa bbb ccc ddd eee fff"),
    ("1002","999", "1", "data", "data"),
    ("1003","900", "1", "stack", "stack over flow .com"),
    ("1003","900", "2", "over", "stack over flow .com"),
    ("1003","900", "3", "flow", "stack over flow .com"),
    ("1003","900", "4", ".com", "stack over flow .com"),
    ("1003","997", "1", "py", "py spark databricks"),
    ("1003","997", "2", "spark", "py spark databricks"),
    ("1003","997", "3", "databricks", "py spark databricks"),
    ("1003","998", "1", "data science", "data science data engineering"),
    ("1003","998", "2", "data engineering", "data science data engineering"),
    ("1003","999", "1", "azure", "azure aws gcp"),
    ("1003","999", "2", "aws", "azure aws gcp"),
    ("1003","999", "3", "gcp", "azure aws gcp")
]

注:第5列为期望输出结果,需将同一分组内的col4值按顺序拼接,每6个为一组(不足6个则取全部),且分组内所有行对应同一拼接结果。

尝试过的代码

window_spec = Window().partitionBy("col1").orderBy("col3").rowsBetween(Window.currentRow - 5, Window.currentRow)

df = df.withColumn("col5", when(col("col3") == 6, concat_ws(" ", collect_list("col2").over(window_spec))))

解决方案

核心问题分析

之前的代码存在两个关键问题:

  1. 分区维度错误:需要按col1+col2联合分区(从数据可见,同一col1下不同col2属于独立分组);
  2. 窗口逻辑错误:不需要滑动窗口,而是要将分组内的行按每6个为一组划分,然后组内所有行关联对应分组的拼接结果。

方法一:PySpark API实现

步骤:

  1. 按col1、col2分区,col3排序,计算每个行所在的6行组编号;
  2. 按分区+组编号聚合,拼接col4的值;
  3. 将聚合结果关联回原表,得到所有行的对应拼接值。

代码示例:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 1. 创建初始DataFrame
df = spark.createDataFrame(data, ["col1", "col2", "col3", "col4", "expected_col5"])

# 2. 计算每个行所属的6行组编号(整数除法,每6行一组)
window_part = Window.partitionBy("col1", "col2").orderBy("col3")
df = df.withColumn("group_id", (F.row_number().over(window_part) - 1) // 6)

# 3. 按分区+组编号聚合,拼接col4的值
grouped_df = df.groupBy("col1", "col2", "group_id") \
               .agg(F.concat_ws(" ", F.collect_list("col4").orderBy("col3")).alias("col5"))

# 4. 关联回原表,得到最终结果
result_df = df.join(grouped_df, on=["col1", "col2", "group_id"], how="left") \
              .select("col1", "col2", "col3", "col4", "col5", "expected_col5")

result_df.show(truncate=False)

方法二:SparkSQL实现

先注册临时视图,再通过SQL语句完成逻辑:

-- 1. 注册临时视图
CREATE OR REPLACE TEMP VIEW source_data AS
SELECT * FROM VALUES
    ("1000","999", "1", "123", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "2", "abc", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "3", "456", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "4", "def", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "5", "789", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "6", "ghi", "123 abc 456 def 789 ghi 101112 jkl 131415"),
    ("1000","999", "7", "101112", "101112 jkl 131415"),
    ("1000","999", "8", "jkl", "101112 jkl 131415"),
    ("1000","999", "9", "131415", "101112 jkl 131415"),
    ("1001","999", "1", "aaa", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "2", "bbb", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "3", "ccc", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "4", "ddd", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "5", "eee", "aaa bbb ccc ddd eee fff"),
    ("1001","999", "6", "fff", "aaa bbb ccc ddd eee fff"),
    ("1002","999", "1", "data", "data"),
    ("1003","900", "1", "stack", "stack over flow .com"),
    ("1003","900", "2", "over", "stack over flow .com"),
    ("1003","900", "3", "flow", "stack over flow .com"),
    ("1003","900", "4", ".com", "stack over flow .com"),
    ("1003","997", "1", "py", "py spark databricks"),
    ("1003","997", "2", "spark", "py spark databricks"),
    ("1003","997", "3", "databricks", "py spark databricks"),
    ("1003","998", "1", "data science", "data science data engineering"),
    ("1003","998", "2", "data engineering", "data science data engineering"),
    ("1003","999", "1", "azure", "azure aws gcp"),
    ("1003","999", "2", "aws", "azure aws gcp"),
    ("1003","999", "3", "gcp", "azure aws gcp")
AS t(col1, col2, col3, col4, expected_col5);

-- 2. 计算组编号并关联拼接结果
WITH data_with_group AS (
    SELECT 
        *,
        (ROW_NUMBER() OVER (PARTITION BY col1, col2 ORDER BY col3) - 1) // 6 AS group_id
    FROM source_data
),
grouped_data AS (
    SELECT 
        col1, col2, group_id,
        CONCAT_WS(' ', COLLECT_LIST(col4) ORDER BY col3) AS col5
    FROM data_with_group
    GROUP BY col1, col2, group_id
)
SELECT 
    d.col1, d.col2, d.col3, d.col4, g.col5, d.expected_col5
FROM data_with_group d
LEFT JOIN grouped_data g 
ON d.col1 = g.col1 AND d.col2 = g.col2 AND d.group_id = g.group_id
ORDER BY d.col1, d.col2, d.col3;

结果验证

两种方法都会生成与示例中expected_col5完全一致的col5列,实现按每6行窗口拼接的需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 00:27:33