如何在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
相关产品推荐
相关产品推荐

