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 和幂等处理。

参考资料