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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 07:57:03