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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 12:25:16