You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 22:15:02