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

Windows环境下Celery无GPU冲突实现及与Flask集成问题

在Windows上部署Celery并避免GPU冲突

部署步骤

Windows对Celery的多进程支持有限,必须用solo进程池,推荐用Redis作为消息中间件(比RabbitMQ更易部署):

  1. 安装依赖:
    pip install celery redis
    
  2. 编写Celery实例配置(比如celery_app.py):
    from celery import Celery
    
    celery = Celery(
        'video_tasks',
        broker='redis://localhost:6379/0',
        backend='redis://localhost:6379/0'
    )
    
  3. 启动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,无需额外扩展:

  1. 编写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)
    
  2. 启动服务:
    • 启动Flask:python app.py
    • 启动Celery Worker:celery -A app.celery worker --pool=solo --loglevel=info

进阶优化

  • 任务状态实时反馈:通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 04:25:22