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

