【发布时间】:2021-09-04 04:10:24
【问题描述】:
我有两个文件:第一列代表id,第二列代表值。
file1.txt 如下:
1 0.1
2 0.2
3 0.3
file2.txt 如下:
1 0.4
2 0.3
3 0.1
如何在spark中合并两个文件并添加对应的值? result.txt 如下:
1 0.5
2 0.5
3 0.4
【问题讨论】:
标签: scala apache-spark
我有两个文件:第一列代表id,第二列代表值。
file1.txt 如下:
1 0.1
2 0.2
3 0.3
file2.txt 如下:
1 0.4
2 0.3
3 0.1
如何在spark中合并两个文件并添加对应的值? result.txt 如下:
1 0.5
2 0.5
3 0.4
【问题讨论】:
标签: scala apache-spark
您可以在id 上通过Union 2 DataFrames 和GroupBy 生成value 列的总和 -
input_list1 = [
(1,0.1)
,(2,0.2)
,(3,0.3)
]
input_list2 = [
(1,0.4)
,(2,0.3)
,(3,0.1)
]
sparkDF1 = sql.createDataFrame(input_list1, ['id','value'])
sparkDF2 = sql.createDataFrame(input_list2, ['id','value'])
sparkDF1.show()
+---+-----+
| id|value|
+---+-----+
| 1| 0.1|
| 2| 0.2|
| 3| 0.3|
+---+-----+
sparkDF2.show()
+---+-----+
| id|value|
+---+-----+
| 1| 0.4|
| 2| 0.3|
| 3| 0.1|
+---+-----+
combinedDF = sparkDF1.union(sparkDF2)
combinedDF.show()
+---+-----+
| id|value|
+---+-----+
| 1| 0.1|
| 2| 0.2|
| 3| 0.3|
| 1| 0.4|
| 2| 0.3|
| 3| 0.1|
+---+-----+
combinedDF.groupBy('id').agg(F.sum(F.col('value')).alias('value')).show()
+---+-----+
| id|value|
+---+-----+
| 1| 0.5|
| 3| 0.4|
| 2| 0.5|
+---+-----+
【讨论】:
您可以将两个文本文件读取为 csv 文件,然后加入两个读取的数据框,将两个值列相加,然后将结果数据框保存为 csv 文件。
import org.apache.spark.sql.functions.col
import org.apache.spark.sql.types.{DoubleType, IntegerType, StructField, StructType}
// Define schema of the two input files
val schema1 = StructType(Seq(
StructField("id", IntegerType),
StructField("value1", DoubleType)
))
val schema2 = StructType(Seq(
StructField("id", IntegerType),
StructField("value2", DoubleType)
))
// Read the two input files to two dataframes
val dataframe1 = sparkSession
.read
.option("delimiter", " ")
.option("header", "false")
.schema(schema1)
.csv("file1.txt")
val dataframe2 = sparkSession
.read
.option("delimiter", " ")
.option("header", "false")
.schema(schema2)
.csv("file2.txt")
// Join the two dataframes and save result as one csv file
dataframe1.join(dataframe2, Seq("id"))
.withColumn("value", col("value1") + col("value2"))
.drop("value1", "value2")
.repartition(1)
.write
.mode("overwrite")
.option("header", "false")
.option("delimiter", " ")
.csv("/tmp/result")
但是,您不能选择输出文件的名称。所以如果你真的想要一个result.txt 文件,你应该在之后重命名它:
mv /tmp/result/*.csv /tmp/result/result.txt
【讨论】: