← 文章 / 编程开发
freeCodeCamp 6小时前 · 2026-09-12 08:15:00 · 5 阅读

如何在 Django 中防止竞态条件

假设你在一个 AI 应用里还剩最后一次生成图片的额度。你在一个浏览器标签页提交请求,然后在第一个请求完成前,又立刻在第二个标签页提交了另一个请求。

应用两次都接受了。

表面上看,每个请求都通过了额度检查。但实际上,应用承接的工作量超过了你的余额所能支付的范围。你刚才经历的这就是一种竞态条件(Race Condition)。

对于开发者来说,这引出了两个问题:

  • 额度只够扣一次,为什么两个请求都能通过检查?

  • 如何防止它们消耗同一份额度?

在这篇教程中,我们将构建一个小型 Django 额度系统来探讨这些问题。我们会复现这个 Bug,然后使用数据库事务和行锁来修复它,并测试当两个请求争抢最后一份额度时会发生什么。

内容大纲

目标读者

本教程面向熟悉基础 Django 但对并发概念陌生的开发者。

了解模型、迁移和视图(Views)有助于你跟进示例。按此处的设置指南,你需要 Python 3.12 以及带 Compose 的 Docker。我们将使用 Django 5.2、Django REST Framework 3.16 和 PostgreSQL 17,因为 SQLite 不支持我们将使用的行锁(Row Lock)功能。

在应用到代码之前,我会先解释并发的相关概念。

整个教程中,图片生成过程将被模拟。

我们的 API 将接收提示词并返回模拟结果。你无需 AI 服务账号或付费的 API 密钥。这样可以让练习专注于额度决策及其数据库变更。

你可以全程跟进示例而无需为图片请求付费。

理解并发与竞态条件

什么是并发?

并发是指两个或多个任务在重叠的时间段内同时推进。

比如,服务器可以在处理请求 A 的同时开始处理请求 B,而请求 A 正在等待数据库响应。两者的指令不必在同一瞬间执行。事件循环可以在某个任务等待时切换到另一个任务,这也是 Python 支持并发的一种方式。

关键在于:第一个操作还没完成,另一个操作就已经可以推进了。

什么是竞态条件?

竞态条件是一种缺陷,指结果的正确性取决于并发操作的执行时机或顺序。

当多个操作共享数据,而应用没有充分协调它们的访问时,就可能出现竞态条件。每个操作单独看都没问题,但某个操作可能基于另一个操作已经修改过的信息行事。换一种执行顺序,就可能产生不同的、错误的结果。

并发创造了操作重叠的机会,而竞态条件就是应用在处理这种重叠时可能出现的 bug。

竞态条件还会发生在哪些地方?

计数器、用户账户、后台任务、浏览器中显示的结果,都可能受竞态条件影响。

竞争的参与者可能是同一个人的两个请求、两个不同用户,甚至是完全不需要用户参与的自动化任务。共享资源可能是一条数据库记录、一个文件、一个内存值,或者页面的当前状态。下面的例子展示了几种顺序会产生影响的情况。

要留意那些共享状态、且可能在彼此完成前相互干扰的操作。

浏览计数丢失更新

假设两个请求都想把一篇文章的浏览量从 100 加一。

两者在任何一方保存之前都读到了 100。各自加一后都写入了 101。计数器本应达到 102,结果一次更新丢失了。

这就是典型的丢失更新,而且和购买、库存毫无关系。

两个注册抢同一个用户名

如果注册表单只依赖一个初步的用户名查重,就可能产生竞态。

两个请求检查了同一个用户名,都发现它可用,于是各自尝试用这个名字创建账户。如果数据库没有唯一性约束,应用就可能创建出重复的用户名。

这里的防范手段是在数据库层面强制唯一性,而不是相信之前那次检查。

两个 Worker 抢同一个任务

后台工作进程可能会意外处理同一个待处理任务。

两个工作进程都在任何一方认领之前读取了任务状态。它们各自判定任务可用并开始工作。随后,应用可能会发送两次通知,或生成两份相同的报告。

任务认领机制需要在工作开始前协调所有权。

较旧的搜索响应覆盖较新的响应

由于响应乱序到达,浏览器可能显示错误的搜索结果。

你输入“Django”,然后在第一个响应返回前将搜索词改为“Django transactions”。第二个响应先到达并显示了你现在想要的结果。如果第一个响应随后到达且未经检查直接覆盖页面,界面就会显示旧查询的结果。

请求标识符或过期响应检查可解决此问题,因此数据库行锁并非恰当的修复方案。

信用额度检查为何会失败

首先,让我们从本教程开头的示例明确规则:每个图像请求消耗一个信用点。信用点代表应用的用量额度,与模型处理的 Token 无关。

这是教程设定的规则。如果账户初始拥有十个信用点,接受九个请求后,额度刚好剩一个。

为接受请求,后端必须:

  1. 读取账户余额。

  2. 检查是否至少还剩一个信用点。

  3. 扣除一个信用点并保存余额。

  4. 记录已接受的请求。

若逐条测试请求,这一流程看似正确。第一个请求将余额保存为零,下一个请求读取到零便停止。

