如何在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
相关产品推荐
相关产品推荐

