如何解决Celery+HuggingFaceEmbedding的WorkerLostError问题
问题描述
使用Celery异步处理Knowledge模型创建时的Qdrant集合生成任务,流程为提取PDF文件内容、生成文本embedding并存储到Qdrant数据库。但在Celery任务内调用HuggingFaceEmbeddings.embed_query时,Worker进程因SIGSEGV信号退出,触发WorkerLostError。不用Celery时API功能正常但响应极慢,尝试拆分为多个小Celery任务又无法传递embeddings等非JSON序列化数据,需要可行的解决思路。
错误日志
celery-dev-1 | [2024-03-27 10:18:27,451: INFO/ForkPoolWorker-19] Load pretrained SentenceTransformer: sentence-transformers/all-mpnet-base-v2 celery-dev-1 | [2024-03-27 10:18:35,856: ERROR/MainProcess] Process 'ForkPoolWorker-19' pid:115 exited with 'signal 11 (SIGSEGV)' celery-dev-1 | [2024-03-27 10:18:35,868: ERROR/MainProcess] Task handler raised error: WorkerLostError('Worker exited prematurely: signal 11 (SIGSEGV) Job: 3.') celery-dev-1 | Traceback (most recent call last): celery-dev-1 | File "/usr/local/lib/python3.10/site-packages/billiard/pool.py", line 1264, in mark_as_worker_lost celery-dev-1 | raise WorkerLostError( celery-dev-1 | billiard.einfo.ExceptionWithTraceback: celery-dev-1 | """ celery-dev-1 | Traceback (most recent call last): celery-dev-1 | File "/usr/local/lib/python3.10/site-packages/billiard/pool.py", line 1264, in mark_as_worker_lost celery-dev-1 | raise WorkerLostError( celery-dev-1 | billiard.exceptions.WorkerLostError: Worker exited prematurely: signal 11 (SIGSEGV) Job: 3. celery-dev-1 | """
相关代码
Knowledge模型代码
class Knowledge(Common): name = models.CharField(max_length=255, blank=True, null=True) file = models.FileField(upload_to=knowledge_path, storage=PublicMediaStorage()) qd_knowledge_id = models.CharField(max_length=255, blank=True, null=True) is_public = models.BooleanField(default=False) def save(self, *args, **kwargs): if self.pk is None: collection_name = f"{self.name}-{datetime.now().strftime('%Y_%m_%d_%H_%M_%S')}" process_files_and_upload_to_qdrant.delay(self.file.name, collection_name) self.qd_knowledge_id = collection_name super().save(*args, **kwargs)
Celery任务及相关函数代码
@shared_task def process_files_and_upload_to_qdrant(file_name, collection_name): file_path = default_storage.open(file_name) result = process_file(file_path, collection_name) return result def process_file(file : InMemoryUploadedFile, collection_name): text = read_data_from_pdf(file) chunks = get_text_chunks(text) embeddings = get_embeddings(chunks) client.create_collection( collection_name=collection_name, vectors_config=qdrant_models.VectorParams( size=768, distance=qdrant_models.Distance.COSINE ), ) client.upsert(collection_name=collection_name, wait=True, points=embeddings) def read_data_from_pdf(file : InMemoryUploadedFile): text = "" pdf_reader = PdfReader(file) for page in pdf_reader.pages: text += page.extract_text() return text def get_text_chunks(texts: str): text_splitter = CharacterTextSplitter( separator="\n", chunk_size=1000, chunk_overlap=200, length_function=len ) chunks = text_splitter.split_text(texts) return chunks def get_embeddings(text_chunks): from langchain_community.embeddings import HuggingFaceEmbeddings from qdrant_client.http.models import PointStruct embeddings = HuggingFaceEmbeddings( model_name="sentence-transformers/all-mpnet-base-v2" ) points = [] for chunk in text_chunks: embedding = embeddings.embed_query(chunk) # 错误发生在此处 point_id = str(uuid.uuid4()) points.append( PointStruct(id=point_id, vector=embedding, payload={"text": chunk}) ) return points
解决思路
1. 修复SIGSEGV崩溃问题
- 根本原因:Celery默认用
fork方式启动worker,而HuggingFace/PyTorch模型在fork进程中加载时,容易因内存页复制、GPU上下文不兼容等问题触发段错误。 - 解决方案:
- 改用单进程worker启动:启动Celery时添加参数
--pool=solo,适合测试或低并发场景,避免fork带来的问题。 - 使用协程池:安装
gevent或eventlet后,启动时用--pool=gevent/--pool=eventlet,协程模式不会触发fork,避免内存错误。 - 预加载模型:将
HuggingFaceEmbeddings的初始化移到任务模块的全局作用域,让worker启动时就加载模型,任务内直接复用,避免每次任务重复加载模型:# 全局初始化,worker启动时加载一次 from langchain_community.embeddings import HuggingFaceEmbeddings embeddings = HuggingFaceEmbeddings(model_name="sentence-transformers/all-mpnet-base-v2") def get_embeddings(text_chunks): from qdrant_client.http.models import PointStruct points = [] for chunk in text_chunks: embedding = embeddings.embed_query(chunk) point_id = str(uuid.uuid4()) points.append(PointStruct(id=point_id, vector=embedding, payload={"text": chunk})) return points
- 改用单进程worker启动:启动Celery时添加参数
2. 解决任务序列化与拆分问题
- 任务间只传递可序列化数据(如文件路径、文本块ID),禁止传递
InMemoryUploadedFile、embeddings等复杂对象:- 修改任务流程,先将拆分后的文本块存储到Redis或临时文件系统,再传递存储标识给后续任务。
- 用Celery的
group实现并行处理文本块:@shared_task def process_pdf_and_split(file_name, collection_name): file_path = default_storage.open(file_name) text = read_data_from_pdf(file_path) chunks = get_text_chunks(text) # 将chunks存入Redis,用collection_name作为key redis_client.set(collection_name, json.dumps(chunks)) # 创建Qdrant集合 client.create_collection( collection_name=collection_name, vectors_config=qdrant_models.VectorParams(size=768, distance=qdrant_models.Distance.COSINE) ) # 生成并行任务处理每个chunk job = group(process_chunk.s(chunk, collection_name) for chunk in chunks) job.apply_async() @shared_task def process_chunk(chunk, collection_name): embedding = embeddings.embed_query(chunk) point_id = str(uuid.uuid4()) client.upsert( collection_name=collection_name, points=[PointStruct(id=point_id, vector=embedding, payload={"text": chunk})] )
- 批量生成embedding:改用
embeddings.embed_documents(text_chunks)批量处理,减少模型调用次数,提升效率:def get_embeddings(text_chunks): from qdrant_client.http.models import PointStruct embeddings_list = embeddings.embed_documents(text_chunks) points = [] for chunk, embedding in zip(text_chunks, embeddings_list): point_id = str(uuid.uuid4()) points.append(PointStruct(id=point_id, vector=embedding, payload={"text": chunk})) return points
3. 其他优化建议
- 给任务添加重试机制:处理临时的网络或服务错误
@shared_task(autoretry_for=(Exception,), retry_backoff=3, retry_kwargs={"max_retries": 3}) def process_files_and_upload_to_qdrant(file_name, collection_name): # 任务逻辑 - 限制worker的并发数:启动Celery时用
--concurrency=2,避免模型加载过多导致内存不足。
内容的提问来源于stack exchange,提问作者OzoneBht
相关产品推荐
相关产品推荐

