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

如何在Flask+Celery+Redis项目中让多Worker执行同一任务

实现多Celery Worker并行执行process_data任务的方案

Celery原生支持多Worker并行处理任务,完全可以解决你当前的串行延迟问题,具体实施步骤如下:

1. 启动多个Celery Worker实例

你可以通过多种方式启动多个Worker,让它们同时监听任务队列:

  • 终端手动启动:打开多个终端窗口,每个窗口执行以下命令(根据CPU核心数调整--concurrency参数,比如4核CPU设为4):
    celery -A liveapp.celery_worker worker --loglevel=info --concurrency=4
    
  • 后台启动多个Worker:用nohup让Worker在后台运行,每个Worker指定唯一名称区分:
    nohup celery -A liveapp.celery_worker worker --loglevel=info --concurrency=4 --name=worker1 &
    nohup celery -A liveapp.celery_worker worker --loglevel=info --concurrency=4 --name=worker2 &
    
  • 生产环境进程管理:用Supervisor或Systemd管理多个Worker进程,确保意外退出后自动重启,示例Supervisor配置:
    [program:celery-worker-1]
    command=celery -A liveapp.celery_worker worker --loglevel=info --concurrency=4 --name=worker1
    directory=/path/to/your/project/new_test
    user=your_system_user
    autostart=true
    autorestart=true
    stdout_logfile=/var/log/celery/worker1.log
    stderr_logfile=/var/log/celery/worker1_err.log
    
    [program:celery-worker-2]
    command=celery -A liveapp.celery_worker worker --loglevel=info --concurrency=4 --name=worker2
    directory=/path/to/your/project/new_test
    user=your_system_user
    autostart=true
    autorestart=true
    stdout_logfile=/var/log/celery/worker2.log
    stderr_logfile=/var/log/celery/worker2_err.log
    

2. 调整Celery配置优化并行效率

修改__init__.py中的Celery配置,让并行处理更适配实时任务场景:

from celery import Celery

def make_celery(app_name = __name__):
    #local
    redis_uri = "redis://172.30.30.224:6379/6"
    #production
    #redis_uri = "redis://192.168.33.122:6379/6"
    celery = Celery(app_name, backend=redis_uri, broker=redis_uri)
    # 优化并行相关配置
    celery.conf.update(
        worker_prefetch_multiplier=1,  # 每个Worker一次只预取1个任务,避免任务堆积在单个Worker
        task_acks_late=True,  # 任务完成后再向Broker确认,防止Worker崩溃丢失任务
        worker_max_tasks_per_child=1000  # 每个Worker进程处理1000个任务后重启,避免内存泄漏
    )
    return celery
celery = make_celery()

3. 任务代码适配多Worker场景

确保process_data任务在多Worker环境下稳定运行:

  • 独立数据库连接:不要使用全局数据库连接,每个任务内部创建独立连接,避免多Worker共享连接导致的线程安全问题:
    @celery.task()
    def process_data(data):
        # 任务内部创建数据库连接
        conn_str = "your_pyodbc_connection_string"
        conn = pyodbc.connect(conn_str)
        try:
            # 数据清洗、预测逻辑
            cleaned_data = your_cleaning_logic(data)
            prediction_result = your_prediction_logic(cleaned_data)
            
            # 插入数据库
            cursor = conn.cursor()
            cursor.execute("INSERT INTO prediction_results (data, result, created_at) VALUES (?, ?, ?)", 
                          (str(data), prediction_result, datetime.now()))
            conn.commit()
        finally:
            # 确保连接关闭
            conn.close()
            gc.collect()
    
  • 添加日志追踪:在任务中添加日志,方便验证多Worker并行执行:
    import logging
    logger = logging.getLogger(__name__)
    
    @celery.task()
    def process_data(data):
        logger.info(f"Task {process_data.request.id} is running on worker: {process_data.request.hostname}")
        # 原有逻辑
    

4. 验证并行效果

启动多个Worker后,向/api/live接口发送多个请求,查看Celery日志会看到不同的hostname(Worker名称)在处理不同的任务ID,说明多Worker并行生效,串行延迟问题会得到缓解。


内容的提问来源于stack exchange,提问作者Amanda Choy Siew Wen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 00:34:53