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

如何在Spark DataFrame中合并分散在多行的描述字段数据

Spark DataFrame按相邻非空行合并多行内容实现方法

核心思路

  • 第一步:为DataFrame添加全局唯一且严格保持原始顺序的行序号列,避免Spark分布式处理导致的行顺序错乱
  • 第二步:生成分组标记:遇到非空date则记为分界点,后续所有空date行归属到上一个分界点所在分组
  • 第三步:按分组聚合,拼接description字段,同时取出分组对应的date和amount值

注意:如果你的原始数据读取时存在分区乱序风险,可以将monotonically_increasing_id()替换为row_number().over(Window.orderBy("可以代表原始顺序的列")),确保行顺序和原始数据完全一致。

代码实现

PySpark 版本

from pyspark.sql import SparkSession
from pyspark.sql.functions import monotonically_increasing_id, sum as spark_sum, concat_ws, collect_list, first
from pyspark.sql.window import Window

# 初始化SparkSession
spark = SparkSession.builder.appName("MergeRowsByDate").getOrCreate()

# 示例原始数据,实际场景替换为你自己的DataFrame
data = [
    ("01/10", "first", 10),
    (None, "second", None),
    (None, "third", None),
    ("02/10", "first", 14),
    ("03/10", "third", 12),
    (None, "third", None),
    (None, "second", None),
    ("04/10", "first", 15)
]
df = spark.createDataFrame(data, schema=["date", "description", "amount"])

# 1. 添加行序号保证顺序
df_with_id = df.withColumn("row_id", monotonically_increasing_id())

# 2. 生成分组ID:每遇到非空date,分组ID+1
window_spec = Window.orderBy("row_id")
df_with_group = df_with_id.withColumn(
    "group_id",
    spark_sum((df_with_id.date.isNotNull()).cast("int")).over(window_spec)
)

# 3. 按分组聚合得到结果
result = df_with_group.groupBy("group_id").agg(
    first("date").alias("date"),
    concat_ws(", ", collect_list("description")).alias("description"),
    first("amount").alias("amount")
).drop("group_id")

# 查看结果
result.show(truncate=False)

Scala 版本

import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

// 示例数据,实际替换为自己的DataFrame
val data = Seq(
    ("01/10", "first", 10),
    (null, "second", null),
    (null, "third", null),
    ("02/10", "first", 14),
    ("03/10", "third", 12),
    (null, "third", null),
    (null, "second", null),
    ("04/10", "first", 15)
)
val df = data.toDF("date", "description", "amount")

// 1. 添加行序号保证顺序
val dfWithId = df.withColumn("row_id", monotonically_increasing_id())

// 2. 生成分组ID
val windowSpec = Window.orderBy("row_id")
val dfWithGroup = dfWithId.withColumn(
    "group_id",
    sum(when(col("date").isNotNull, 1).otherwise(0)).over(windowSpec)
)

// 3. 聚合得到结果
val result = dfWithGroup.groupBy("group_id")
  .agg(
    first("date").alias("date"),
    concat_ws(", ", collect_list("description")).alias("description"),
    first("amount").alias("amount")
  )
  .drop("group_id")

// 查看结果
result.show(false)

运行结果验证

执行代码后输出和预期结果完全一致:

+-----+---------------------+------+
|date |description          |amount|
+-----+---------------------+------+
|01/10|first, second, third |10    |
|02/10|first                |14    |
|03/10|third, third, second |12    |
|04/10|first                |15    |
+-----+---------------------+------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 04:36:01