下一步是检查这些操作在执行重叠时的情形。

两个请求重叠时发生什么

假设请求 A 读取账户发现剩余一个信用点。在它更新数据库前,请求 B 读取同一账户,也发现剩余一个信用点。

此时每个请求都持有余额的副本。两者都通过检查,并各自计算出 1 - 1 = 0。请求 A 保存零值并记录一次生成,请求 B 随后也保存零值并记录另一次生成。

最终余额为零,但应用接受了两个请求。仅检查余额是否为负数无法发现这种特定故障。

错误在于:在另一个请求有机会修改某值之后,仍盲目信任该值。为了直观说明这一点,我们先构建包含此缺陷的版本。

如何初始化 Django 项目

创建项目目录和虚拟环境。下方的激活命令适用于 Linux 和 macOS:

mkdir django-credit-demo
cd django-credit-demo
python3.12 -m venv .venv
source .venv/bin/activate

这些命令为示例代码提供了独立的目录和 Python 环境。

mkdir 创建目录,cd 进入该目录。venv 命令创建一个名为 .venv 的隔离环境。source 命令激活它,确保后续安装的包归属于此项目。

在执行以下命令时,请保持该环境处于激活状态。

在 Windows 上,使用 py -3.12 -m venv .venv 创建环境,并在 PowerShell 中运行 .venv\Scripts\Activate.ps1 进行激活。

创建 requirements.txt 文件:

Django>=5.2,<5.3
djangorestframework>=3.16,<3.17
psycopg[binary]>=3.2,<3.3

该文件列出了项目所需的三个包。

Django 提供模型和数据库工具,Django REST Framework 处理接口端点。Psycopg 提供 PostgreSQL 连接,其中 [binary] 参数指定使用预编译的二进制实现。每个版本范围允许在选定的发行系列内进行更新,同时排除下一个系列。

使用单一的 requirements 文件能让依赖关系对任何跟随本指南的人一目了然。

安装依赖,创建项目,并添加一个名为 credits 的应用:

python -m pip install -r requirements.txt
python -m django startproject config .
python manage.py startapp credits

这些命令用于安装依赖并创建应用结构。

pip install -rrequirements.txt 读取包列表。startproject 命令创建 config 包和 manage.py,末尾的点号表示选择当前目录。startapp 命令创建 credits 包,模型、服务函数、视图和测试都将放置于此。

现在,我们拥有一个准备就绪、可连接数据库的 Django 项目。

启动 PostgreSQL

manage.py 旁边创建 compose.yaml

services:
  db:
    image: postgres:17
    environment:
      POSTGRES_DB: credit_demo
      POSTGRES_USER: credit_demo
      POSTGRES_PASSWORD: local-demo-only
    ports:
      - "127.0.0.1:5433:5432"
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U credit_demo -d credit_demo"]
      interval: 2s
      timeout: 5s
      retries: 15

这份配置定义了一个用于本地开发的 PostgreSQL 容器。

image 指定使用 PostgreSQL 17,环境变量设置了数据库名和本地凭据。端口映射让数据库只暴露在你本机的回环地址上,端口为 5433。健康检查每两秒运行一次 pg_isready,单次检查超时五秒,连续失败十五次后才将容器标记为不健康。

Django 稍后会在配置文件中使用相同的连接信息。

启动命令:

docker compose up -d --wait

这条命令会启动 compose.yaml 中定义的数据库。

up 负责创建并启动服务;-d 让它在后台运行,你仍可继续使用终端;--wait 则会等待服务通过健康检查变为健康状态。

本教程全程请保持该数据库容器运行。

配置 Django

config/settings.py 替换为下面这份精简配置:

import os

SECRET_KEY = "local-tutorial-only-do-not-use-in-production"
DEBUG = True
ALLOWED_HOSTS = ["localhost", "127.0.0.1", "testserver"]
INSTALLED_APPS = [
    "django.contrib.auth",
    "django.contrib.contenttypes",
    "rest_framework",
    "credits",
]
MIDDLEWARE = []
ROOT_URLCONF = "config.urls"
DEFAULT_AUTO_FIELD = "django.db.models.BigAutoField"
USE_TZ = True
DATABASES = {
    "default": {
        "ENGINE": "django.db.backends.postgresql",
        "NAME": os.environ.get("DB_NAME", "credit_demo"),
        "USER": os.environ.get("DB_USER", "credit_demo"),
        "PASSWORD": os.environ.get("DB_PASSWORD", "local-demo-only"),
        "HOST": os.environ.get("DB_HOST", "127.0.0.1"),
        "PORT": os.environ.get("DB_PORT", "5433"),
        "OPTIONS": {"options": "-c lock_timeout=5000 -c statement_timeout=10000"},
    }
}
REST_FRAMEWORK = {
    "DEFAULT_AUTHENTICATION_CLASSES": [
        "rest_framework.authentication.BasicAuthentication",
    ],
    "DEFAULT_PERMISSION_CLASSES": [
        "rest_framework.permissions.IsAuthenticated",
    ],
}

这份 settings 文件将我们小型应用的各个部分连接起来。

