
如何在 Slurm 集群上运行 Faiss 分布式 k-means 训练大量质心【免费下载链接】faissA library for efficient similarity search and clustering of dense vectors.项目地址: https://gitcode.com/GitHub_Trending/fa/faiss当你需要给上亿向量级的训练集聚出百万级质心例如用 IVF 索引切 100 万个倒排表时单机 k-means 不够用。Faiss 仓库的benchs/distributed_ondisk/目录提供了一套可直接在 Slurm 集群上跑的分布式 k-means 实现把训练向量按行切分到多台机器上由各自节点上的 Faiss 索引CPU 或 GPU完成到质心的 assignment主节点把各节点结果合成后更新质心。本文给出从本地 sanity check 到提交 Slurm 作业、并验证质心输出文件的完整操作路径。分布式 k-means 的工作方式与相关文件这套实现基于两个核心文件说明见 READMEdistributed_kmeans.pyk-means 主循环的 Python 实现紧跟 Faiss C 版本计算大头是 assignment由一个DatasetAssign抽象对象完成。该对象可以是本地 Faiss CPU 索引、GPU 索引或一组通过 RPC 连接的远端 CPU/GPU 索引DatasetAssignDispatch。run_on_cluster.bash在集群上以 Slurm 方式调度上述脚本的 shell 入口通过第一个参数todo选择执行哪一步test_kmeans_0/1/2、slurm_distributed_kmeans、deep1b_clustering等。运行前提README 明确列出Python 3不兼容 Python 2Python 环境里装好faiss和scipy稀疏矩阵用训练数据可以是 fvecs、bvecs 或 npy 格式向量规模大没关系读取走 memory-mapping数据文件必须能被所有参与节点访问例如通过分布式文件系统。先做本地 sanity check提交集群任务前可以先在单机上验证脚本能跑通。distributed_kmeans.py内置了 4 级测试用--test指定数字可逗号分隔# reference Faiss C run python distributed_kmeans.py --test 0 # using the Python implementation python distributed_kmeans.py --test 1 # use the dispatch object (on local datasets) python distributed_kmeans.py --test 2 # same, with GPUs python distributed_kmeans.py --test 3测试数据取自脚本内的testdata路径默认是/datasets01_101/simsearch/041218/bigann/bigann_learn.bvecs可以按 README 的提示把它改成本地数据副本如果该路径不存在脚本会自动退回SyntheticDataset(128, 100000, 0, 0)合成数据并打印 using synthetic dataset。README 说明输出应对照作者给出的示例 gist 判断是否正常具体输出样例链接见 README 的 Local tests 一节。如果要在一台多卡机器上模拟分布式流程每个 GPU 起一个 server通过 rpc 协议连接可以用# non distributed baseline bash run_on_cluster.bash test_kmeans_0 # using all the machines GPUs bash run_on_cluster.bash test_kmeans_1 # distributed run, with one local server per GPU bash run_on_cluster.bash test_kmeans_2test_kmeans_2会在本机为每块 GPU 各起一个 server 进程端口从 12012 递增然后跑 client 完成聚类。注意这一步的副作用脚本用trap kill -HUP 0 0在退出时杀掉它自己起的后台 server 进程只影响本脚本启动的进程。修改 run_on_cluster.bash 的头部参数所有集群参数集中在 run_on_cluster.bash 开头提交任务前需要按自己的环境修改这几处# the training data of the Deep1B dataset deep1bdir/datasets01_101/simsearch/041218/deep1b traindata$deep1bdir/learn.fvecs # this is for small tests nvec1000000 k4000 # for the real run # nvec50000000 # k1000000 # working directory for the real run workdir/checkpoint/matthijs/ondisk_distributed mkdir -p $workdir/{vslices,hslices}需要替换的项deep1bdir/traindata改成你的训练数据目录和文件fvecs/bvecs/npy 均可。如果你用 Deep1B 演示数据下载与拼接方式见 benchs/README.md 的 Getting Deep1B 一节。nvec参与聚类的训练向量数k质心数。脚本默认是小测试用的 100 万向量、4000 质心正式运行的数值在注释行里。workdir输出目录。注意脚本每次运行都会执行mkdir -p $workdir/{vslices,hslices}创建该目录含vslices、hslices子目录请把它改为你有写权限的路径。在 Slurm 上提交分布式作业README 说明该目录的集群命令是面向 Slurm 写的换其他调度器也可以只要它能分配一组机器并在每台机器上启动同一个可执行文件即可。提交入口是bash run_on_cluster.bash slurm_distributed_kmeans该分支内部调用srun请求 5 个任务-n5每个任务 40 核、4 张 GPU--gresgpu:4、100G 内存、最长运行 48 小时分区名prioritysrun -n$nserv \ --time48:00:00 \ --cpus-per-task40 --gresgpu:4 --mem100G \ --partitionpriority --commentpriority is the only one that works \ -l bash $( realpath $0 ) slurm_within_kmeans_server所有srun启动的任务都会执行slurm_within_kmeans_server分支其逻辑是用SLURM_NPROCS得到 server 总数、SLURM_PROCID得到自己的 server idrank并用这两个值把[i0, i1)向量区间和端口12012 rank切分给各节点每个节点以--server --gpu -1 --port $port --ipv4模式启动distributed_kmeans.py对自己负责的那段数据做 assignmentrank 0 额外承担 client 角色它先把本节点自己的 server 挂到后台再用SLURM_TASKS_PER_NODE和SLURM_JOB_NODELIST通过scontrol show hostnames展开拼出全部host:port列表sleep 20s等 server 就绪后以--client --servers $hostports --k $k --ipv4跑聚类client 跑完后脚本打印 Done, kill the job 并执行scancel $SLURM_JOBID取消整个 Slurm 作业——也就是说作业结束时由脚本自己主动取消其他节点的 server 进程随之被回收。使用scancel前请确认你对该 Slurm 作业有取消权限作业本身由你的账号提交通常满足。两个必须对照自己集群核对的点--partitionpriority是文档作者环境的分区名注释里写明 priority is the only one that works你需要改成自己集群上可用的分区必要时加--constraint--cpus-per-task、--gresgpu:4、--mem、--time都是按其集群配置填写的资源请求按实际配额调整。nserv的值决定训练向量如何切分server 数变多时每台机器负责的向量区间自动变小。正式训练5000 万向量聚出 100 万质心README 的 Run used for deep1B 一节给出了实际生产用的跑法对 5000 万向量聚出 100 万质心把质心写入 npy 文件。对应操作是在脚本头部把正式运行参数改为nvec50000000、k1000000文件中已给出这两行注释并把workdir指向有足够空间的目录运行bash run_on_cluster.bash deep1b_clustering与slurm_distributed_kmeans分支的区别该分支用nserv20请求更多任务并在传给 client 的参数里追加--out $workdir/1M_centroids.npy聚类结束后由distributed_kmeans.py执行np.save(args.out, centroids)把质心矩阵写成 npy 文件。如果不想一步跳到 20 个节点也可以先保持slurm_distributed_kmeans分支的 5 节点配置做小规模验证再切到deep1b_clustering。结果验证判断训练是否完成看 client 端输出的最后几行。README 给出的文档示例数值来自作者的 Deep1B 实际运行不要当作固定预期Iteration 19 (898.92 s, search 875.71 s): objective1.33601e07 imbalance1.303 nsplit0 0: writing centroids to /checkpoint/matthijs/ondisk_distributed/1M_centroids.npy即每轮迭代会打印耗时、objective 和 imbalance出现 writing centroids to 你的 --out 路径 说明质心已落盘。之后检查workdir下的1M_centroids.npy是否存在、可用np.load读出的形状是否与k一致聚类产物即完成。README 同时给出了这次运行的性能说明总训练时间 899s 中 876s 用于计算并指出这个纯 Python RPC 的实现里数据传输和质心更新的开销不可忽略RPC 协议不像 MPI 那样为 broadcast/gather 优化——如果你的数据集更大这部分开销会随节点数上升。已知限制与适配点数据文件必须对所有参与节点可见如分布式文件系统否则 server 端 mmap 失败端口固定从 12012 起按 rank 递增client 侧有sleep 20s的固定等待来等 server 就绪机器很多或启动慢时这个等待时间可能需要自行调大--gpu -1表示使用节点全部 GPU--gpu id可指定单卡本地 CPU 基线用--gpu -2不启用 GPU分区、资源参数40 核 / 4 GPU / 100G / 48h和workdir路径均为文档作者环境的取值提交前必须按自己集群改写。聚类产出的质心是后续建索引的输入README 指出下一步是用 make_trained_index.py 构造带质心的空索引含随机旋转预处理、HNSW 包裹质心、训练 6-bit 标量量化器等步骤那属于索引构建场景与本文的 k-means 训练任务分开执行即可。【免费下载链接】faissA library for efficient similarity search and clustering of dense vectors.项目地址: https://gitcode.com/GitHub_Trending/fa/faiss创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考