Celery DjangoTask 指南:用 delay_on_commit 在 Django 事务提交后安全触发任务

Celery DjangoTask 指南:用 delay_on_commit 在 Django 事务提交后安全触发任务 Celery DjangoTask 指南用 delay_on_commit 在 Django 事务提交后安全触发任务【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celerycelery.contrib.django.task是 Celery 5.4 起提供的一个 Django 专用任务基类模块核心价值在于解决任务在数据库事务提交前就被触发这一经典竞态问题。本文基于 celery.contrib.django.task API 参考 展开结合 模块源码、Django fixup 自动装配逻辑 与 单元测试 深入讲解delay_on_commit/apply_async_on_commit两个新 API 的实现原理、使用场景、行为差异与自定义任务基类时的继承方式。读完本文你将能够在 Django 项目中正确、安全地编排事务提交后再执行的异步任务并理解 Celery 与 Django 事务生命周期打交道的底层机制。一、DjangoTask为 Django 扩展的任务基类celery.contrib.django.task模块在 Celery 5.4 中引入对应 API 参考页中的versionadded:: 5.4标记。整个模块只定义了一个类——DjangoTask它的定位非常明确扩展 Celery 基础任务类celery.app.task.Task为 Django 场景提供更友好的 API。源码中的类文档字符串给出了精确定义Extend the base~celery.app.task.Taskfor Django. Provide a nicer API to trigger tasks at the end of the DB transaction.完整实现见 celery/contrib/django/task.pyimport functools from django.db import transaction from celery.app.task import Task class DjangoTask(Task): Extend the base :class:~celery.app.task.Task for Django. Provide a nicer API to trigger tasks at the end of the DB transaction. def delay_on_commit(self, *args, **kwargs) - None: Call :meth:~celery.app.task.Task.delay with Djangos on_commit(). transaction.on_commit(functools.partial(self.delay, *args, **kwargs)) def apply_async_on_commit(self, *args, **kwargs) - None: Call :meth:~celery.app.task.Task.apply_async with Djangos on_commit(). transaction.on_commit(functools.partial(self.apply_async, *args, **kwargs))从源码结构看DjangoTask本身并不重写delay、apply_async等既有行为而是新增了两个以_on_commit为后缀的触发方法。这种设计保持了与普通Task完全兼容的继承关系同时把 Django 事务钩子封装进任务对象本身。二、核心 API 解析delay_on_commit 与 apply_async_on_commit模块为开发者暴露了两个新增方法二者在 Celery 侧分别对应两个最常用的触发入口方法内部委托对象语义delay_on_commit(*args, **kwargs)Task.delay事务提交后执行delay等价于异步发送任务apply_async_on_commit(*args, **kwargs)Task.apply_async事务提交后执行apply_async可传入完整发送选项2.1 实现原理functools.partial transaction.on_commit两个方法的实现如出一辙可以拆成三层理解functools.partial(self.delay, *args, **kwargs)把调用delay并携带全部参数这件事包装成一个零参数的可调用对象。partial在此处的作用是冻结参数让回调在事务提交时无需再传入任何参数即可执行。transaction.on_commit(...)Django 提供的钩子注册的回调会在当前事务成功提交后被执行如果事务回滚则回调不会执行。返回值None两个方法的类型注解均为- None即调用方拿不到任何句柄详见下文行为差异一节。也就是说delay_on_commit(user.pk)的完整语义是等到 Django 当前数据库事务真正提交、数据落库之后再向 Broker 发送任务消息。2.2 使用方式一行替换在定义了任务之后调用方式与delay几乎一致示例取自 Django 官方示例项目文档 from demoapp.tasks import add res add.delay_on_commit(2, 3) res.get() 5注意上例中res实际为None方法不返回结果这里仅为展示调用语法。真正需要拿到结果时请使用delay。三、为什么需要它事务未提交即触发任务的经典竞态DjangoTask要解决的问题在 docs/django/first-steps-with-django.rst 与 docs/userguide/tasks.rst 中被反复强调是 Django Celery 项目中的常见陷阱。3.1 错误示范任务可能跑在数据落库之前# views.py def create_user(request): # Note: simplified example, use a form to validate input user User.objects.create(usernamerequest.POST[username]) send_email.delay(user.pk) return HttpResponse(User created) # tasks.py shared_task def send_email(user_pk): user User.objects.get(pkuser_pk) # send email ...视图里的User.objects.create(...)位于一个隐式事务中send_email.delay(user.pk)会立即把消息发送到 Broker。如果 worker 的执行速度快于视图事务的提交任务线程查询User.objects.get(pkuser_pk)时记录可能尚不存在——任务因此查不到刚创建的数据这就是典型的任务先于事务提交竞态。3.2 更隐蔽的例子transaction.atomic 装饰的视图docs/userguide/tasks.rst 中还给出了一个使用transaction.atomic的变体from django.db import transaction from django.http import HttpResponseRedirect transaction.atomic def create_article(request): article Article.objects.create() expand_abbreviations.delay(article.pk) return HttpResponseRedirect(/articles/)事务的原子性意味着Article.objects.create()产生的数据要等到视图函数返回后事务提交时才会真正持久化。此时delay触发的异步任务一旦先跑起来就可能查询到尚不存在的 article 对象。文档明确指出要防止这种情况必须确保事务提交后再触发任务。四、两种解决方案对比4.1 方案一Celery 5.4手工调用 Django 的 on_commit在DjangoTask出现之前官方推荐的做法是手工使用 Django 的on_commit钩子import functools from django.db import transaction from django.http import HttpResponseRedirect transaction.atomic def create_article(request): article Article.objects.create() transaction.on_commit( functools.partial(expand_abbreviations.delay, article.pk) ) return HttpResponseRedirect(/articles/)4.2 方案二Celery 5.4直接用 delay_on_commitDjangoTask把上述样板封装成了现成 APIdiff 视角看就是一行替换- send_email.delay(user.pk) send_email.delay_on_commit(user.pk)- transaction.on_commit(lambda: send_email.delay(user.pk)) send_email.delay_on_commit(user.pk)与手工functools.partial写法相比delay_on_commit在语义上更直白——延迟到提交后再发送也避免了在业务代码里反复引入django.db.transaction和functools。五、关键行为差异不返回 task ID发送时机延后docs/django/first-steps-with-django.rst 明确列出了delay_on_commit与delay的两个本质差异这也是使用前必须理解的行为约束不返回任务 IDdelay_on_commit不会把任务 ID 返回给调用方源码注解- None即印证这一点。原因在于方法内部只是向transaction.on_commit注册了一个回调调用瞬间任务并未真正发送自然也没有生成可供返回的AsyncResult。如果业务上需要任务 ID 做后续轮询或关联应继续使用delay。发送时机延后到事务提交任务消息不是在调用delay_on_commit时进入 Broker而是在Django 事务成功提交之后才被发送若事务回滚回调不会触发任务也就不会发送。这两个差异在 单元测试 中得到验证——测试通过 patch 将django.db.transaction.on_commit替换为同步执行side_effectlambda f: f()并断言delay_on_commit()与apply_async_on_commit()的返回值均为None。六、自动生效机制DjangoFixup 如何替换任务基类一个很自然的问题是既然要用DjangoTask是不是每个任务都要改成from celery.contrib.django.task import DjangoTask答案是否定的——在标准配置下Celery 会自动完成替换。Django fixup 源码 中的安装逻辑如下def install(self) - DjangoFixup: # Need to add project directory to path. # The project directory has precedence over system modules, # so we prepend it to the path. sys.path.insert(0, os.getcwd()) self._settings symbol_by_name(django.conf:settings) self.app.loader.now self.now if not self.app._custom_task_cls_used: self.app.task_cls celery.contrib.django.task:DjangoTask ...关键点在于if not self.app._custom_task_cls_used:这一行只要你的 Celery 应用没有显式自定义task_clsfixup 就会把应用的任务基类替换为celery.contrib.django.task:DjangoTask。于是所有通过app.task/shared_task定义的任务自动获得delay_on_commit能力无需逐个类去改继承。fixup的触发条件也值得注意环境变量DJANGO_SETTINGS_MODULE已设置、且 loader 非 django loader 时才会安装 Django fixup见fixup(app, envDJANGO_SETTINGS_MODULE)同时会校验 Django 版本——Celery 5.x 要求Django 1.11 或更高版本。七、自定义任务基类时如何继承 DjangoTask官方文档特别提醒了一种例外情况如果项目 使用了自定义任务基类即手动设置了task_clsfixup 中的_custom_task_cls_used为真自动替换不会发生。此时要获得delay_on_commit行为需要显式让你的自定义基类继承DjangoTaskfrom celery.contrib.django.task import DjangoTask class MyBaseTask(DjangoTask): # 你的自定义逻辑…… pass这一约束在 fixup 单元测试 中也有覆盖测试断言使用默认配置时app.task_cls celery.contrib.django.task:DjangoTask且issubclass(f.app.Task, DjangoTask)、hasattr(f.app.Task, delay_on_commit)同时验证了自定义task_cls时专用 DjangoTask 不会被使用。集成测试 t/integration/test_worker.py 中还有一个补充场景Celery 子类如果保持task_cls不变未自定义依然会得到DjangoTask说明自动装配对继承场景同样生效。八、测试与示例项目佐证单元测试t/unit/contrib/django/test_task.py通过pytest.mark.patched_module屏蔽真实 Django 依赖patchtransaction.on_commit后验证两个新方法均可正常调用且返回None是理解方法行为最直接的代码证据。fixup 测试t/unit/fixups/test_django.py覆盖DjangoTask自动替换的启用与禁用分支。官方示例项目examples/django/README.rst完整展示在 Django 项目中启动 workercelery -A proj worker -l INFO后通过manage.py shell调用add.delay_on_commit(2, 3)的全流程示例项目源码位于 examples/django其中 demoapp/tasks.py 展示了使用shared_task解耦任务定义的标准姿势。文档出处本文主题对应的 API 参考页为 docs/reference/celery.contrib.django.task.rst其内容通过automodule直接由 celery/contrib/django/task.py 的 docstring 生成更完整的实战讲解可参阅 docs/django/first-steps-with-django.rst 与 docs/userguide/tasks.rst。九、适用前提与限制小结版本要求DjangoTask及delay_on_commit/apply_async_on_commit自Celery 5.4起可用Celery 5.x 要求 Django ≥ 1.11。仅限 Django 环境该模块直接from django.db import transaction只有在 Django 应用上下文中才能正常导入与使用示例文档亦注明delay_on_commit仅在使用 Django 时可用。行为取舍需要任务 ID 或需要在调用点立即发送任务时请继续使用delaydelay_on_commit专为事务提交后再执行的语义设计。自定义基类注意若项目自定义了task_cls必须显式继承DjangoTask才能获得新 API。总而言之DjangoTask是 Celery 为 Django 事务边界提供的一层轻量但关键的封装它没有改变任务的执行模型只是把发送时机与数据库事务生命周期绑定在一起让开发者用一行代码规避掉最容易踩的竞态陷阱。【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考