INSTALLED_APPS 启用了 Django 的用户支持、Django REST Framework 以及我们的 credits 应用。ROOT_URLCONF 指向路由定义,DEFAULT_AUTO_FIELDUSE_TZ 分别配置自动标识符和时区感知日期。空的 MIDDLEWARE 列表保持了该 API 示例的极简性,ALLOWED_HOSTS 则接受本地地址和测试客户端主机。

这些设置是针对本教程端点量身定制的。

DATABASES 部分告诉 Django 如何连接 PostgreSQL。

ENGINE 选择 PostgreSQL 后端。每个 os.environ.get() 都会读取可选的环境变量,若缺失则回退到对应的容器默认值。选项设置了 5 秒锁超时和 10 秒语句超时,因此停滞的操作会抛出错误,而非无限期等待。

超时属于错误路径,而此小型端点并未为其提供自定义响应。

REST_FRAMEWORK 部分要求用户必须经过身份验证。

BasicAuthentication 读取请求中携带的凭据,IsAuthenticated 在视图处理请求之前拒绝匿名用户。示例中的密钥和 DEBUG = True 仅适用于本地开发环境。

部署到生产环境时,需配置正式密钥、合适的认证机制以及加密连接。

暂时将 config/urls.py 替换为空的路由列表。待信用逻辑修复后,我们再来添加端点:

urlpatterns = []

这个空列表暂时让 Django 没有应用路由。

Django 会从 ROOT_URLCONF 指定的模块中读取 urlpatterns。由于此极简配置未启用 admin 应用,我们移除了生成的管理路由。即使端点尚未存在,仍可执行数据库命令。

待受保护的视图就绪后,再替换此列表。

创建 Models

credits/models.py 中添加以下 models:

from django.conf import settings
from django.db import models

class CreditAccount(models.Model):
    user = models.OneToOneField(settings.AUTH_USER_MODEL, on_delete=models.CASCADE)
    balance = models.PositiveIntegerField(default=0)

class Generation(models.Model):
    account = models.ForeignKey(CreditAccount, on_delete=models.CASCADE)
    prompt = models.CharField(max_length=500)
    status = models.CharField(max_length=20, default="reserved")
    created_at = models.DateTimeField(auto_now_add=True)

CreditAccount 存储与用户关联的余额。

settings.AUTH_USER_MODEL 指向项目中配置的用户模型。OneToOneField 确保每个用户最多一个账户,on_delete=models.CASCADE 指示 Django 在删除用户时级联删除对应账户。PositiveIntegerField(default=0) 用于存储非负整数余额,初始值为 0。

balance 字段代表我们的配额,但仅靠其类型无法强制实现每次生成消耗一个信用。

Generation 记录应用已接受的请求。

外键把每条记录关联到其付费账户,使一个账户可以拥有多个 generation。prompt 最多保存 500 个字符,created_at 记录该行的创建时间。status 初始值为 reserved,表示我们的服务已记录该请求,但尚未完成模拟的图像生成工作。

保留 generation 记录,是为了能查证扣除的积分到底用在了哪里。

创建并应用迁移:

python manage.py makemigrations credits
python manage.py migrate

这两条命令会把模型定义转换成数据库表。

makemigrations credits 生成描述模型变更的迁移文件;migrate 应用待处理的迁移,包括 Django 的用户表和我们新建的两张表。数据库连接来自刚才配置的数据库设置。

至此,数据库已经可以存储账户及其 generation 请求了。

如何复现竞态条件

创建 credits/services.py,添加下面这个刻意不加保护、不安全的函数:

from .models import CreditAccount, Generation

class InsufficientCredits(Exception):
    pass

def reserve_generation_unsafe(user_id, prompt):
    account = CreditAccount.objects.get(user_id=user_id)

    if account.balance < 1:
        raise InsufficientCredits

    account.balance -= 1
    account.save(update_fields=["balance"])

    return Generation.objects.create(account=account, prompt=prompt)

这个函数实现了积分检查逻辑,但没有做任何并发保护。

objects.get() 根据给定的 user_id 获取账户,如果余额不足 1,raise InsufficientCredits 会中断函数执行。减法操作只修改了 Python 对象,随后 save(update_fields=["balance"]) 把这个值写回数据库。最后,Generation.objects.create() 把接受的请求写入数据库,并返回对应的模型对象。

读取和写入是两个独立的操作,因此其他请求可能在这两者之间插入执行。

InsufficientCredits 让调用方能够针对具体的失败情况做处理。

它是一个继承自 Python Exception 的自定义异常类。pass 语句表示我们不为这个类添加任何行为。稍后,视图会捕获这个异常并返回有用的响应。

保留那个故意存在缺陷的函数以便对照,但不要让端点调用它。

确保两次读操作都先于写操作完成

同时打开两个浏览器标签页并不能可靠地复现该 bug。因为一个请求可能在另一个请求读取余额之前就已执行完毕。因此,我们将使用两个 Django shell 会话,并在每次读取账户后暂停执行。

在项目目录中打开终端,激活虚拟环境,然后启动 shell:

python manage.py shell

该命令会启动一个已加载 Django 项目的 Python shell。

