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

Celery任务卡在Received状态:Docker Redis作为Broker时无法执行

问题背景

需求目标

在Python/Flask项目中,通过Celery实现语音转写任务的异步后台处理,采用Docker部署的Redis作为消息中间件(Broker),并用Flower监控Celery任务。

问题现象

Flower控制面板显示任务已接收(Received)但始终未执行;Broker页面显示Celery从未向Redis发送消息。仅当启动Celery Worker时设置--pool参数为solo,任务才能运行,但该模式不符合生产环境规范,且仅能执行1个任务,之后无法接收新任务。

代码详情

项目包含app.py、transcription.py、transcribejob.py三个文件:

transcribejob.py

from transcription import Transcriber, WhisperStrategy
from celery import Celery
import os
import json
import sys

transcriber = Transcriber(WhisperStrategy())
#to start the celery worker, enter in the terminal:
#celery -A transcribejob:celery worker --loglevel=info --max-memory-per-child=1000000 --concurrency=4

#to monitor proces: celery -A transcribejob flower
redis_url = 'redis://127.0.0.1:6379/0'
celery = Celery('transcribejob', broker=redis_url, backend=redis_url)
cwd = os.getcwd()
print(f"transcribejob.py, cwd: {cwd}", file=sys.stdout)

@celery.task
def transcribe_export(temp_audio_path, audio_file_name):
    result = transcriber.transcribe(os.path.join(temp_audio_path, audio_file_name))
    txt_fn = audio_file_name.replace(".wav", "") + '.txt'
    with open(os.path.join(cwd, txt_fn), 'w') as f:
        json.dump(result, f)

transcription.py

import whisper


class TranscribeStrategy:
    def __init__(self, vendor):
        self.vendor = vendor

    def do_transcribe(self, audio_file_path):
        pass


class WhisperStrategy(TranscribeStrategy):
    model = whisper.load_model("small")

    def __init__(self):
        super().__init__("openai-whisper-small")

    def do_transcribe(self, audio_file_path):
        return self.model.transcribe(audio_file_path, language="nl")


class Transcriber:
    def __init__(self, transcribe_strategy):
        self.strategy = transcribe_strategy

    def transcribe(self, audio_file_path):
        return {"vendor": self.strategy.vendor, "result": self.strategy.do_transcribe(audio_file_path)}

app.py

from flask import Flask, request
from io import BytesIO
from transcribejob import transcribe_export
import os
import uuid
from pydub import AudioSegment

app = Flask(__name__)
temp_audio_path = os.path.join(os.getcwd(), 'tempAudioStorage')

@app.route('/transcribe', methods=['POST'])
def transcribe_audio_file():
    if 'file' not in request.files:
        return "No file uploaded", 400

    file = request.files['file']

    if file.content_type not in {'audio/wav', 'audio/mpeg'}:
        return "Invalid file format. Only WAV and MP3 files are allowed. Received Format: "+file.content_type, 400

    file_data = BytesIO(file.read())
    audio_format = file.content_type.split('/')[-1]
    if audio_format == 'mpeg':
        audio_format = 'mp3'
    audio = AudioSegment.from_file(file_data, format=audio_format)
    audio_file_name = str(uuid.uuid4())+".wav"
    audio.export(os.path.join(temp_audio_path, audio_file_name), format='wav')

    transcribe_export.delay(temp_audio_path, audio_file_name)

    return 'Transcription request received, transcription in process', 200


if __name__ == '__main__':
    app.run()

运行输出

启动Flask后,执行Celery命令的输出:

(venv) PS C:\Users\SeanS\PycharmProjects\isampTranscribe> celery -A transcribejob:celery worker --loglevel=info --max-memory-per-child=1000000 --concurrency=4

transcribejob.py, cwd: C:\Users\SeanS\PycharmProjects\isampTranscribe


 -------------- celery@SeanS01 v5.2.7 (dawn-chorus)

--- *****

-- ******* ---- Windows-10-10.0.22621-SP0 2023-04-13 14:26:11

- *** --- * ---

- ** ---------- [config]

- ** ---------- .> app:         transcribejob:0x22c05f52c10

- ** ---------- .> transport:   redis://127.0.0.1:6379/0

- ** ---------- .> results:     redis://127.0.0.1:6379/0

- *** --- * --- .> concurrency: 4 (prefork)

-- ******* ---- .> task events: OFF (enable -E to monitor tasks in this worker)

--- *****

-------------- [queues]

                .> celery           exchange=celery(direct) key=celery

 
 
[tasks]

  . transcribejob.transcribe_export

 
[2023-04-13 14:26:11,887: INFO/MainProcess] Connected to redis://127.0.0.1:6379/0

[2023-04-13 14:26:11,891: INFO/MainProcess] mingle: searching for neighbors

[2023-04-13 14:26:12,429: INFO/SpawnPoolWorker-1] child process 18524 calling self.run()

[2023-04-13 14:26:12,450: INFO/SpawnPoolWorker-3] child process 29760 calling self.run()

[2023-04-13 14:26:12,453: INFO/SpawnPoolWorker-2] child process 30472 calling self.run()

[2023-04-13 14:26:12,460: INFO/SpawnPoolWorker-4] child process 30600 calling self.run()

[2023-04-13 14:26:12,921: WARNING/MainProcess] C:\Users\SeanS\PycharmProjects\isampTranscribe\venv\lib\site-packages\celery\app\control.py:56: DuplicateNodenameWarning: Received multiple replies from node name: celery@SeanS01.

