Django 新订单偶发“查询不到”:Celery 任务早于事务提交的排查与修复
适用场景与现象
下单接口先写入订单,再让 Celery Worker 发送确认邮件。接口返回成功,Worker 却偶发 Order.DoesNotExist;在管理后台查询,同一订单已经存在。失败多发生在数据库写入较慢、Worker 较空闲时。本文假设订单和任务分别由不同进程访问同一数据库,示例中的 orders 是 Django 应用名。
一个容易出错的实现是:
from django.db import transaction
from orders.models import Order
from orders.tasks import send_order_receipt
@transaction.atomic
def create_order(email: str) -> int:
order = Order.objects.create(email=email)
send_order_receipt.delay(order.pk)
return order.pk
delay() 立即向消息代理投递任务。此时 atomic() 尚未退出,订单对另一个数据库连接通常不可见。Worker 抢先执行查询就会报“记录不存在”;等接口返回后再人工查询,事务已经提交,因此现场看起来互相矛盾。
先确认是不是提交时序
在接口和任务日志中同时记录 order_id、任务 ID、事件名和时间,核对“投递任务、Worker 查询、事务提交”三者的顺序。固定日志文案保持简洁,不记录完整邮箱地址。排查时还要确认 Worker 与 Web 使用相同环境和数据库配置;读写分离时,读副本延迟也可能表现为同样的异常。若主库提交后能查到、只在读副本查不到,应处理副本延迟或让这类立即执行的任务读主库。
不要把 sleep(1) 当修复方案:数据库耗时和队列调度没有固定上界。也不要只为掩盖首次查不到而无限重试,否则真正删除、错误 ID 或库配置错误会变成难以发现的积压任务。
修复:提交成功后再投递
Celery 5.4 起,Django 集成的 DjangoTask 提供 delay_on_commit()。下面给出可放进现有 orders 应用的最小示例;项目需已按 Celery 的 Django 集成文档配置应用和 Broker。
# orders/models.py
from django.db import models
class Order(models.Model):
email = models.EmailField()
created_at = models.DateTimeField(auto_now_add=True)
# orders/tasks.py
from celery import shared_task
from django.conf import settings
from django.core.mail import send_mail
from orders.models import Order
@shared_task
def send_order_receipt(order_id: int) -> None:
order = Order.objects.get(pk=order_id)
send_mail(
subject=f"订单 {order.pk} 已创建",
message="我们已收到您的订单。",
from_email=settings.DEFAULT_FROM_EMAIL,
recipient_list=[order.email],
fail_silently=False,
)
# orders/services.py
from django.db import transaction
from orders.models import Order
from orders.tasks import send_order_receipt
@transaction.atomic
def create_order(email: str) -> int:
order = Order(email=email)
order.full_clean()
order.save()
send_order_receipt.delay_on_commit(order.pk)
return order.pk
关键点是传主键,而不是把尚未提交的模型对象序列化给任务。delay_on_commit() 注册回调:外层事务成功提交后才调用 delay();若事务回滚,回调被丢弃。嵌套 atomic() 的内层退出只释放保存点,仍要等最外层事务提交。若项目启用了 ATOMIC_REQUESTS,回调会等请求事务提交后执行。当前没有活动事务时,回调会立即执行。
旧版 Celery 或自定义任务基类没有这个方法时,可以使用 Django 自身的接口:
from functools import partial
from django.db import transaction
from orders.tasks import send_order_receipt
transaction.on_commit(partial(send_order_receipt.delay, order.pk))
把这行放在创建订单的事务内,替换原来的 delay()。partial 在注册时绑定当前 order.pk,避免循环或后续变量变化导致回调使用错误的 ID。自定义 Celery 任务基类若要使用 delay_on_commit(),需继承 celery.contrib.django.task.DjangoTask。
用测试锁住提交与回滚行为
Django 的 TestCase 会把每个测试包在事务里,默认不会真的提交;直接断言任务已投递会得到误导性结果。captureOnCommitCallbacks(execute=True) 可以模拟执行提交回调。以下测试不会连接真实 Broker,也不会发邮件:
# orders/tests/test_services.py
from unittest.mock import patch
from django.db import transaction
from django.test import TestCase
from orders.services import create_order
from orders.tasks import send_order_receipt
class CreateOrderTests(TestCase):
@patch("orders.services.send_order_receipt.delay")
def test_dispatches_only_after_commit_callback(self, mocked_delay):
with self.captureOnCommitCallbacks(execute=True) as callbacks:
order_id = create_order("reader@example.test")
mocked_delay.assert_not_called()
self.assertEqual(len(callbacks), 1)
mocked_delay.assert_called_once_with(order_id)
@patch("orders.services.send_order_receipt.delay")
def test_rollback_discards_callback(self, mocked_delay):
with self.captureOnCommitCallbacks(execute=True) as callbacks:
try:
with transaction.atomic():
send_order_receipt.delay_on_commit(123)
raise ValueError("模拟业务回滚")
except ValueError:
pass
self.assertEqual(callbacks, [])
mocked_delay.assert_not_called()
在项目根目录执行 python manage.py test orders.tests.test_services。第一项检查投递发生在回调执行时,第二项检查回滚不会投递。若测试报任务对象没有 delay_on_commit,先核对 Celery 版本与自定义任务基类。
上线前的边界与预防
on_commit() 解决的是“消息先于数据可见”的竞态,不保证消息一定送达:数据库已提交后,Broker 仍可能不可用。此时订单不会因投递失败而回滚,业务需要记录告警并补偿;若要求提交成功后最终必定投递,应在同一数据库事务内写入待发送事件,使用 outbox 表和独立发布器重试。
邮件这类外部副作用还要考虑任务重复投递或 Worker 重试。生产实现应按订单和通知类型设计幂等记录或供应商幂等键,避免重复发信。delay_on_commit() 不返回任务 ID;需要在响应里展示任务 ID 时,应另外设计状态标识,不能假设注册回调时已有 Celery 结果对象。
总结
遇到“接口成功、Worker 偶发查不到新订单”,先对齐投递、查询和提交时间,并排除数据库配置与读副本延迟。确认是事务竞态后,将投递移到 on_commit();同时用提交与回滚测试固定行为。对必须可靠送达的事件,再引入 outbox 和幂等处理。
Discussion
评论