它使用的是与 manage.py 关联的配置。你可以导入模型并直接查询配置好的数据库。每打开一个终端,就相当于为一个独立的 shell 会话做准备,用于我们的实验。

我们将利用这些 shell 会话来控制数据库操作的执行顺序。

创建一个新用户和一个含有一次点数的账户:

from django.contrib.auth import get_user_model
from credits.models import CreditAccount, Generation

user = get_user_model().objects.create_user(username="race-demo")
CreditAccount.objects.create(user=user, balance=1)

account_a = CreditAccount.objects.get(user__username="race-demo")
print(account_a.balance)

这段代码为第一次操作做准备,初始余额设为 1。

get_user_model() 获取配置的用户模型类,create_user() 插入我们的演示用户。创建账户时为该用户分配了 1 点。查询时使用 user__username 沿用户关系追踪,并将获取到的账户对象存储在 account_a 中。

打印结果应显示为1。保持此 shell 会话开启,但不要执行账户更新操作。

在第二个终端中激活相同的虚拟环境,再次运行 python manage.py shell。在那里也读取账户:

from credits.models import CreditAccount, Generation

account_b = CreditAccount.objects.get(user__username="race-demo")
print(account_b.balance)

第二个 shell 会话将同一数据库行读取到了一个不同的 Python 对象中。

account_b 属于当前 shell 会话,与 account_a 相互独立。由于第一个 shell 会话尚未保存扣减操作,此处的查询结果也应返回余额为 1。后续在另一个 shell 中发生的更改不会自动刷新此对象。

此时,两个操作都各自持有一份它打算消耗的额度副本。

回到第一个 shell,运行:

if account_a.balance >= 1:
    account_a.balance -= 1
    account_a.save(update_fields=["balance"])
    Generation.objects.create(account=account_a, prompt="A garden")

第一个 shell 检查并扣除了它所持有的余额副本。

if 条件判断通过,因为 account_a.balance 当前为 1。减法将其变为 0,随后 save() 将 0 写入数据库中的账户行。最后一行为“花园”提示词创建了一条生成记录。

在空行按下 Enter 结束代码块,然后再切换到另一个终端。

接着在第二个 shell 中运行对应的代码块:

if account_b.balance >= 1:
    account_b.balance -= 1
    account_b.save(update_fields=["balance"])
    Generation.objects.create(account=account_b, prompt="A beach")

第二个 shell 基于它之前加载的对象来做决定。

由于尚未刷新 account_b,它的 if 条件判断仍看到余额为 1。它执行减法并保存 0,用同样的值覆盖了第一个操作写入的余额。随后,它为“海滩”提示词创建了另一条独立的生成记录。

第二个操作接受了工作,但它使用的额度实际上已被第一个操作消耗。

最后,在第二个 shell 中检查数据库状态:

account_b.refresh_from_db()
print(account_b.balance)
print(Generation.objects.filter(account=account_b).count())

这段代码用于查看两个操作在数据库中留下的最终结果。

refresh_from_db() 重新加载账户,使 print 输出反映数据库中的真实值。经过过滤的 count() 仅统计关联到该账户的生成记录数。你应该看到额度为 0,而生成记录数为 2。

生成记录数揭示了仅看余额时会掩盖的问题。

通过两个 shell,我们可以刻意重现这一不安全序列。

我们在每次读取后暂停,随后允许两个写入同时发生。当请求重叠时,服务器也可能产生同样的顺序,尽管并非每次尝试都会如此。若要重复该实验,请新建用户名和账户,以免之前的记录影响计数结果。

现在,我们可以针对观察到的具体缺口构建防护机制。

如何保护余额

这个函数涉及两个数据库问题:扣费和生成记录必须同时成功;并发请求之间也需要协调的方式来检查和更新账户。

下面按这个顺序逐一解决。

把相关变更放进同一个事务

数据库事务会把一组操作打包成一个工作单元。Django 的 transaction.atomic() 会在代码块成功执行完后提交变更,一旦异常导致代码块中断,则回滚所有操作。

这一点很重要:不安全的写法先保存余额、再创建生成记录,一旦创建记录失败,扣掉的额度就可能凭空消失,却没有对应的生成记录。

作为演示,把操作包进事务后是这样的:

from django.db import transaction

@transaction.atomic
def reserve_generation_atomic_only(user_id, prompt):
    account = CreditAccount.objects.get(user_id=user_id)

    if account.balance < 1:
        raise InsufficientCredits

    account.balance -= 1
    account.save(update_fields=["balance"])

    return Generation.objects.create(account=account, prompt=prompt)

装饰器把函数内的数据库操作放进了一个原子事务。

from django.db import transaction 引入了 Django 的事务工具。函数的逻辑不变:先读取账户、检查余额、保存扣费,最后创建生成记录。如果创建记录时抛出异常,事务会连同扣费操作一起回滚。

这样就保住了扣费和生成记录之间的对应关系。

但普通的账户查询仍然无法保护余额检查环节。

PostgreSQL 默认的隔离级别是 Read Committed。在这个级别下,事务中的普通读取不会让其他事务等待。所以两个请求都可能在任何一方更新之前读到同一个余额。

