【发布时间】:2020-11-15 11:23:11
【问题描述】:
我有两个数据集df1 和df2,我需要在其中检测df2 与df1 中不同的任何记录,并创建一个结果数据集,其中包含一个标记不同记录的附加列。这是一个例子。
package playground
import org.apache.log4j.{Level, Logger}
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{col, sum}
object sample4 {
val spark = SparkSession
.builder()
.appName("Sample app")
.master("local")
.getOrCreate()
val sc = spark.sparkContext
final case class Owner(a: Long,
b: String,
c: Long,
d: Short,
e: String,
f: String,
o_qtty: Double)
final case class Result(a: Long,
b: String,
c: Long,
d: Short,
e: String,
f: String,
o_qtty: Double,
isDiff: Boolean)
def main(args: Array[String]): Unit = {
Logger.getLogger("org").setLevel(Level.OFF)
import spark.implicits._
val data1 = Seq(
Owner(11, "A", 666, 2017, "x", "y", 50),
Owner(11, "A", 222, 2018, "x", "y", 20),
Owner(33, "C", 444, 2018, "x", "y", 20),
Owner(33, "C", 555, 2018, "x", "y", 120),
Owner(22, "B", 555, 2018, "x", "y", 20),
Owner(99, "D", 888, 2018, "x", "y", 100),
Owner(11, "A", 888, 2018, "x", "y", 100),
Owner(11, "A", 666, 2018, "x", "y", 80),
Owner(33, "C", 666, 2018, "x", "y", 80),
Owner(11, "A", 444, 2018, "x", "y", 50),
)
val data2 = Seq(
Owner(11, "A", 666, 2017, "x", "y", 50),
Owner(11, "A", 222, 2018, "x", "y", 20),
Owner(33, "C", 444, 2018, "x", "y", 20),
Owner(33, "C", 555, 2018, "x", "y", 55),
Owner(22, "B", 555, 2018, "x", "y", 20),
Owner(99, "D", 888, 2018, "x", "y", 100),
Owner(11, "A", 888, 2018, "x", "y", 100),
Owner(11, "A", 666, 2018, "x", "y", 80),
Owner(33, "C", 666, 2018, "x", "y", 80),
Owner(11, "A", 444, 2018, "x", "y", 50),
)
val expected = Seq(
Result(11, "A", 666, 2017, "x", "y", 50, isDiff = false),
Result(11, "A", 222, 2018, "x", "y", 20, isDiff = false),
Result(33, "C", 444, 2018, "x", "y", 20, isDiff = false),
Result(33, "C", 555, 2018, "x", "y", 55, isDiff = true),
Result(22, "B", 555, 2018, "x", "y", 20, isDiff = false),
Result(99, "D", 888, 2018, "x", "y", 100, isDiff = false),
Result(11, "A", 888, 2018, "x", "y", 100, isDiff = false),
Result(11, "A", 666, 2018, "x", "y", 80, isDiff = false),
Result(33, "C", 666, 2018, "x", "y", 80, isDiff = false),
Result(11, "A", 444, 2018, "x", "y", 50, isDiff = false),
)
val df1 = spark
.createDataset(data1)
.as[Owner]
.cache()
val df2 = spark
.createDataset(data2)
.as[Owner]
.cache()
}
}
最有效的方法是什么?
【问题讨论】:
-
这能回答你的问题吗? Compare two Spark dataframes
-
不是真的,这个问题没有考虑隔离记录不同的标志。
-
我认为下面的 2 个答案效率不高,因为
join和intersect为所有记录和所有分区创建哈希表并进行比较。至少您可以尝试最简单的解决方案:df1.rdd.zip(df2.rdd).map {case (x,y) => (x, x != y)}并比较真实数据集的速度。 PS:将单字符字符串替换为char是个好主意,因为char比较非常快 -
@MikhailIonkin 您能否提供完整的答案,以便我接受您的答案。你在这里说得很好!
-
完成了。我没有真正的数据集,所以我不能证明我的答案更快,但我认为这是根据对小数据集的测试和文档
标签: sql scala apache-spark dataset