Windows环境下Celery无GPU冲突实现及与Flask集成问题
在Windows上部署Celery并避免GPU冲突
部署步骤
Windows对Celery的多进程支持有限,必须用solo进程池,推荐用Redis作为消息中间件(比RabbitMQ更易部署):
- 安装依赖:
pip install celery redis - 编写Celery实例配置(比如
celery_app.py):from celery import Celery celery = Celery( 'video_tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0' ) - 启动Worker:
注意:celery -A celery_app worker --pool=solo --loglevel=info--pool=solo是Windows下必须的参数,否则会出现进程崩溃问题。
避免GPU冲突的方案
GPU冲突本质是多个任务抢占显存或计算资源,可通过以下方式解决:
- 限制GPU任务并发数:给GPU任务单独创建队列,启动Worker时指定
--concurrency=1,确保同一时间只有一个GPU任务运行:
任务定义时指定队列:celery -A celery_app worker --pool=solo --concurrency=1 --queues=gpu_tasks --loglevel=info@celery.task(queue='gpu_tasks') def process_video_with_gpu(video_path): # GPU处理逻辑,比如用PyTorch/TensorFlow import torch torch.cuda.empty_cache() # 任务结束后释放显存 return "GPU处理完成" - 显式指定GPU设备:在任务代码里固定使用某块GPU,避免自动分配冲突:
torch.cuda.set_device(0) # 指定使用第0块GPU - 显存预分配:如果用TensorFlow,可设置显存按需分配:
tf.config.experimental.set_memory_growth(tf.config.list_physical_devices('GPU')[0], True)
Celery与Flask集成
基础集成方式
直接在Flask项目中初始化Celery,无需额外扩展:
- 编写Flask应用(
app.py):from flask import Flask, request, jsonify from celery import Celery import os app = Flask(__name__) # 配置Celery app.config.update( CELERY_BROKER_URL='redis://localhost:6379/0', CELERY_RESULT_BACKEND='redis://localhost:6379/0' ) # 初始化Celery celery = Celery(app.name, broker=app.config['CELERY_BROKER_URL']) celery.conf.update(app.config) # 定义视频处理任务 @celery.task(bind=True) def process_video_task(self, video_path): # 模拟长耗时视频处理(比如转码、AI分析) import time for i in range(10): time.sleep(1) self.update_state(state='PROGRESS', meta={'current': i, 'total': 10}) os.remove(video_path) # 处理完删除临时文件 return "视频处理完成" # 上传接口 @app.route('/upload', methods=['POST']) def upload(): if 'video' not in request.files: return jsonify({'error': '未上传视频文件'}), 400 video_file = request.files['video'] temp_path = f'temp_{video_file.filename}' video_file.save(temp_path) # 提交异步任务 task = process_video_task.delay(temp_path) return jsonify({'task_id': task.id}), 202 # 查询任务状态接口 @app.route('/task/<task_id>') def get_task_status(task_id): task = process_video_task.AsyncResult(task_id) if task.state == 'PENDING': response = {'state': task.state, 'status': '任务等待中...'} elif task.state == 'PROGRESS': response = {'state': task.state, 'status': f'处理中:{task.info["current"]}/{task.info["total"]}'} elif task.state == 'SUCCESS': response = {'state': task.state, 'result': task.result} else: response = {'state': task.state, 'error': str(task.info)} return jsonify(response) if __name__ == '__main__': app.run(debug=True) - 启动服务:
- 启动Flask:
python app.py - 启动Celery Worker:
celery -A app.celery worker --pool=solo --loglevel=info
- 启动Flask:
进阶优化
- 任务状态实时反馈:通过
self.update_state给前端返回处理进度,配合WebSocket可实现实时推送 - 异常防护:在任务中加入try-except块,记录错误日志,避免Worker崩溃
- 超时控制:在Celery配置中设置
CELERY_TASK_TIME_LIMIT=3600(1小时超时),防止任务无限挂起
长耗时视频处理场景的可行性与价值
可行性
完全可行,Celery的异步架构天生适合处理长耗时任务:
- Windows环境下只要用
solo进程池,Worker就能稳定运行,不会出现多进程兼容问题 - 消息中间件(Redis)保证任务不会丢失,即使Worker重启,未完成的任务也能重新执行
- 可以通过队列隔离GPU/CPU任务,合理分配硬件资源
价值
- 提升用户体验:用户上传视频后立即收到任务ID,无需等待处理完成,避免HTTP请求超时
- 服务解耦:Flask只负责接收请求和返回状态,视频处理逻辑交给Celery Worker,降低服务压力
- 可扩展性:如果视频处理需求增加,只需新增Worker节点(比如带更高性能GPU的机器),无需修改Flask代码
- 资源复用:Worker可以持续处理任务,避免每次请求都启动新的处理进程,提升资源利用率
内容的提问来源于stack exchange,提问作者Hossein
相关产品推荐
相关产品推荐