在做出余额判断之前,我们需要先协调对数据的访问。

在检查余额之前先锁定账户

行锁是一种数据库机制,在事务持有锁期间,限制其他事务对所选行的冲突操作。

在这个例子中,一行记录代表一个信用账户。SELECT FOR UPDATE 锁会让竞争的更新和冲突的锁请求等待,直到锁释放。普通读取操作仍可继续,事务也能处理其他账户的行。

这让我们能在检查并修改余额时保护单个账户。

credits/services.py 顶部添加 from django.db import transaction。保留原不安全函数以便对比,然后添加此受保护版本:

@transaction.atomic
def reserve_generation(user_id, prompt):
    account = CreditAccount.objects.select_for_update().get(user_id=user_id)

    if account.balance < 1:
        raise InsufficientCredits

    account.balance -= 1
    account.save(update_fields=["balance"])

    return Generation.objects.create(account=account, prompt=prompt)

该函数在检查余额前先获取账户的行锁。

select_for_update() 请求锁,.get(user_id=user_id) 在事务内执行该账户的查询。一旦持有锁,函数就检查余额,必要时抛出 InsufficientCredits。否则,它会在事务完成前保存扣减并创建生成记录。

现在的决策及其数据库更改都在账户受保护状态下进行。

对此函数的竞争调用必须在锁查询处等待。

如果第一个调用提交了扣减,等待的调用会在 Read Committed 设置下检查更新后的余额。如果第一个调用回滚,其扣减不会保留。我们配置的超时也能通过错误终止等待。

对于只有一个信用额度且首次提交成功的情况,第二次调用会读取到零并拒绝请求。

这种保护需要覆盖所有消耗余额的路径。

较旧的函数可能仍会在不请求锁的情况下读取账户。其后续更新会在我们持有锁时等待,但随后可能用陈旧值覆盖余额。因此,管理调整和后台作业也需要安全更新策略。

一个受保护的函数无法修正其他地方不安全的写入者。

将图像生成放在事务外

事务应当涵盖积分预留。远程图片请求可能耗时远长于这些数据库操作,因此将其置于事务内会导致其他竞争请求不必要地等待。

在本示例中,向 credits/services.py 添加此函数:

def simulate_generation(generation):
    # No external AI request is made in this tutorial.
    generation.status = "completed"
    generation.save(update_fields=["status"])
    return "Simulated image generation completed."

该函数模拟已接受生成任务的完成状态。

它将传入的 generation 对象的 status 更改为 completed。调用 save() 仅持久化该字段。返回值是一条消息,因此该函数不生成任何图片,也不发起外部请求。

我们会在预留函数返回后调用此函数。

这种调用顺序将图片处理工作保留在预留事务之外。

在此配置下且无外层事务包裹时,预留会在模拟开始前提交。此后若发生故障,系统不会自动返还积分。集成真实服务提供商时,需要制定专门的失败重试或退款策略。

我们稍后在测试预留逻辑本身时,再回过头来讨论这些限制。

如何接收并测试图片请求

现在我们可以将受保护的函数连接到一个端点。向 credits/views.py 添加以下代码:

from rest_framework import serializers, status
from rest_framework.response import Response
from rest_framework.views import APIView
from .models import CreditAccount
from .services import InsufficientCredits, reserve_generation, simulate_generation

class GenerationInput(serializers.Serializer):
    prompt = serializers.CharField(max_length=500)

class GenerateView(APIView):
    def post(self, request):
        serializer = GenerationInput(data=request.data)
        serializer.is_valid(raise_exception=True)
        try:
            generation = reserve_generation(
                request.user.pk, serializer.validated_data["prompt"]
            )
        except CreditAccount.DoesNotExist:
            return Response({"detail": "Credit account not found."}, status=404)
        except InsufficientCredits:
            return Response({"detail": "Not enough credits."}, status=409)

        result = simulate_generation(generation)
        return Response(
            {"id": generation.pk, "status": generation.status, "result": result},
            status=status.HTTP_201_CREATED,
        )

在视图扣除积分之前,serializer 会先校验 prompt。

GenerationInput 声明了一个必填文本字段,最长 500 个字符。is_valid(raise_exception=True) 会在输入缺失、为空或非法时直接返回校验错误。视图随后从 validated_data 中取出清洗后的值,连同 request.user.pk(已认证用户的数据库 ID)一起传给预留函数。

prompt 由客户端提供,而账户则由服务端根据已认证用户来确定。

视图把预留结果转换为响应。

账户不存在返回 404InsufficientCredits 返回 409,这是我们为余额冲突选定的状态码。预留成功后,视图再调用模拟函数,返回 201,附带记录 ID、完成状态和模拟结果。

通过这些分支,调用方就能区分任务是已接受还是被拒绝。

config/urls.py 替换为:

from django.urls import path
from credits.views import GenerateView

urlpatterns = [
    path("api/generate/", GenerateView.as_view()),
]

这条路由把请求地址映射到我们的视图。

path() 匹配地址中的 api/generate/ 部分。GenerateView.as_view() 将基于类的视图转换为可被 Django 调用的对象。该视图的 post() 方法负责处理此路由下的 POST 请求。

