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

如何对单字符串列的表进行每两行聚合操作?

实现方法说明

你需要将单列数据按固定行数分组合并,flatten函数不适用的原因是它仅用于展平嵌套数组,和当前分组聚合的需求不匹配。以下是两种符合你需求的实现方式:

1. 准备测试数据

先通过以下代码创建测试DataFrame,方便验证逻辑:

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

spark = SparkSession.builder.appName("group_names").getOrCreate()

data = [("Alice",), ("Bob",), ("Carol",), ("David",)]
df = spark.createDataFrame(data, ["name"])

2. 合并为逗号分隔字符串(对应第一个期望结果)

通过生成分组ID将每2行划分为一组,再聚合合并为字符串:

# 生成行号并计算分组ID
df_with_group = df.withColumn("row_num", F.row_number().over(Window.orderBy(F.lit(1)))) \
                 .withColumn("group_id", (F.col("row_num") - 1) // 2)

# 分组聚合生成逗号分隔的names列
result_str = df_with_group.groupBy("group_id") \
                          .agg(F.concat_ws(", ", F.collect_list("name")).alias("names")) \
                          .drop("group_id")

result_str.show(truncate=False)

输出结果:

+------------+
|names       |
+------------+
|Alice, Bob  |
|Carol, David|
+------------+

3. 合并为数组(对应第二个期望结果)

复用分组ID逻辑,聚合时直接生成数组类型的列:

# 分组聚合生成数组类型的names列
result_array = df_with_group.groupBy("group_id") \
                            .agg(F.collect_list("name").alias("names")) \
                            .drop("group_id")

result_array.show(truncate=False)

输出结果:

+----------------+
|names           |
+----------------+
|[Alice, Bob]    |
|[Carol, David]  |
+----------------+

补充说明

  • 代码中Window.orderBy(F.lit(1))用于保留数据的原有输入顺序,若你的表有明确排序依据(比如自增ID列),可替换为对应列以保证分组逻辑的稳定性。
  • 调整(F.col("row_num") - 1) // 2中的除数,可修改每组包含的行数,比如改为3即可实现每3行一组。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 23:22:38