大型pd.DataFrame余弦相似度高效计算方案(Vaex/Dask优先)
Vaex 加速方案
核心思路是将嵌入向量列展开为多列,通过向量化运算批量计算余弦相似度,避免逐行调用cosine_similarity的额外开销。
- 步骤1:将Pandas DataFrame转为Vaex DataFrame,展开(1,50)向量为50个独立列
- 步骤2:手动实现余弦相似度的向量化计算(Vaex支持广播运算)
import vaex import numpy as np # Pandas转Vaex vx_df = vaex.from_pandas(df) # 展开嵌入向量列 for idx in range(50): vx_df[f"emb_{idx}"] = vx_df["embedding"].apply(lambda arr: arr[0, idx]) # 处理目标向量 target_vec = array2.flatten() norm_target = np.linalg.norm(target_vec) # 批量计算点积与范数 emb_matrix = vx_df[[f"emb_{i}" for i in range(50)]].values dot_product = emb_matrix @ target_vec norm_emb = np.sqrt((emb_matrix ** 2).sum(axis=1)) # 计算余弦相似度 vx_df["cos_sim"] = dot_product / (norm_emb * norm_target) # 获取结果 cos_sim_results = vx_df["cos_sim"].values
Dask 加速方案
利用Dask的并行计算能力,将嵌入向量转为Dask Array后做矩阵级运算,充分利用多核CPU。
- 步骤1:Pandas转Dask DataFrame,合理设置分区数(建议为CPU核心数的1-2倍)
- 步骤2:将向量列转为Dask Array,并行计算点积、范数与余弦相似度
import dask.dataframe as dd import dask.array as da import numpy as np # Pandas转Dask,设置分区 ddf = dd.from_pandas(df, npartitions=8) # 将向量列转为(300000, 50)的Dask Array embeddings = ddf["embedding"].apply( lambda arr: arr.flatten(), meta=("embedding", "float64", (50,)) ).to_dask_array(lengths=True) # 处理目标向量 target_vec = da.from_array(array2.flatten(), chunks=50) norm_target = da.linalg.norm(target_vec) # 并行计算点积与范数 dot_product = da.matmul(embeddings, target_vec) norm_emb = da.sqrt(da.sum(embeddings ** 2, axis=1)) # 计算余弦相似度 cos_sim = dot_product / (norm_emb * norm_target) # 触发计算获取结果 cos_sim_results = cos_sim.compute()
额外优化建议
- 避免逐行调用
sklearn.metrics.pairwise.cosine_similarity:该函数针对成对矩阵优化,单向量逐行计算会产生大量冗余开销 - Vaex优化:若数据存储在磁盘(如HDF5、Parquet),直接用
vaex.open()加载,避免Pandas中转的内存拷贝 - Dask优化:根据CPU核心数调整分区数,分区过小会增加调度开销,过大则无法充分利用并行性
内容的提问来源于stack exchange,提问作者Johnny
相关产品推荐
相关产品推荐

