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:
单独计算相同标题的嵌入可以成功,但多进程代码无法正常运行,请求帮助。
问题分析与修复
核心问题
- 跨进程共享资源异常:代码中
df_text、embedding_model、np都是父进程的全局对象,在fork模式的多进程环境下,子进程无法正确继承这类带有状态的资源(尤其是模型类),容易导致执行异常或中断。 - 中断后资源未清理:手动中断内核后,进程池的任务队列被破坏,主进程尝试从空队列获取结果触发
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
相关产品推荐
相关产品推荐

