【发布时间】:2020-10-03 17:56:37
【问题描述】:
对 spark 和 RDD 来说非常新,所以我希望我能很好地解释我所追求的东西,以便有人理解和帮助:)
我有两组非常大的数据,比如说 300 万行和 50 列存储在 hadoop hdfs 中。 我想要做的是将这两个读入 RDD,以便它使用并行性 & 我想返回包含所有不匹配的记录(来自任一 RDD)的第三个 RDD。
希望下面有助于显示我想要做什么... 只是试图以最快最有效的方式找到所有不同的记录......
数据的顺序不一定相同 - rdd1 的第 1 行可能是 rdd2 的第 4 行。
提前非常感谢!!
所以...这似乎在做我想做的事,但似乎很容易正确...
%spark
import org.apache.spark.sql.DataFrame
import org.apache.spark.rdd.RDD
import sqlContext.implicits._
import org.apache.spark.sql._
//create the tab1 rdd.
val rdd1 = sqlContext.sql("select * FROM table1").withColumn("source",lit("tab1"))
//create the tab2 rdd.
val rdd2 = sqlContext.sql("select * FROM table2").withColumn("source",lit("tab2"))
//create the rdd of all misaligned records between table1 and the table2.
val rdd3 = rdd1.except(rdd2).unionAll(rdd2.except(rdd1))
//rdd3.printSchema()
//val rdd3 = rdd1.except(rdd2)
//drop the temporary table that was used to create a hive compatible table from the last run.
sqlContext.dropTempTable("table3")
//register the new temporary table.
rdd3.toDF().registerTempTable("table3")
//drop the old compare table.
sqlContext.sql("drop table if exists data_base.compare_table")
//create the new version of the s_asset compare table.
sqlContext.sql("create table data_base.compare_table as select * from table3")
这是我到目前为止完成的最后一段代码,它似乎正在完成这项工作 - 不确定完整数据集的性能,让我的手指交叉......
非常感谢所有花时间帮助这个可怜的人 :)
附言如果有人有性能更高的解决方案,我很想听听! 或者如果您可以看到一些问题,这可能意味着它会返回错误的结果。
【问题讨论】:
-
可以使用rdd.fullOuterJoin。
-
有什么特别的东西可以让它分布在所有节点上吗?
-
不是很大恕我直言
-
fullOuterJoin无法正常工作,您可能需要df1.leftOuter.df2 union df2.leftOuter.df1条件 -
哈哈@thebluephantom - 它实际上比这大得多,只是我没有想到实际的数字:)
标签: scala apache-spark pyspark rdd