【问题标题】:Effective record linkage有效的记录联动
【发布时间】:2021-11-30 13:32:23
【问题描述】:

我今天早些时候问了一个类似的问题。 Here 是。 很快:我需要为两个大型数据集(1.6M 和 6M)进行记录链接。我打算使用 Sparks,认为我被警告过的笛卡尔积不会是一个大问题。但它是。它对性能的打击如此之大,以至于联动过程没有在 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")
      }
    }

可能条件太复杂了,但一定是这样的。

【问题讨论】:

    标签: java scala apache-spark record-linkage


    【解决方案1】:

    您可以通过稍微更改执行链接的逻辑来提高当前解决方案的性能:

    • 首先对您知道匹配的列执行ab 数据帧的内连接。在您的情况下,它似乎是 LAST_NAMEFIRST_NAME 列。
    • 然后根据您的特定复杂条件过滤生成的数据框,在您的情况下,出生日期相等或街道匹配条件。
    • 最后,如果您还需要保留未链接的记录,请对a 数据框执行右连接

    你的代码可以改写如下:

    import org.apache.spark.sql.functions.{col, substring, to_date}
    import org.apache.spark.sql.SparkSession
    
    import java.time.Duration
    
    object App {
    
      def main(args: Array[String]): Unit = {
    
        val spark = SparkSession.builder()
          .master("local[4]")
          .appName("MatchingApp")
          .getOrCreate()
    
        val a = spark.read
          .format("csv")
          .option("header", true)
          .option("delimiter", ";")
          .load("/home/helveticau/workstuff/a.csv")
          .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("BIRTH_DATE", to_date(col("BIRTH_DATE"), "dd.MM.yyyy"))
    
        val condition = a.col("BIRTH_DATE").equalTo(b.col("BIRTH_DATE"))
          .or(a.col("STREET").startsWith(substring(b.col("STR"), 1, 3)))
    
        val startMillis = System.currentTimeMillis();
        val res = a.join(b, Seq("LAST_NAME", "FIRST_NAME"))
          .filter(condition)
          // two following lines optional if you want to only keep records with not null B_ID
          .select("B_ID", "A_ID")
          .join(a, Seq("A_ID"), "right_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")
      }
    }
    

    因此,您将以两个连接而不是一个连接的价格来避免笛卡尔积。

    示例

    文件a.csv包含以下数据:

    "A_ID";"FIRST_NAME";"LAST_NAME";"BIRTH_DATE";"STREET"
    10;John;Doe;1965-10-21;Johnson Road
    11;Rebecca;Davis;1977-02-27;Lincoln Road
    12;Samantha;Johns;1954-03-31;Main Street
    13;Roger;Penrose;1987-12-25;Oxford Street
    14;Robert;Smith;1981-08-26;Canergie Road
    15;Britney;Stark;1983-09-27;Alshire Road
    

    b.txt 具有以下数据:

    "B_ID";"FIRST_NAME";"LAST_NAME";"BIRTH_DATE";"STR"
    29;John;Doe;21.10.1965;Johnson Road
    28;Rebecca;Davis;28.03.1986;Lincoln Road
    27;Shirley;Iron;30.01.1956;Oak Street
    26;Roger;Penrose;25.12.1987;York Street
    25;Robert;Dayton;26.08.1956;Canergie Road
    24;Britney;Stark;22.06.1962;Algon Road
    

    res 数据框将是:

    +----+----+----------+---------+----------+-------------+
    |A_ID|B_ID|FIRST_NAME|LAST_NAME|BIRTH_DATE|STREET       |
    +----+----+----------+---------+----------+-------------+
    |10  |29  |John      |Doe      |1965-10-21|Johnson Road |
    |11  |28  |Rebecca   |Davis    |1977-02-27|Lincoln Road |
    |12  |null|Samantha  |Johns    |1954-03-31|Main Street  |
    |13  |26  |Roger     |Penrose  |1987-12-25|Oxford Street|
    |14  |null|Robert    |Smith    |1981-08-26|Canergie Road|
    |15  |null|Britney   |Stark    |1983-09-27|Alshire Road |
    +----+----+----------+---------+----------+-------------+
    

    注意:如果您的FIRST_NAMELAST_NAME 列不完全相同,您可以尝试使它们与Spark's built-in functions 匹配,例如:

    • trim 删除字符串开头和结尾的空格
    • lower 将列转换为小写(因此忽略大小写)

    真正重要的是拥有完全匹配的最大列数。

    【讨论】:

    • 一个不错的主意。虽然我的名字(LAST_NAME、FIRST_NAME)很难仅使用完全相等来比较。有些人在他们的姓名字段之一(有点 2 个姓名)中有 2-3 个单词,而在相应的列中只有一个单词。基本上是这样的:10|Johannes Diderik|van der Waals|1837-10-17|Random Road 在另一组中它可以更短。然后'包含'会起作用,但它很慢。
    猜你喜欢
    • 2016-05-10
    • 1970-01-01
    • 2022-12-18
    • 1970-01-01
    • 1970-01-01
    • 2014-08-05
    • 1970-01-01
    • 1970-01-01
    • 2020-09-20
    相关资源
    最近更新 更多