批量处理CSV嵌入Pinecone时遇‘无法启动新线程’错误求助
解决"can't start a new thread"错误:批量处理CSV生成嵌入并上传Pinecone
问题背景
我用代码逐行处理数百个CSV文件,通过OpenAI Embeddings生成嵌入后上传到Pinecone,但运行一段时间后总会抛出“can't start a new thread”错误。试过线程和多进程两种方式,结果都一样。
原代码
def process_file(file_path): print(f'file: {file_path}') def process_row(row): text = row['text'] row2data = row['row2data'] year = row['year'] group_id = row['group_id'] docs = embedder(text, text, year, group_id) my_index = pc_store.from_documents(docs, embeddings, index_name=PINECONE_INDEX_NAME) with open(file_path, 'r') as file: reader = csv.DictReader(file) for row in reader: process_row(row) if __name__ == '__main__': file_paths = ['file1', 'file2', 'file3'] processes = [] for file_path in file_paths: p = Process(target=process_file, args=(file_path,)) p.start() processes.append(p) for p in processes: p.join()
错误栈追踪
Traceback (most recent call last): File "/Library/Frameworks/Python.framework/Versions/3.12/lib/python3.12/multiprocessing/pool.py", line 215, in __init__ self._repopulate_pool() File "/Library/Frameworks/Python.framework/Versions/3.12/lib/python3.12/multiprocessing/pool.py", line 306, in _repopulate_pool return self._repopulate_pool_static(self._ctx, self.Process, ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/Library/Frameworks/Python.framework/Versions/3.12/lib/python3.12/multiprocessing/pool.py", line 329, in _repopulate_pool_static w.start() File "/Library/Frameworks/Python.framework/Versions/3.12/lib/python3.12/multiprocessing/dummy/__init__.py", line 51, in start threading.Thread.start(self) File "/Library/Frameworks/Python.framework/Versions/3.12/lib/python3.12/threading.py", line 971, in start _start_new_thread(self._bootstrap, ()) RuntimeError: can't start new thread
问题分析
核心问题是逐行调用pc_store.from_documents():这个方法内部会创建线程池或新线程处理嵌入上传,每一行都调用一次意味着短时间内会生成大量线程,很快耗尽系统线程资源上限。加上多进程模式下每个进程都重复这个行为,资源消耗呈倍数增长,最终触发线程创建失败的错误。
修复方案
1. 复用Pinecone索引实例
不要在每一行都初始化索引,每个进程只初始化一次,后续复用该实例。
2. 批量处理文档
积累一定数量的文档后再批量上传,减少add_documents()的调用次数,降低线程创建频率。
3. 控制并发进程数
使用进程池限制同时运行的进程数量,避免无限制占用系统资源。
修改后的代码
from multiprocessing import Pool import csv def init_process(): # 每个进程启动时初始化一次Pinecone索引 global my_index # 传入空列表初始化索引连接,后续用add_documents上传 my_index = pc_store.from_documents([], embeddings, index_name=PINECONE_INDEX_NAME) def process_file(file_path): print(f'file: {file_path}') batch_size = 50 # 可根据系统性能和Pinecone限制调整 batch_docs = [] with open(file_path, 'r') as file: reader = csv.DictReader(file) for row in reader: text = row['text'] row2data = row['row2data'] year = row['year'] group_id = row['group_id'] docs = embedder(text, text, year, group_id) batch_docs.extend(docs) # 达到批量大小则上传 if len(batch_docs) >= batch_size: my_index.add_documents(batch_docs) batch_docs = [] # 处理最后一批不足batch_size的文档 if batch_docs: my_index.add_documents(batch_docs) if __name__ == '__main__': file_paths = ['file1', 'file2', 'file3'] # 根据CPU核心数设置进程数,比如4个,避免资源过载 with Pool(processes=4, initializer=init_process) as pool: pool.map(process_file, file_paths)
额外优化建议
- 调整
batch_size:如果仍出现线程问题,可适当减小批量大小;若想提高效率,可在系统资源允许的范围内增大。 - 监控系统资源:运行时观察CPU、内存和线程数,确保不超出系统承载能力。
- 异常处理:添加try-except块捕获上传过程中的异常,避免单个文件处理失败导致整个进程崩溃。
内容的提问来源于stack exchange,提问作者John Taylor
相关产品推荐
相关产品推荐