Please make sure you give each node a unique nodename using

the celery worker `-n` option.

  warnings.warn(DuplicateNodenameWarning(

 
[2023-04-13 14:26:12,922: INFO/MainProcess] mingle: all alone

[2023-04-13 14:26:12,944: INFO/MainProcess] celery@SeanS01 ready.

[2023-04-13 14:26:15,346: INFO/MainProcess] Events of group {task} enabled by remote.

transcribejob.py, cwd: C:\Users\SeanS\PycharmProjects\isampTranscribe

transcribejob.py, cwd: C:\Users\SeanS\PycharmProjects\isampTranscribe

transcribejob.py, cwd: C:\Users\SeanS\PycharmProjects\isampTranscribe

transcribejob.py, cwd: C:\Users\SeanS\PycharmProjects\isampTranscribe

[2023-04-13 14:27:09,247: INFO/MainProcess] Task transcribejob.transcribe_export[653178f2-c67d-4f96-afa8-f1cbcd470d91] received

[2023-04-13 14:27:10,041: INFO/SpawnPoolWorker-5] child process 25092 calling self.run()

[2023-04-13 14:27:10,051: INFO/SpawnPoolWorker-6] child process 26020 calling self.run()

transcribejob.py, cwd: C:\Users\SeanS\PycharmProjects\isampTranscribe

transcribejob.py, cwd: C:\Users\SeanS\PycharmProjects\isampTranscribe

[2023-04-13 14:27:15,740: INFO/SpawnPoolWorker-7] child process 30576 calling self.run()

transcribejob.py, cwd: C:\Users\SeanS\PycharmProjects\isampTranscribe

Flower仪表盘显示任务状态为Received,且无消息发送至Redis Broker;Docker部署的Redis已正确映射6379端口。


问题分析与解决

核心原因

问题根源在于Whisper模型在Celery prefork进程池初始化时被提前加载,导致子进程无法正常启动:

  1. 在transcribejob.py中,模块级别直接初始化了transcriber = Transcriber(WhisperStrategy()),而WhisperStrategy类定义时就加载了Whisper模型(model = whisper.load_model("small"))。
  2. Celery默认使用prefork进程池,主进程启动后会fork出子Worker进程。Windows不支持真正的fork,prefork实现会重新导入整个模块,导致每个子进程都要重新加载Whisper模型——这不仅耗时,还会因模型加载时的资源占用或死锁,导致子进程无法完成初始化,进而无法处理任务。
  3. 使用solo模式时只有单个进程,不会触发多进程重新加载模块的问题,所以任务能执行,但单进程无法处理多任务,不符合生产需求。

解决方案

1. 延迟模型加载,避免模块级初始化

修改transcribejob.py,将Transcriber和模型的初始化移到任务函数内部:

from transcription import Transcriber, WhisperStrategy
from celery import Celery
import os
import json
import sys

redis_url = 'redis://127.0.0.1:6379/0'
celery = Celery('transcribejob', broker=redis_url, backend=redis_url)
cwd = os.getcwd()
print(f"transcribejob.py, cwd: {cwd}", file=sys.stdout)

@celery.task
def transcribe_export(temp_audio_path, audio_file_name):
    # 延迟初始化,确保每个Worker进程仅在首次执行任务时加载模型
    transcriber = Transcriber(WhisperStrategy())
    result = transcriber.transcribe(os.path.join(temp_audio_path, audio_file_name))
    txt_fn = audio_file_name.replace(".wav", "") + '.txt'
    with open(os.path.join(cwd, txt_fn), 'w') as f:
        json.dump(result, f)

同时修改transcription.py,将模型加载移到__init__方法中,避免类级别静态变量导致的提前加载:

import whisper


class TranscribeStrategy:
    def __init__(self, vendor):
        self.vendor = vendor

    def do_transcribe(self, audio_file_path):
        pass


class WhisperStrategy(TranscribeStrategy):
    def __init__(self):
        super().__init__("openai-whisper-small")
        # 移到__init__中,确保实例化时才加载模型
        self.model = whisper.load_model("small")

    def do_transcribe(self, audio_file_path):
        return self.model.transcribe(audio_file_path, language="nl")


class Transcriber:
    def __init__(self, transcribe_strategy):
        self.strategy = transcribe_strategy

    def transcribe(self, audio_file_path):
        return {"vendor": self.strategy.vendor, "result": self.strategy.do_transcribe(audio_file_path)}

2. 针对Windows环境调整Celery进程池

Windows对prefork的支持有限,可改用gevent或eventlet作为进程池:

  • 安装依赖:pip install gevent
  • 启动Worker时指定池类型:celery -A transcribejob:celery worker --loglevel=info --max-memory-per-child=1000000 --concurrency=4 --pool=gevent

3. 解决重复节点名警告

启动Worker时指定唯一节点名,避免DuplicateNodenameWarning:

celery -A transcribejob:celery worker --loglevel=info --max-memory-per-child=1000000 --concurrency=4 -n worker1@%h

验证步骤

  1. 重启Redis服务,确保连接正常。
  2. 按修改后的代码重新启动Celery Worker。
  3. 发送测试请求,检查Flower中任务状态是否变为SUCCESS或STARTED,同时查看Redis是否有任务消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 10:22:35