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

如何解决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
      

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 22:17:34