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进程池初始化时被提前加载,导致子进程无法正常启动:
- 在
transcribejob.py中,模块级别直接初始化了transcriber = Transcriber(WhisperStrategy()),而WhisperStrategy类定义时就加载了Whisper模型(model = whisper.load_model("small"))。 - Celery默认使用prefork进程池,主进程启动后会fork出子Worker进程。Windows不支持真正的fork,prefork实现会重新导入整个模块,导致每个子进程都要重新加载Whisper模型——这不仅耗时,还会因模型加载时的资源占用或死锁,导致子进程无法完成初始化,进而无法处理任务。
- 使用
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
验证步骤
- 重启Redis服务,确保连接正常。
- 按修改后的代码重新启动Celery Worker。
- 发送测试请求,检查Flower中任务状态是否变为
SUCCESS或STARTED,同时查看Redis是否有任务消息。
内容的提问来源于stack exchange,提问作者Estoo
相关产品推荐
相关产品推荐

