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

Spark中如何迭代合并两行生成新表(每两行对应一行)

Spark实现两行分组生成新表方案

核心思路

先按id排序保证行序正确,通过行号将每两行划分为一组,再对每组聚合提取所需字段,最后按指定格式拼接结果。

步骤与代码实现

1. 排序并添加分组标识

先给原表按id升序排序,生成连续行号后,通过行号计算分组键(每两行一组)。

2. 分组聚合生成目标字段

对每个分组提取第一行的id作为新表id,格式化两行id为两位补零格式,取第二行的col b和message,最后拼接成指定格式的full message。

Scala版本代码

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

// 原表DataFrame假设为df
val windowSpec = Window.orderBy("id")
val groupedDf = df.withColumn("row_num", row_number().over(windowSpec))
  .withColumn("group_id", (col("row_num") - 1) / 2)

// 分组聚合并生成结果
val resultDf = groupedDf.groupBy("group_id")
  .agg(
    first("id").alias("id"),
    concat(
      lpad(first("id"), 2, "0"),
      lit(":"),
      lpad(last("id"), 2, "0")
    ).alias("id_part"),
    last("col b").alias("col_b"),
    last("message").alias("second_message")
  )
  .withColumn("full message", concat_ws(",", col("id_part"), col("col_b"), col("second_message")))
  .drop("group_id", "id_part", "col_b", "second_message")
  .orderBy("id")

// 查看结果
resultDf.show()

Python版本代码

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

# 原表DataFrame假设为df
window_spec = Window.orderBy("id")
grouped_df = df.withColumn("row_num", F.row_number().over(window_spec)) \
    .withColumn("group_id", (F.col("row_num") - 1) // 2)

# 分组聚合并生成结果
result_df = grouped_df.groupBy("group_id") \
    .agg(
        F.first("id").alias("id"),
        F.concat(
            F.lpad(F.first("id"), 2, "0"),
            F.lit(":"),
            F.lpad(F.last("id"), 2, "0")
        ).alias("id_part"),
        F.last("col b").alias("col_b"),
        F.last("message").alias("second_message")
    ) \
    .withColumn("full message", F.concat_ws(",", F.col("id_part"), F.col("col_b"), F.col("second_message"))) \
    .drop("group_id", "id_part", "col_b", "second_message") \
    .orderBy("id")

# 查看结果
result_df.show()

关键注意事项

  • 必须先按id排序:否则行号混乱,分组结果不符合预期
  • id格式化:使用lpad函数将id补为两位,不足时前面补0,确保输出如01、03的格式
  • 处理奇数行场景:如果原表总行数为奇数,最后一组仅一行,可添加过滤条件保留仅含两行的分组,示例如下(Scala):
val resultDf = groupedDf.groupBy("group_id")
  .agg(
    count("id").alias("row_count"),
    first("id").alias("id"),
    concat(lpad(first("id"),2,"0"), lit(":"), lpad(last("id"),2,"0")).alias("id_part"),
    last("col b").alias("col_b"),
    last("message").alias("second_message")
  )
  .filter(col("row_count") === 2)
  .withColumn("full message", concat_ws(",", col("id_part"), col("col_b"), col("second_message")))
  .drop("group_id", "row_count", "id_part", "col_b", "second_message")
  .orderBy("id")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 08:40:26