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

Python多进程计算句嵌入时遇IndexError: pop from empty deque问题求助

多进程计算句子嵌入报错问题排查与修复

问题描述

我编写了如下Python代码,尝试通过multiprocessing并行计算句子嵌入:

import multiprocessing
from tqdm import tqdm

# Define the function to be executed in parallel
def process_data(chunk):
    results = []
    for row in tqdm(chunk):
        work_id = row[1]
        mentioning_work_id = row[3]

        if work_id in df_text and mentioning_work_id in df_text:
            title1 = df_text[work_id]['title']
            title2 = df_text[mentioning_work_id]['title']
            print(title1 + '/n' + title2)
            embeddings_title1 = embedding_model.encode(title1,convert_to_numpy=True)
            print(embeddings_title1)
            embeddings_title2 = embedding_model.encode(title2,convert_to_numpy=True)
            print(embeddings_title2)
            results.append(np.matmul(embeddings_title1, embeddings_title2.T))
            print(results)
        else:
            continue
    return results

from multiprocessing import Pool

# Define the data to be processed
data = df_rud_labels

# Split the data into chunks
chunk_size = len(data) // num_cores
chunks = [data[i:i+chunk_size] for i in range(0, len(data), chunk_size)]

# Create a pool of worker processes
pool = multiprocessing.Pool(processes=num_cores)

results = []
with tqdm(total=len(chunks)) as pbar:
    for i, result_chunk in enumerate(pool.imap_unordered(process_data, chunks)):
        # Update the progress bar
        pbar.update()
        # Add the results to the list
        results += result_chunk

# Concatenate the results
final_result = results

中断内核后遇到如下错误:

0%|          | 0/2500 [00:00<?, ?it/s]
Financialization and Institutional Change in Capitalisms: A Comparison of the US and Germany/nSingle domain antibodies: promising experimental and therapeutic tools in infection and immunity
  0%|          | 0/2500 [00:00<?, ?it/s]
Encyclopedia of India-China Cultural Contacts, vol I/nToll-like receptors as a key regulator of mesenchymal stem cell function: An up-to-date review
  0%|          | 0/2500 [00:00<?, ?it/s]
Ioannis ROMANIDES, Dogmatica patristica ortodoxa, traducere de Dragos Dasca, Editura Ecclesiast, editie de protos Vasile Bîrzu, 2011/nANTHROPOMETRIC MEASUREMENTS, SOMATOTYPES AND PHYSICAL ABILITIES AS A FUNCTION TO PREDICT THE SELECTION OF TALENTS JUNIOR WEIGHTLIFTERS 
A prophet of old: Jesus the “public theologian”/nCurriculum alignment at undergraduate level: military geography at the South African Military Academy
  0%|          | 0/4 [00:09<?, ?it/s]Process ForkPoolWorker-72:
Process ForkPoolWorker-71:
Traceback (most recent call last):
Traceback (most recent call last):
  File "/opt/conda/lib/python3.7/multiprocessing/process.py", line 297, in _bootstrap
    self.run()
  File "/opt/conda/lib/python3.7/multiprocessing/process.py", line 99, in run
    self._target(*self._args, **self._kwargs)
  File "/opt/conda/lib/python3.7/multiprocessing/process.py", line 297, in _bootstrap
    self.run()
  File "/opt/conda/lib/python3.7/multiprocessing/process.py", line 99, in run
    self._target(*self._args, **self._kwargs)
  File "/opt/conda/lib/python3.7/multiprocessing/pool.py", line 110, in worker
    task = get()
  File "/opt/conda/lib/python3.7/multiprocessing/queues.py", line 351, in get
    with self._rlock:
  File "/opt/conda/lib/python3.7/multiprocessing/synchronize.py", line 95, in __enter__
    return self._semlock.__enter__()
KeyboardInterrupt
  File "/opt/conda/lib/python3.7/multiprocessing/pool.py", line 110, in worker
    task = get()
  File "/opt/conda/lib/python3.7/multiprocessing/queues.py", line 351, in get
    with self._rlock:
  File "/opt/conda/lib/python3.7/multiprocessing/synchronize.py", line 95, in __enter__
    return self._semlock.__enter__()
KeyboardInterrupt

---------------------------------------------------------------------------
IndexError                                Traceback (most recent call last)
/opt/conda/lib/python3.7/multiprocessing/pool.py in next(self, timeout)
    732             try:
--> 733                 item = self._items.popleft()
    734             except IndexError:

IndexError: pop from an empty deque

During handling of the above exception, another exception occurred:

KeyboardInterrupt                         Traceback (most recent call last)
<ipython-input-48-fcb9ab74a032> in <module>
     31 results = []
     32 with tqdm(total=len(chunks)) as pbar:
--> 33     for i, result_chunk in enumerate(pool.imap_unordered(process_data, chunks)):
     34         # Update the progress bar
     35         pbar.update()

