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

如何在Spark中合并两个文件并对相同id对应的数值求和

Spark 按id合并双文件并对数值求和实现方案

实现思路

  • 同时读取两个目标文本文件
  • 将每行内容拆分,转换为(id, 数值)的键值对结构,注意把数值从字符串转为浮点类型
  • 按id分组,对同id的数值执行求和计算
  • 整理结果格式后输出到文件

PySpark 代码示例

from pyspark.sql import SparkSession
from pyspark.sql.functions import split, sum, col

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

# 同时读取两个文件
df = spark.read.text(["file1.txt", "file2.txt"])

# 拆分列、转换类型,按id求和
result_df = df.select(
    split(col("value"), "\\s+").getItem(0).alias("id").cast("int"),
    split(col("value"), "\\s+").getItem(1).alias("num").cast("double")
).groupBy("id").agg(sum("num").alias("total_num")).orderBy("id")

# 拼接为空格分隔的字符串后输出,coalesce(1)可选,用于合并为单个输出文件
result_df.selectExpr("concat(id, ' ', total_num) as value").coalesce(1).write.text("result.txt")

spark.stop()

Spark Scala 代码示例

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{split, sum, col, concat, lit}

object MergeSumById {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder().appName("MergeSumById").master("local[*]").getOrCreate()
    
    // 读取两个文件
    val df = spark.read.text("file1.txt", "file2.txt")
    
    // 数据处理、按id求和
    val resultDf = df.select(
      split(col("value"), "\\s+")(0).cast("int").alias("id"),
      split(col("value"), "\\s+")(1).cast("double").alias("num")
    ).groupBy("id").agg(sum("num").alias("total_num")).orderBy("id")
    
    // 输出结果
    resultDf.select(concat(col("id"), lit(" "), col("total_num")).alias("value"))
      .coalesce(1)
      .write.text("result.txt")
      
    spark.stop()
  }
}

注意事项

  • Spark默认输出是result.txt目录,目录下名为part-00000的文件就是最终结果,和预期输出内容完全一致
  • 如果使用RDD API实现逻辑完全一致:读取文件后拆分为(id, 数值)键值对,调用reduceByKey(_+_)后排序输出即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 19:18:02