Python线程池执行FAISS相似性搜索时应用程序挂起问题
解决多线程运行FAISS MMR搜索时随机挂起的问题
当前使用ThreadPoolExecutor对DataFramespend_sheet_uniques每行执行基于FAISS的MMR相似性搜索,程序会随机在处理部分行后挂起,无固定停止位置,极少能完成全部行处理。
可能的原因
- FAISS索引非线程安全:底层FAISS索引的操作(包括LangChain封装的
max_marginal_relevance_search方法)大多不具备线程安全性,多线程并发访问时易引发内部死锁或资源竞争,导致程序挂起。 - 共享资源竞争:
indexed_taxonomy_described_cleaned作为多线程共享对象,并发调用其方法时可能出现未定义的冲突行为。 - 低层级并发异常未捕获:虽然
process_mmr_search包含异常捕获,但死锁等底层并发异常可能无法被常规try-except捕获,导致线程永久阻塞。
解决方案
方案1:改用进程池规避线程安全问题
进程池的每个进程拥有独立内存空间,不会共享FAISS索引等资源,从根源上避免线程安全冲突。需确保相关对象可被序列化(pickle):
from concurrent.futures import ProcessPoolExecutor, as_completed def process_mmr_search(row, itemdesc, indexed_taxonomy): try: formatted_itemdesc = str(row[itemdesc]) docs = indexed_taxonomy.max_marginal_relevance_search(formatted_itemdesc, 20) return [doc.page_content for doc in docs] except Exception as e: print(f"Error in MMR search: {e}") return [] def process_row(args): index, row, itemdesc, indexed_taxonomy = args mmr_matches = process_mmr_search(row, itemdesc, indexed_taxonomy) return index, mmr_matches # 初始化进程池并执行任务 with ProcessPoolExecutor(max_workers=4) as executor: task_args = [ (index, row, 'Material Description', indexed_taxonomy_described_cleaned) for index, row in spend_sheet_uniques.iterrows() ] futures = {executor.submit(process_row, args): args for args in task_args} for future in as_completed(futures): index, mmr_matches = future.result() spend_sheet_uniques.at[index, 'Best_Matches_MMR'] = str(mmr_matches)
方案2:给FAISS访问加线程锁
若必须使用线程池,通过线程锁确保同一时间只有一个线程访问FAISS索引方法,避免并发冲突:
from concurrent.futures import ThreadPoolExecutor, as_completed import threading # 创建全局锁控制FAISS访问 faiss_access_lock = threading.Lock() def process_mmr_search(row, itemdesc): try: formatted_itemdesc = str(row[itemdesc]) # 加锁后执行FAISS搜索操作 with faiss_access_lock: docs = indexed_taxonomy_described_cleaned.max_marginal_relevance_search(formatted_itemdesc, 20) return [doc.page_content for doc in docs] except Exception as e: print(f"Error in MMR search: {e}") return [] def threaded_mmr_search(index, row, itemdesc): mmr_matches = process_mmr_search(row, itemdesc) return index, mmr_matches # 带锁执行线程池任务 with ThreadPoolExecutor(max_workers=4) as executor: futures = { executor.submit(threaded_mmr_search, index, row, 'Material Description'): index for index, row in spend_sheet_uniques.iterrows() } for future in as_completed(futures): index, mmr_matches = future.result() spend_sheet_uniques.at[index, 'Best_Matches_MMR'] = str(mmr_matches)
额外优化建议
- 替换
print为线程安全的logging模块,避免多线程输出混乱或缓冲区阻塞。 - 若FAISS索引占用内存过大,可尝试减小
max_workers数量,或对索引进行量化压缩以降低内存占用。 - 使用系统监控工具(如
htop、psutil)观察进程/线程的CPU、内存使用情况,排查是否存在资源耗尽导致的阻塞。
内容的提问来源于stack exchange,提问作者Ashwin Baldawa
相关产品推荐
相关产品推荐