/opt/conda/lib/python3.7/multiprocessing/pool.py in next(self, timeout)
    735                 if self._index == self._length:
    736                     raise StopIteration from None
--> 737                 self._cond.wait(timeout)
    738                 try:
    739                     item = self._items.popleft()

/opt/conda/lib/python3.7/threading.py in wait(self, timeout)
    294         try:    # restore state no matter what (e.g., KeyboardInterrupt)
    295             if timeout is None:
--> 296                 waiter.acquire()
    297                 gotit = True
    298             else:

KeyboardInterrupt: 

单独计算相同标题的嵌入可以成功,但多进程代码无法正常运行,请求帮助。

问题分析与修复

核心问题

  1. 跨进程共享资源异常:代码中df_text、embedding_model、np都是父进程的全局对象,在fork模式的多进程环境下,子进程无法正确继承这类带有状态的资源(尤其是模型类),容易导致执行异常或中断。
  2. 中断后资源未清理:手动中断内核后,进程池的任务队列被破坏,主进程尝试从空队列获取结果触发IndexError,但根源还是资源共享问题。

修复步骤

1. 子进程内独立初始化资源

每个子进程需要单独加载模型和数据,避免依赖父进程的全局对象:

def process_data(chunk):
    # 子进程内独立导入依赖、初始化模型和数据
    import numpy as np
    from sentence_transformers import SentenceTransformer  # 替换为你使用的嵌入模型库
    
    # 重新加载df_text(或通过Manager共享,见后续建议)
    def load_df_text():
        # 替换成你加载df_text的实际代码
        import pandas as pd
        df = pd.read_csv("your_data_path.csv")
        return df.set_index('work_id')['title'].to_dict()
    
    df_text = load_df_text()
    embedding_model = SentenceTransformer('your_model_name')  # 替换为你的模型名称
    
    results = []
    for row in chunk:
        work_id = row[1]
        mentioning_work_id = row[3]

        if work_id in df_text and mentioning_work_id in df_text:
            title1 = df_text[work_id]
            title2 = df_text[mentioning_work_id]
            embeddings_title1 = embedding_model.encode(title1, convert_to_numpy=True)
            embeddings_title2 = embedding_model.encode(title2, convert_to_numpy=True)
            results.append(np.matmul(embeddings_title1, embeddings_title2.T))
    return results

2. 避免子进程进度条混乱

多个子进程同时输出tqdm进度条会导致终端显示错乱,建议只保留主进程的全局进度条,去掉子进程内的tqdm(chunk)。

3. 正确关闭进程池

在任务执行完成后,必须关闭并等待进程池释放资源:

results = []
with tqdm(total=len(chunks)) as pbar:
    for i, result_chunk in enumerate(pool.imap_unordered(process_data, chunks)):
        pbar.update()
        results += result_chunk

# 关闭进程池并等待所有子进程结束
pool.close()
pool.join()

final_result = results

4. 处理手动中断场景(可选)

添加异常捕获,在手动中断时安全终止进程池:

results = []
try:
    with tqdm(total=len(chunks)) as pbar:
        for i, result_chunk in enumerate(pool.imap_unordered(process_data, chunks)):
            pbar.update()
            results += result_chunk
except KeyboardInterrupt:
    print("任务中断,正在终止进程池...")
    pool.terminate()
    pool.join()
    raise
else:
    pool.close()
    pool.join()

final_result = results

额外优化建议

  • 共享大字典减少内存消耗:如果df_text体积较大,可通过multiprocessing.Manager创建共享字典,避免每个子进程重复加载:
    from multiprocessing import Manager
    
    manager = Manager()
    shared_df_text = manager.dict(df_text.to_dict())  # 假设df_text是DataFrame
    
    # 修改process_data,接收shared_df_text作为参数
    def process_data(chunk, shared_df_text):
        import numpy as np
        from sentence_transformers import SentenceTransformer
        
        embedding_model = SentenceTransformer('your_model_name')
        results = []
        for row in chunk:
            work_id = row[1]
            mentioning_work_id = row[3]
            
            if work_id in shared_df_text and mentioning_work_id in shared_df_text:
                title1 = shared_df_text[work_id]
                title2 = shared_df_text[mentioning_work_id]
                # 后续计算逻辑不变
        return results
    
    # 调用时传入共享字典
    for result_chunk in pool.imap_unordered(process_data, chunks, [shared_df_text]*len(chunks)):
        # ...
    
  • 使用模型自带的并行能力:多数句子嵌入库(如sentence-transformers)的encode方法本身支持多线程和批量处理,比手动用multiprocessing更稳定,可直接尝试:
    titles = [df_text[work_id]['title'] for work_id in df_rud_labels[1] if work_id in df_text]
    embeddings = embedding_model.encode(titles, batch_size=32, show_progress_bar=True, convert_to_numpy=True)
    

内容的提问来源于stack exchange,提问作者Patthebug

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 03:12:00