服务器运行时,该端点即通过 /api/generate/ 访问。

逐次发送请求测试

打开 Django shell,为端点演示创建一个独立用户:

from django.contrib.auth import get_user_model
from credits.models import CreditAccount

user = get_user_model().objects.create_user(
    username="api-demo",
    password="local-example-password",
)
CreditAccount.objects.create(user=user, balance=1)

此代码块为受认证的端点示例创建了一个用户。

create_user() 会保存用户名并对提供的密码进行哈希处理。CreditAccount.objects.create() 为新用户赋予一信用额度。使用独立的用户名可确保该检查与双 shell 实验中使用的账户相互独立。

以下凭据仅适用于此本地演示账户。

退出 shell,启动开发服务器:

python manage.py runserver

此命令启动 Django 开发服务器。

若不指定地址参数,它默认监听 127.0.0.1:8000。发送到该地址的请求将经过我们刚添加的路由配置。在从另一个终端发送请求期间,请保持此终端处于打开状态。

此服务器仅用于本地演示。

在另一个终端中,发送请求:

curl -i -u api-demo:local-example-password \
  -H "Content-Type: application/json" \
  -d '{"prompt": "A garden at sunrise"}' \
  http://127.0.0.1:8000/api/generate/

此命令使用演示用户的凭据,向该端点提交提示词。

-i 包含响应头,-u 提供 Basic 认证所需的用户名和密码。-H 将请求体格式声明为 JavaScript Object Notation (JSON)。-d 提供该请求体,并让 curl 向指定地址发送 POST 请求。

第一次请求应返回 201 状态码及模拟的完成结果。

再次运行该命令,可检查首次扣费后的账户状态。

下一个请求的查询结果应为零积分。视图应返回 409 状态码,并附带 "detail": "Not enough credits."。由于我们在两次请求之间做了等待,这一步验证了常规的串行执行流程,并未涉及并发场景。

接下来,我们将测试并发重叠的操作。

编写自动化测试

用以下导入语句和辅助类替换 credits/tests.py 中的现有内容。我们稍后在后续步骤中添加具体的测试方法。

from concurrent.futures import ThreadPoolExecutor
from threading import Barrier
from unittest.mock import patch

from django.contrib.auth import get_user_model
from django.db import (
    OperationalError, close_old_connections, connection, connections, transaction,
)
from django.test import TransactionTestCase
from rest_framework.test import APIClient

from .models import CreditAccount, Generation
from .services import InsufficientCredits, reserve_generation, reserve_generation_unsafe

class ConcurrencyTests(TransactionTestCase):
    def setUp(self):
        if connection.vendor != "postgresql":
            self.skipTest("Run these concurrency tests on PostgreSQL.")

        self.user = get_user_model().objects.create_user(username="parallel-reader")
        self.account = CreditAccount.objects.create(user=self.user, balance=1)

    def run_two(self, action):
        def worker():
            close_old_connections()
            try:
                return action()
            finally:
                connections.close_all()

        with ThreadPoolExecutor(max_workers=2) as pool:
            futures = [pool.submit(worker) for _ in range(2)]
            return [future.result(timeout=15) for future in futures]

该测试类会在每个并发测试开始前准备一个全新的账户。

setUp() 会在数据库连接不是 PostgreSQL 时跳过测试。随后,它会创建一个用户和一个仅含 1 个积分的账户。使用 TransactionTestCase 而非普通的 TestCase,是因为前者允许真正的事务边界,而后者内部的事务可能会掩盖锁使用方面的错误。

在 SQLite 上跳过的测试无法验证 PostgreSQL 的锁机制行为。

run_two() 辅助方法会将同一个操作分配给两个工作线程执行。

ThreadPoolExecutor(max_workers=2) 提供工作线程,每次 pool.submit(worker) 调度一个任务。future 代表该调用的最终结果,通过 future.result(timeout=15) 获取结果或在出错时抛出异常。每个 worker 在执行操作前清理不可用的旧连接,并在操作结束时(包括失败的情况)关闭自己的连接。

这样,两个操作就能通过各自的连接访问同一个数据库。

验证不安全行为

ConcurrencyTests 中添加以下方法,与 run_two() 保持相同的缩进层级:

 def test_reproduce_unsafe_spending(self):
        both_have_read = Barrier(2)
        original_get = CreditAccount.objects.get

        def read_then_wait(*args, **kwargs):
            account = original_get(*args, **kwargs)
            both_have_read.wait(timeout=5)
            return account

        with patch(
            "credits.services.CreditAccount.objects.get",
            side_effect=read_then_wait,
        ):
            self.run_two(
                lambda: reserve_generation_unsafe(self.user.pk, "A garden").pk
            )

        self.account.refresh_from_db()
        self.assertEqual(self.account.balance, 0)
        self.assertEqual(Generation.objects.count(), 2)

这个测试会强制两个不安全操作都先完成读取,再继续执行。

original_get 保存了真实的查询方法,read_then_wait() 先调用它,然后在 Barrier(2) 处等待。栅栏只有在两个 worker 都到达后才放行,如果等待超时则抛出错误。patch() 会在 with 代码块执行期间,用这个包装函数临时替换服务中的查询方法。

