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

