基于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
相关产品推荐
相关产品推荐