数据库读取依然是真实执行的,只是测试控制了读取之后的暂停时机。

最后的断言记录了这个刻意制造的错误结果。

那个小小的 lambda 调用不安全的函数,并为每个 worker 返回所创建记录的 ID。两个操作都返回后,refresh_from_db() 重新加载账户数据。在这个隔离的测试中,断言期望余额为零,且有两条 generation 记录。

这里测试通过,说明成功复现了 bug,而不是不安全函数的行为正确。

检查受保护的接口

在同一个类中添加以下方法:

def concurrent_api_requests(self):
        ready = Barrier(2)

        def send_request():
            client = APIClient()
            client.force_authenticate(self.user)
            ready.wait(timeout=5)
            return client.post(
                "/api/generate/",
                {"prompt": "A garden"},
                format="json",
            ).status_code

        return self.run_two(send_request)

    def test_last_credit_accepts_only_one_request(self):
        self.assertEqual(sorted(self.concurrent_api_requests()), [201, 409])
        self.account.refresh_from_db()
        self.assertEqual(self.account.balance, 0)
        self.assertEqual(Generation.objects.count(), 1)

    def test_two_credits_accept_both_requests(self):
        self.account.balance = 2
        self.account.save(update_fields=["balance"])

        self.assertEqual(self.concurrent_api_requests(), [201, 201])
        self.account.refresh_from_db()
        self.assertEqual(self.account.balance, 0)
        self.assertEqual(Generation.objects.count(), 2)

请求辅助函数通过受保护的视图发送两个经过身份验证的请求。

每个工作线程创建独立的 APIClient 实例,并通过 force_authenticate() 注入测试用户,无需进行密码交换。Barrier 屏障置于 post() 调用之前,确保两个工作线程同时到达请求起始点。每次调用将响应状态码返回给 run_two()

该测试专门验证端点的积分逻辑,而不涉及身份验证机制本身。

两个测试方法分别覆盖了积分不足和积分充足这两种共享积分场景。

单积分测试对状态码进行排序,因为任意一个工作线程都可能有先完成请求。测试预期结果为一个 201 和一个 409。同时检查剩余积分是否为零,以及是否恰好生成一条记录。双积分测试将初始余额设为 2,预期两个请求均成功,生成两条记录,且余额归零。

除了检查响应状态码,还检查数据库记录,有助于捕获被错误接受的生成任务。

Barrier 用于协调请求的发起时机,而非控制所有的数据库操作。

其中一个请求的执行速度仍可能快于另一个。若将 Barrier 移至锁获取之后,第一个工作线程将不得不等待第二个获取不到锁的工作线程。当前的放置方式避免了这种人为死锁,但并不能穷尽所有可能的执行顺序。

将该回归测试与受控复现步骤及记录的数据库行为一并使用。

验证服务遇到被持有的锁

我们还可以测试预订函数在检查余额之前是否会等待锁。

下面的测试从零积分开始,这样余额检查会立即拒绝未保护的读取操作。主连接在启动工作线程前获取账户锁。该工作线程调用实际的预订函数,并设置一秒钟的锁超时。

ConcurrencyTests内部添加此方法:

    def test_reservation_waits_for_account_lock(self):
        self.account.balance = 0
        self.account.save(update_fields=["balance"])
        user_id = self.user.pk

        def attempt_reservation():
            close_old_connections()
            try:
                try:
                    with transaction.atomic():
                        with connection.cursor() as cursor:
                            cursor.execute("SET LOCAL lock_timeout = '1s'")
                        reserve_generation(user_id, "A garden")
                except OperationalError as error:
                    return error.__cause__.sqlstate
                except InsufficientCredits:
                    return "insufficient_credits"
                return "accepted"
            finally:
                connections.close_all()

        with ThreadPoolExecutor(max_workers=1) as pool:
            with transaction.atomic():
                CreditAccount.objects.select_for_update().get(pk=self.account.pk)
                result = pool.submit(attempt_reservation).result(timeout=5)
                self.assertEqual(result, "55P03")

            result = pool.submit(attempt_reservation).result(timeout=5)
            self.assertEqual(result, "insufficient_credits")

        self.account.refresh_from_db()
        self.assertEqual(self.account.balance, 0)
        self.assertEqual(Generation.objects.count(), 0)

该测试控制竞争锁出现的时间。

主连接在等待工作线程结果期间保持其事务开启。在工作线程的事务内,SET LOCAL 临时缩短锁超时。锁等待超时后,工作线程应返回 PostgreSQL 的 55P03 错误码,该码表示 lock_not_available

此时余额校验被拒绝,说明服务没有先获取账户锁就做了检查。

第二次调用是在主事务释放锁之后,对同一操作再次检查。

这一次 worker 能成功拿到账户锁并读取到零余额,应该返回 insufficient_credits 而不是数据库错误。最后的断言确认两次尝试中扣款记录和生成记录都没有发生变化。

两次调用共同验证了竞争处理和锁释放逻辑,而不必依赖两个请求碰巧同时发生。

把余额设为零是刻意为之,这样测试更有针对性。

