1.6M与6M数据集Record Linkage实现及Spark关联方案性能优化求助
我今日早些时候问过一个类似问题。
需求简述:我需要对两个规模分别为160万、600万的大型数据集做记录关联。最初选用Spark实现,原本以为此前被提醒的笛卡尔积性能问题影响不大,但实际性能影响极为严重,关联流程运行7小时仍未完成。
请问是否有其他更高效的库/框架/工具可实现该需求?或者有没有办法优化下述解决方案的性能?
我最终使用的代码如下:
object App { def left(col: Column, n: Int) = { assert(n > 0) substring(col, 1, n) } def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .master("local[4]") .appName("MatchingApp") .getOrCreate() import spark.implicits._ val a = spark.read .format("csv") .option("header", true) .option("delimiter", ";") .load("/home/helveticau/workstuff/a.csv") .withColumn("FULL_NAME", concat_ws(" ", col("FIRST_NAME"), col("LAST_NAME"))) .withColumn("BIRTH_DATE", to_date(col("BIRTH_DATE"), "yyyy-MM-dd")) val b = spark.read .format("csv") .option("header", true) .option("delimiter", ";") .load("/home/helveticau/workstuff/b.txt") .withColumn("FULL_NAME", concat_ws(" ", col("FIRST_NAME"), col("LAST_NAME"))) .withColumn("BIRTH_DATE", to_date(col("BIRTH_DATE"), "dd.MM.yyyy")) // @formatter:off val condition = a .col("FULL_NAME").contains(b.col("FIRST_NAME")) .and(a.col("FULL_NAME").contains(b.col("LAST_NAME"))) .and(a.col("BIRTH_DATE").equalTo(b.col("BIRTH_DATE")) .or(a.col("STREET").startsWith(left(b.col("STR"), 3)))) // @formatter:on val startMillis = System.currentTimeMillis(); val res = a.join(b, condition, "left_outer") val count = res .filter(col("B_ID").isNotNull) .count() println(s"Count: $count") val executionTime = Duration.ofMillis(System.currentTimeMillis() - startMillis) println(s"Execution time: ${executionTime.toMinutes}m") } }
当前关联条件逻辑虽然较为复杂,但出于业务要求必须保留,无法调整。
内容的提问来源于stack exchange,提问作者Sebastian Lore
相关产品推荐
相关产品推荐

