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

基于sparklyr的数据集全量样本最近邻查询方案优化问询

基于sparklyr批量计算全样本最近邻的优化方案

原有方案的问题

你的判断是正确的,用lapply遍历调用ml_approx_nearest_neighbors()的方案效率极低:

  • 每一次遍历调用都会触发一次独立的Spark作业,700次调用就需要700次任务调度,开销远大于单次作业
  • 需要先把全量特征从Spark集群拉取到R本地,再逐个传回Spark计算,产生大量不必要的跨节点IO
  • 数据量如果进一步上升,这种写法很容易出现内存溢出、作业超时的问题

优化实现思路

直接使用LSH模型自带的ml_approx_similarity_join()方法完成批量最近邻查询,该方法可以在单次Spark作业内完成两个数据集全量样本的近似相似配对,不需要遍历调用,性能提升非常明显。

优化后代码

你原有代码中数据集预处理、向量组装、LSH模型训练的部分可以完全保留,只需替换最后拉取特征和lapply遍历的部分即可:

# 拆分查询集(前700个样本)和候选集(全量待匹配样本),重命名列避免后续join冲突
sdf_query <- sdf_titanic_va %>%
  dplyr::filter(id %in% 1:700) %>%
  dplyr::select(query_id = id, query_features = features)

sdf_candidate <- sdf_titanic_va %>%
  dplyr::select(candidate_id = id, candidate_features = features)

# 执行批量近似相似度连接,仅返回距离小于阈值的配对
knn_pairs <- sparklyr::ml_approx_similarity_join(
  x = brp_fit,
  dataset_a = sdf_query,
  dataset_b = sdf_candidate,
  dist_col = "dist_col",
  threshold = 100 # 可根据你的特征值域调整,只要能覆盖需要的近邻范围即可
)

# 按距离排序,取每个查询样本的前2个最近邻(包含自身)
knn_result <- knn_pairs %>%
  dplyr::group_by(query_id) %>%
  dplyr::arrange(dist_col, .by_group = TRUE) %>%
  dplyr::filter(dplyr::row_number() <= 2) %>%
  dplyr::ungroup()

# 如果需要关联回原始属性字段,直接和原表join即可
full_result <- knn_result %>%
  dplyr::left_join(sdf_titanic, by = c("query_id" = "id")) %>%
  dplyr::left_join(sdf_titanic, by = c("candidate_id" = "id"), suffix = c("_query", "_neighbor"))

额外优化建议

  • 可以在向量组装完成后调用sdf_titanic_va <- sdf_persist(sdf_titanic_va)将特征表缓存到Spark内存,避免重复计算特征,进一步提升速度
  • 如果不需要把自身作为最近邻,可在过滤行号前添加dplyr::filter(query_id != candidate_id),再取top1即可得到排除自身的最近邻
  • 可根据精度要求调整ft_bucketed_random_projection_lsh()的bucket_length和num_hash_tables参数:哈希表数量越多、桶长越小,查询精度越高但计算速度越慢

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 10:54:02