如果余额是一笔信用额度,缺乏保护的函数在最终更新该行时仍可能碰到锁;而余额为零时,它会在任何更新之前就拒绝请求,导致预期的超时断言失败。因此这个测试验证的是服务“先加锁后检查”的行为,而前面的端点测试验证的则是消费结果。

这个受控锁测试已作为配套项目 PostgreSQL 测试套件的一部分通过。

检查失败后的回滚

ConcurrencyTests 之外再加一个类:

class CreditRollbackTests(TransactionTestCase):
    def test_failed_record_creation_restores_credit(self):
        user = get_user_model().objects.create_user(username="rollback-reader")
        account = CreditAccount.objects.create(user=user, balance=1)

        with patch(
            "credits.services.Generation.objects.create",
            side_effect=RuntimeError("Simulated record creation failure"),
        ):
            with self.assertRaises(RuntimeError):
                reserve_generation(user.pk, "A garden")

        account.refresh_from_db()
        self.assertEqual(account.balance, 1)
        self.assertEqual(Generation.objects.count(), 0)

这个测试故意在创建生成记录时抛出异常。

它先创建一个只有一笔信用额度的账户。通过 patch 让 Generation.objects.create() 抛出 RuntimeErrorassertRaises() 确认这个错误会从预约函数中抛出。事务回滚之后,重新加载的账户应该仍有一笔信用额度,且没有任何生成记录。

这验证了预约失败不会留下已扣款的状态。

在 PostgreSQL 仍可用的前提下运行测试:

python manage.py test credits -v 2

该命令运行 credits 应用内的测试。

-v 2 参数让 Django 展示每个测试及其结果。Django 会创建一个独立的测试数据库,因此数据库用户需要具备创建数据库的权限。我们本地容器中的用户拥有此权限,但现有的数据库配置可能需要调整。

在标准的 PostgreSQL 环境下,预期会有五个测试通过。在实际依赖该示例前,请通过运行来确认结果。

常见错误与后续步骤

示例代码已为信用决策提供了保护,但在改编时,有几个细节容易被忽视。

在锁内读取账户

如果在事务之前获取账户,并继续使用该对象,可能会导致使用过时的余额。受保护的函数特意通过 select_for_update() 在做出决策前获取账户。

将代码迁移到其他服务或端点时,请保持这一顺序。传入用户标识符后,函数自身负责执行最新且已加锁的读取。

使用数据库作为共享协调点

禁用提交按钮可以减少误点,但第二个标签页或其他客户端仍可能发送请求。Python 线程锁也仅协调同一进程中共享该锁的代码。

如果运行多个应用工作进程,信用决策仍需在共享数据库层面进行保护。请审查其他余额更新操作,包括管理端调整,不要假设该端点是唯一的写入者。

简单计数器考虑使用条件更新

行锁是一种方案。对于简单的扣减操作,Django 还可以使用 F() 表达式在一条数据库更新语句中同时表达余额条件和减法操作。

条件很关键:无条件减法无法确保信用充足。如果同时创建生成记录,请将扣减操作和记录创建放在同一个事务中。

此处使用显式锁是因为它能清晰展示读取、决策和更新的流程。理解操作需保护什么之后,你可以进一步探索条件更新。

区分并发请求与重试

我们的示例将两次提交视为两个不同的请求。如果用户拥有两份信用额度,两者都应成功。

重试引入了一个新的需求。应用可能已经接受了生成请求,但在客户端收到结果前丢失了响应。如果客户端重新发送同一个逻辑请求,你可能希望直接返回原始结果,而无需再次扣费。

这需要幂等性:一种识别重复操作并避免重复应用其结果的方法。典型设计包括请求键、数据库唯一性约束以及存储的结果。仅靠余额锁无法识别重复的意图。

应对供应商故障

一旦预扣款提交,后续供应商的故障并不会自动恢复信用。你需要制定策略,决定是重试生成、退款,还是保持待定状态以便后续恢复。

实际实现还需处理应用在预扣信用后、启动任务前停止的情况。reserved 记录提供了可追踪的对象,但本教程并未实现持久化任务队列或恢复工作进程。

在扩展示例时,请将这些考量保持可见。防止并发超支只是完整信用体系的一部分。

总结与下一步

在本指南开头,两个请求都能通过余额检查并消耗同一笔信用。余额最终归零,这使得该故障很容易被忽视。

我们复现了该序列,随后通过事务和行锁保护了决策逻辑。我们还检查了生成计数,测试了两个请求均具备足够信用的情况,并添加了记录创建失败的回滚测试。

审查自己应用中类似功能时:

  1. 明确数据必须满足的规则,例如每接受一次生成请求只扣减一笔信用。

  2. 识别不同请求可能读取和修改同一数值的位置。

  3. 保护决策及其相关的数据库变更。

  4. 在实际上使用的数据库上测试重叠操作及失败场景。

同样的逻辑也适用于在线商店的最后一件商品或任何共享配额。只有当应用能安全地基于检查结果执行动作时,成功的检查才有意义。

参考

若想深入了解重试机制,请参见 Making retries safe with idempotent APIs

原始来源: freeCodeCamp

评论 (0)