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

基于Arrow与DuckDB实现大数据集精确+模糊联合查询

问题描述

我有一个4300万行、3.84GB的大型数据集(first_data),以及一个6000行、459KB的小型数据集(second_data)。需要执行内连接,匹配条件为:

  1. 基于id列的精确匹配
  2. 基于fullname列的模糊匹配

此前尝试用stringdist_join实现,但出现内存溢出问题,代码如下:

fuzzy_join<-stringdist_join(first_data, second_data, 
                by = "fullname",
                mode = "inner",
                ignore_case = FALSE, 
                method = "jw", 
                max_dist = 0.5) %>% 
  filter(id == id)

两个数据集的前10行结构:

first_data<- structure(list(user_id = c(441391106, 441514065, 442060539, 442158489, 
438197192, 438206034, 438689594, 438881971, 440386286, 440479235
), fullname = c("Siva Kumar", "Kalyani M", "James Bigler", "Arthur Stephens", 
"guy guy", "Rick Schlieper", "Tony Klemencic", "baiyu xu", "Michael Fritts", 
"Daniel Wolf Roemele"), f_prob = c(0, 1, 0.005, 0.006, 0.005, 
0.002, 0.011, 0.389, 0.005, 0.004), m_prob = c(1, 0, 0.995, 0.994, 
0.995, 0.998, 0.989, 0.611, 0.995, 0.996), white_prob = c(0.021, 
0.001, 0.994, 0.792, 0.547, 0.949, 0.948, 0.001, 0.995, 0.795
), black_prob = c(0.013, 0.003, 0.001, 0.198, 0.398, 0.004, 0.003, 
0.001, 0.001, 0.097), api_prob = c(0.904, 0.991, 0, 0, 0.001, 
0.002, 0.003, 0.994, 0.001, 0.061), hispanic_prob = c(0.005, 
0.001, 0.001, 0.002, 0, 0.001, 0.039, 0.001, 0, 0.012), native_prob = c(0.006, 
0.002, 0, 0, 0, 0.005, 0, 0, 0, 0.003), multiple_prob = c(0.051, 
0.002, 0.004, 0.008, 0.054, 0.039, 0.007, 0.003, 0.003, 0.032
), degree = c("", "", "Bachelor", "", "", "Master", "Associate", 
"", "", ""), other_id = c(1212616, 1212616, 1212616, 1212616, 1212991, 
1212991, 1212991, 1212991, 1212991, 1212991), id = c(62399, 
62399, 62399, 62399, 63907, 63907, 63907, 63907, 63907, 63907
)), row.names = c(NA, 10L), class = "data.frame")

second_data<- structure(list(id = c(12825, 12945, 12945, 12945, 16456, 16456, 
16456, 12136, 12136, 17254), another_id = c(7879, 8587, 18070, 40634, 
13142, 17440, 41322, 899, 27199, 26604), fname = c("Gerald", 
"John", "Dean", "Todd", "Thomas", "Ivan", "Vinit", "Scott", "Jonathan", 
"William"), mname = c("B.", "L.", "A.", "P.", "F.", "R.", 
"K.", "G.", "I.", "Jensen"), lname = c("Shreiber", "Nussbaum", 
"Foate", "Kelsey", "Kirk", "Sabel, CPO", "Asar", "McNealy", "Schwartz", 
"Gedwed"), companyname = c(NA, "Plexus Corp.", "Plexus Corp.", 
NA, NA, NA, NA, "Oracle America, Inc.", "Oracle America, Inc.", 
NA), fullname = c("Gerald Shreiber", "John Nussbaum", 
"Dean Foate", "Todd Kelsey", "Thomas Kirk", "Ivan Sabel", "Vinit Asar", 
"Scott McNealy", "Jonathan Schwartz", "William Gedwed")), row.names = c(NA, 
-10L), class = c("tbl_df", "tbl", "data.frame"))

需要用DuckDB或Arrow实现同时满足精确+模糊的内连接,解决内存占用问题。


解决方案

利用DuckDB或Arrow的磁盘外处理能力,先按id精确分区,再在分区内做模糊匹配,减少单批次处理的数据量。

方案1:使用DuckDB

DuckDB支持直接处理磁盘数据集(如Parquet),也可将数据导入内存数据库后执行SQL查询,内置jaro_winkler/levenshtein函数实现模糊匹配。

操作步骤:

  1. 安装加载依赖包
install.packages("duckdb")
library(duckdb)
library(dplyr)
  1. 将数据导入DuckDB实例
# 连接内存中的DuckDB
con <- dbConnect(duckdb())

# 写入数据集到数据库
dbWriteTable(con, "first_data", first_data)
dbWriteTable(con, "second_data", second_data)
  1. 执行SQL查询(精确+模糊匹配)
result <- dbGetQuery(con, "
  SELECT f.*, s.*
  FROM first_data f
  INNER JOIN second_data s
    ON f.id = s.id
    AND jaro_winkler(f.fullname, s.fullname) >= 0.5
")

# 关闭连接
dbDisconnect(con, shutdown = TRUE)

注:jaro_winkler返回相似度(0-1),与stringdist的距离逻辑相反,因此用>=0.5对应原代码的max_dist=0.5;若用Levenshtein距离,替换为levenshtein(f.fullname, s.fullname) <= [阈值]即可。

方案2:使用Arrow

Arrow支持大数据集的分区处理和矢量化操作,结合fuzzyjoin分批次执行模糊匹配。

操作步骤:

  1. 安装加载依赖包
install.packages("arrow")
library(arrow)
library(dplyr)
library(fuzzyjoin)
  1. 转换大型数据集为Arrow Dataset
# 将大型数据集转为Arrow Dataset(支持磁盘存储)
first_dataset <- arrow_table(first_data) %>% 
  to_dataset()

# 小型数据集保留为data.frame
second_df <- as.data.frame(second_data)
  1. 按id分组执行模糊匹配
# 按id拆分大型数据集,逐个批次匹配
result <- first_dataset %>%
  group_by(id) %>%
  group_map(function(group, key) {
    current_id <- key$id
    # 过滤小型数据集中同id的行
    second_subset <- filter(second_df, id == current_id)
    
    if(nrow(second_subset) == 0) return(tibble())
    
    # 执行模糊连接
    stringdist_join(
      as.data.frame(group),
      second_subset,
      by = "fullname",
      mode = "inner",
      ignore_case = FALSE,
      method = "jw",
      max_dist = 0.5
    )
  }) %>%
  bind_rows()

注:通过group_by(id)拆分大型数据集,每个批次仅与对应id的小型子集匹配,大幅降低内存占用;若原始数据为Parquet格式,可直接按id分区存储,进一步提升效率。


优化建议

  • 优先将大型数据集存储为Parquet格式,DuckDB和Arrow对Parquet的处理效率远高于CSV。
  • 调整模糊匹配阈值,过滤不必要的匹配,减少计算量。
  • 若id基数较高,可进一步优化批次大小,避免单个批次数据量过大。

内容的提问来源于stack exchange,提问作者bear_525

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 04:57:02