Celery 5.0.5如何将任务分发至多worker并行执行以降低批量任务耗时
目前所有关于Celery的相关问题几乎都是5到10年前的,基于旧版Celery及对应的设计模式给出的方案
你当前任务耗时指数级上升的核心原因有两个:一是循环内的N+1数据库查询,每遍历一条租赁记录就单独查一次用户表,数据量越大查询开销指数级增长;二是单任务串行处理全量数据,没有充分利用多worker并行能力。以下是适配Celery 5.0.5版本的最优解决方案:
1 先优化数据库查询消除N+1问题
优先把租赁表和用户表做联表查询,一次性拉取所有需要的字段,避免循环内重复查库,这一步能直接把你前置处理的耗时降到线性增长。
2 采用Celery原生批量任务机制实现多worker分发
根据你的数据规模选对应方案即可:
方案A:10万条数据以内用group批量提交任务
你当前循环内逐个调用delay()提交任务会频繁和broker交互,效率很低,改成先组装所有子任务,再一次性提交,broker会自动将任务分发给空闲的多worker并行执行。
优化后代码示例:
from celery import group @celery.task() def send_sms(to, body): from twilio.rest import Client import os account_sid = os.environ["ACCOUNT_SID"] auth_token = os.environ["AUTH_TOKEN"] from_ = os.environ["NUMBER"] client = Client(account_sid, auth_token) message = client.messages.create( to=to, from_=from_, body=body, ) @celery.task() def notify_users(): from datetime import datetime session = create_session() # 联表查询一次拿到所有需要的租赁、用户数据,消除N+1查询 query_result = session.query(Rentals, Users)\ .join(Users, Rentals.user_id == Users.id)\ .filter(Rentals.enabled == True).all() today = datetime.now() task_list = [] for rental, user in query_result: if rental.returned_date is not None: if (today - rental.returned_date).total_seconds() < rental.rental_period: continue to = send_notification_get_to.get(rental.notification_method)(user) body = f"sending notification to {user.email}" # 先组装任务签名,暂不提交 task_list.append(send_sms.s(to, body)) # 一次性提交所有任务,自动分发到多worker并行执行 group(task_list).apply_async()
方案B:10万条数据以上用chunks分块处理
如果数据量太大,一次性组装所有子任务会占用过多内存,可以把数据拆成固定大小的块,每个块交给一个子任务处理,再批量分发块任务到多worker:
from celery import chunks @celery.task() def send_sms(to, body): from twilio.rest import Client import os account_sid = os.environ["ACCOUNT_SID"] auth_token = os.environ["AUTH_TOKEN"] from_ = os.environ["NUMBER"] client = Client(account_sid, auth_token) message = client.messages.create( to=to, from_=from_, body=body, ) @celery.task def process_rental_chunk(rental_user_pairs): from datetime import datetime today = datetime.now() for rental, user in rental_user_pairs: if rental.returned_date is not None: if (today - rental.returned_date).total_seconds() < rental.rental_period: continue to = send_notification_get_to.get(rental.notification_method)(user) body = f"sending notification to {user.email}" send_sms.delay(to, body) @celery.task() def notify_users(): session = create_session() query_result = session.query(Rentals, Users)\ .join(Users, Rentals.user_id == Users.id)\ .filter(Rentals.enabled == True).all() # 每1000条数据拆成一个块,自动分发到多worker处理 chunk_tasks = process_rental_chunk.chunks(query_result, 1000) chunk_tasks.apply_async()
3 worker配置优化
因为你的短信发送任务是IO密集型,建议用gevent worker池提升并发能力,启动命令示例:celery -A 你的项目模块名 worker --pool=gevent --concurrency=100 -l INFO
如果是CPU密集型任务,用默认的prefork池即可,并发数设为CPU核数的2-4倍。
附加注意点
如果notify_users是定时触发的任务,一定要加分布式锁,避免上一轮任务还没跑完下一轮又触发,导致重复发送通知,用Redis或数据库实现锁都可以。
内容的提问来源于stack exchange,提问作者uberrebu

