一家猫粮和宠物用品网店的顾客在火车上用手机购买三文鱼猫粮和猫砂。应用发送 POST /orders/ord_1001/pay,服务器已经通过支付服务商扣款,但响应没有到达手机。应用无法区分请求丢失和响应丢失,因此重试。缺少保护时,重试可能以两种方式出错:
- 重试在第一次请求完成后到达:订单已经是
paid,顾客明明付款成功,却收到错误。 - 重试在第一次请求仍等待支付服务商时到达:两次请求都通过“订单是否仍待支付”的检查,银行卡被扣款两次。 幂等键可以解决这两种情况。客户端为每次付款尝试生成唯一键,并在所有重试的
Idempotency-Key请求头中发送同一个键。服务器针对一个键只运行一次处理器,存储结果,并用该结果回答后续重试。@nestjs/idempotency通过@Idempotent()装饰器为 HTTP 路由、GraphQL mutation 和微服务处理器提供该能力。 其 HTTP 行为遵循 IETF 的 Idempotency-Key 请求头草案。 本教程将为付款端点加入幂等性,让移动应用安全重试,把键迁移到 PostgreSQL 或 Redis,加密已保存的收据,保护 GraphQL mutation,并对配送服务消费的order.paid事件去重。
前提条件
教程从已经具备以下功能的订单 API 开始:
OrdersModule中包含OrdersService和OrdersController。Order具有id、userId、items、以分为单位的total,以及取值为'pending'、'paid'、'shipped'或'refunded'的status。PaymentProviderClient封装支付服务商 SDK。银行卡被拒绝时抛出状态码 402 的HttpException;服务商不可用时抛出ServiceUnavailableException(503)。AuthGuard验证 bearer token 并设置req.user,参见身份验证;另有读取该值的@CurrentUser()参数装饰器。- 全局
ValidationPipe,以及以SHIPPING_SERVICE注册的配送服务 TCP 客户端。 - 数据通过 Drizzle 存在 PostgreSQL 中,数据库由
@nestjs/drizzle注册,参见 Drizzle。若使用 TypeORM,后面的 PostgreSQL 小节也给出了对应实现。 需要保护的方法如下:
// 文件:orders/orders.service
async pay(orderId: string, user: User, dto: PayOrderDto): Promise<PaymentReceipt> {
const order = this.ordersRepository.findOwned(orderId, user.id);
if (order.status !== 'pending') {
throw new ConflictException(`Order ${order.id} is already ${order.status}`);
}
const charge = await this.paymentProviderClient.charge({
amount: order.total,
paymentMethod: dto.paymentMethod,
reference: order.id,
});
this.ordersRepository.markPaid(order.id);
this.shipping.emit('order.paid', {
eventId: `order-paid:${order.id}`,
orderId: order.id,
});
return {
orderId: order.id,
receiptId: charge.id,
amount: charge.amount,
cardLast4: charge.cardLast4,
paidAt: charge.createdAt,
};
}
仅检查状态不能防止重复扣款:同时到达的两次请求,可能在任何一次调用 markPaid() 之前都看到 pending。
安装包:
$ npm i --save @nestjs/idempotency
让付款端点具有幂等性
在根模块中注册一次 IdempotencyModule:
// 文件:app.module
import { Module } from '@nestjs/common';
import { IdempotencyModule } from '@nestjs/idempotency';
import type { User } from './orders/order.js';
import { OrdersModule } from './orders/orders.module.js';
@Module({
imports: [
IdempotencyModule.forRoot({
// Runs after guards, so req.user is set. Keys never collide across users.
scope: (req: { user?: User }) => req.user?.id,
}),
OrdersModule,
],
})
export class AppModule {}
该模块是全局模块,会注册一个应用级拦截器。对于没有标记 @Idempotent() 的处理器,它不做任何处理,因此不会改变其他路由的行为。 scope 为每位用户的键提供独立命名空间。没有它,所有客户端共用一个命名空间,一位用户选择的键可能重放另一位用户的收据。使用它后,Alice 和 Bob 都可以发送 Idempotency-Key: 1,而不会看到对方结果。守卫先于拦截器运行,因此调用 scope 时 req.user 已设置;未认证请求会在占用键之前被拒绝。 已登录用户调用没有配置 scope 的 @Idempotent() 处理器时,模块会按处理器记录一次警告。若确实需要共享命名空间,例如使用服务商事件 ID 作为键的 webhook,可明确设置 scope: false。 接下来标记处理器:
// 文件:orders/orders.controller
import { Body, Controller, Param, Post, UseGuards } from '@nestjs/common';
import { Idempotent } from '@nestjs/idempotency';
import { AuthGuard } from '../auth/auth.guard.js';
import { CurrentUser } from '../auth/current-user.decorator.js';
import { PayOrderDto } from './dto/pay-order.dto.js';
import type { User } from './order.js';
import { OrdersService } from './orders.service.js';
@Controller('orders')
@UseGuards(AuthGuard)
export class OrdersController {
constructor(private readonly ordersService: OrdersService) {}
@Post(':id/pay')
@Idempotent({ required: true })
pay(
@Param('id') id: string,
@Body() dto: PayOrderDto,
@CurrentUser() user: User,
) {
return this.ordersService.pay(id, user, dto);
}
}
required: true 会在处理器运行前,拒绝未提供键的请求,返回 400 和 IDEMPOTENCY_KEY_REQUIRED。没有键的付款无法安全重试,因此应拒绝这种请求。
携带键的请求按以下过程处理:
- 拦截器计算请求的指纹:对 scope、方法、包含订单 ID 的 URL,以及按对象键排序的请求体求 SHA-256 哈希。
- 原子地获取记录键
<scope>:<key>的锁,例如usr_alice:5e0f6d0e-…。如果记录已存在,不执行处理器,具体响应见下表。处理器运行期间持续续租,避免耗时处理器的锁被重试抢走。 - 执行处理器。完成后,锁转为已完成记录,保存状态码、响应体,以及处理器或中间件设置的少量允许重放的响应头:
Location、Content-Type、Content-Language、Content-Location、ETag、Last-Modified。这些头描述的是已存储的响应体,所以重试即使发送了不同的Accept-Language,仍会得到原响应体对应的Content-Language。 Cookie、CORS 和限流响应头属于原始请求,不会重放。 后续使用相同键的请求会得到以下响应之一: | 情况 | 响应 | | ————————————– | ————————————————————————————— | | 首次请求仍在运行 | 409IDEMPOTENCY_KEY_IN_USE,附带Retry-After| | 首次请求已完成 | 保存的状态码和响应体,以及Idempotent-Replayed: true;不执行处理器 | | 键相同,但请求体或 URL 不同 | 422IDEMPOTENCY_KEY_REUSED| 错误也可以被存储。后文“调整 TTL 和失败处理”会说明具体规则。 因为拦截器是全局的,它在控制器和处理器拦截器的外层运行。存储的是序列化器及其他拦截器处理后客户端实际收到的响应体;重放会跳过这些处理。管道在它内部运行,因此ValidationPipe的 400 与其他客户端错误一样会被存储。其他全局拦截器则需另外考虑,详见生产检查部分。 在改用 PostgreSQL 前,记录保存在默认InMemoryIdempotencyStore中。它适合开发时的单实例,但各 API 实例各自持有记录,重启后也会丢失。NODE_ENV=production时,应用会拒绝以这种方式启动,并在错误中说明如何注册存储。
从移动应用发送键
服务器只能对客户端一致发送的内容去重。应用遵守两条规则:
- 顾客点击付款时创建键,而不是每个请求都创建。即使应用被杀死再重启,同一次付款的重试也使用同一个键,因此需要把键和待完成付款一起保存。
- 只有顾客做出新的决定时才创建新键,例如银行卡被拒后选择另一张卡。旧键搭配不同请求体会得到 422。
// 文件:mobile/pay-order
export interface Receipt {
orderId: string;
receiptId: string;
amount: number;
cardLast4: string;
paidAt: string;
}
export interface PendingPayment {
orderId: string;
paymentMethod: string;
idempotencyKey: string;
}
export class PaymentError extends Error {
constructor(
readonly status: number,
readonly body: { code?: string; message?: string },
) {
super(body.message ?? `Payment failed with status ${status}`);
}
}
// Call when the customer taps "Pay", and save the result with the order
// (for example in AsyncStorage). Retries, even after an app restart, reuse it.
export function startPayment(orderId: string, paymentMethod: string): PendingPayment {
return { orderId, paymentMethod, idempotencyKey: crypto.randomUUID() };
}
export async function payOrder(
apiUrl: string,
token: string,
payment: PendingPayment,
maxAttempts = 5,
): Promise<Receipt> {
for (let attempt = 1; ; attempt++) {
const backoff = 250 * 2 ** (attempt - 1);
let res: Response;
try {
res = await fetch(`${apiUrl}/orders/${payment.orderId}/pay`, {
method: 'POST',
headers: {
Authorization: `Bearer ${token}`,
'Content-Type': 'application/json',
'Idempotency-Key': payment.idempotencyKey,
},
body: JSON.stringify({ paymentMethod: payment.paymentMethod }),
signal: AbortSignal.timeout(15_000),
});
} catch (err) {
// No response: the card may or may not have been charged. Same key, try again.
if (attempt >= maxAttempts) {
throw err;
}
await sleep(backoff);
continue;
}
if (res.ok) {
return res.json(); // the first result, even if Idempotent-Replayed: true
}
const body = await res.json().catch(() => ({})); // a proxy's 502 page isn't JSON
const retryable =
res.status >= 500 || // the server released the key
res.status === 429 ||
body.code === 'IDEMPOTENCY_KEY_IN_USE'; // the first attempt is still running
if (!retryable || attempt >= maxAttempts) {
throw new PaymentError(res.status, body);
}
const retryAfter = Number(res.headers.get('Retry-After'));
await sleep(retryAfter > 0 ? retryAfter * 1000 : backoff);
}
}
const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));
浏览器和 Node.js 可使用 crypto.randomUUID();React Native 应使用 UUID 库,例如 expo-crypto 的 randomUUID()。
应用可能得到以下响应,它们对重试的含义如下: | 响应 | 含义 | 是否使用同一个键重试 | | ————————————————————- | ———————————————————————— | ————————————————- | | 无响应:网络错误或超时 | 可能已扣款,也可能未扣款 | 是 | | 2xx | 已完成;若带 Idempotent-Replayed: true,这是首次尝试的结果 | 否,显示收据 | | 409,代码为 IDEMPOTENCY_KEY_IN_USE | 首次尝试仍在运行 | 是,等待 Retry-After 指定的秒数 | | 5xx | 尝试失败,服务器释放了键 | 是,退避后重试 | | 429 | 在处理器运行前被限流 | 是,等待 Retry-After | | 422 IDEMPOTENCY_KEY_REUSED | 客户端错误:把键用于不同请求 | 否 | | 其他 4xx:400、402(拒付)、404、409(已付款) | 最终结果;同一个键的重试仍得到相同答案 | 否;显示错误,新的付款尝试使用新键 | 应检查 code,不能只看状态码。409 也可能表示订单已付款,这种结果是最终结果。
把键存入 PostgreSQL
内存存储不足以用于生产。进程停止、部署或更换容器都会丢失记录,之后到达的重试可能再次付款。多实例场景中,重试落到另一实例也找不到记录,两个实例甚至可能同时获取锁。因此需要能原子执行“不存在才创建”的共享存储。 @nestjs/idempotency 不附带这样的实现,而是定义 IdempotencyStore 接口,由应用 provider 使用现有数据库或 Redis 客户端实现。接口有四个方法,前两个参数都是记录键和所有者,时长以整数毫秒传入:
acquire(key, owner, fingerprint, lockTtl):创建锁,或返回现有的执行中记录/包含响应的已完成记录。必须原子执行:任意多个请求竞争空闲键时只能有一个获胜。complete(key, owner, response, ttl):仅在owner仍持有锁时,将其转为已完成记录。release(key, owner):仅在owner仍持有锁时删除锁,让重试可以再次运行处理器。extend(key, owner, lockTtl):仅在owner仍持有锁时,将过期时间重置为从现在起lockTtl。处理器执行和结果保存期间,拦截器每隔lockTtl / 3调用一次。owner是每次尝试的随机 ID,也称 fencing token。如果实例崩溃等原因导致停止续租,重试在lockTtl后接管,旧尝试的extend()、complete()、release()就应无效,不能覆盖新尝试的记录。因此所有者检查必须与写入处于同一个原子操作;先单独读取检查再写入,会给接管留下竞态窗口。 本节先展示教程采用的 Drizzle 实现,再展示同一存储的 TypeORM 实现。根据应用 ORM 二选一。
先在 Drizzle schema 中为 PostgreSQL 记录建立表:
// 文件:database/schema
import type { IdempotencyStoredPayload } from '@nestjs/idempotency';
import { bigint, index, json, pgTable, text } from 'drizzle-orm/pg-core';
/** One row per idempotency record: an in-flight lock, or a completed response. */
export const idempotencyKeys = pgTable(
'idempotency_keys',
{
// SHA-256 of the record key: keys can be longer than an index entry allows.
keyHash: text('key_hash').primaryKey(),
key: text('key').notNull(),
fingerprint: text('fingerprint').notNull(),
// The attempt holding the lock; null once the record is completed.
owner: text('owner'),
// null while in flight. json, not jsonb: jsonb rejects "\u0000" in strings.
response: json('response').$type<IdempotencyStoredPayload>(),
// Epoch milliseconds.
expiresAt: bigint('expires_at', { mode: 'number' }).notNull(),
},
(table) => [index('idempotency_keys_expires_at_idx').on(table.expiresAt)],
);
主键使用记录键的 SHA-256 哈希,因为记录键包含 scope、客户端键,以及 GraphQL 字段路径,长度可能超过 PostgreSQL 索引项限制。响应使用 json,而不是 jsonb:后者拒绝字符串中的 \x00,但响应体可能包含它。用 drizzle-kit 生成迁移,再照常执行 npx drizzle-kit migrate:
$ npx drizzle-kit generate --name=idempotency
// 文件:drizzle/0000_idempotency.sql
CREATE TABLE "idempotency_keys" (
"key_hash" text PRIMARY KEY NOT NULL,
"key" text NOT NULL,
"fingerprint" text NOT NULL,
"owner" text,
"response" json,
"expires_at" bigint NOT NULL
);
--> statement-breakpoint
CREATE INDEX "idempotency_keys_expires_at_idx" ON "idempotency_keys" USING btree ("expires_at");
存储的各个方法使用 Drizzle 查询构建器实现,通过 @nestjs/drizzle 的 @InjectDrizzle() 注入数据库。构造函数向 IdempotencyModule 导出的注册表 IdempotencyStorage 注册存储:
// 文件:idempotency/drizzle-idempotency.store
import { Injectable } from '@nestjs/common';
import { InjectDrizzle } from '@nestjs/drizzle';
import {
IdempotencyStorage,
type IdempotencyAcquireResult,
type IdempotencyStore,
type IdempotencyStoredPayload,
} from '@nestjs/idempotency';
import { and, eq, gt, inArray, lte } from 'drizzle-orm';
import type { PgDatabase, PgQueryResultHKT } from 'drizzle-orm/pg-core';
import { createHash } from 'node:crypto';
import { idempotencyKeys } from '../database/schema.js';
/** A Drizzle database on PostgreSQL: node-postgres in the app, PGlite in its tests. */
export type Database = PgDatabase<PgQueryResultHKT>;
const { keyHash, owner: ownerColumn, expiresAt } = idempotencyKeys;
@Injectable()
export class DrizzleIdempotencyStore implements IdempotencyStore {
constructor(
@InjectDrizzle() private readonly db: Database,
storage: IdempotencyStorage,
) {
storage.registerSource(this);
}
async acquire(
key: string,
owner: string,
fingerprint: string,
lockTtl: number,
): Promise<IdempotencyAcquireResult> {
const hash = hashOf(key);
// A second round only happens when the row was released or expired
// between the two statements below.
for (let attempt = 0; attempt < 5; attempt++) {
const now = Date.now();
const lock = { fingerprint, owner, response: null, expiresAt: now + lockTtl };
// Insert the lock, or take over an expired row. Of many concurrent
// callers, PostgreSQL lets one write, and re-checks `setWhere` for the
// others against the winner's row: they get no row back.
const taken = await this.db
.insert(idempotencyKeys)
.values({ keyHash: hash, key, ...lock })
.onConflictDoUpdate({
target: keyHash,
set: lock,
setWhere: lte(expiresAt, now),
})
.returning({ keyHash });
if (taken.length) {
return { state: 'acquired' };
}
// Someone else holds the key: report their lock or record.
const [row] = await this.db
.select()
.from(idempotencyKeys)
.where(eq(keyHash, hash));
if (row && row.expiresAt > now) {
return row.response === null
? { state: 'in-flight', fingerprint: row.fingerprint }
: { state: 'completed', fingerprint: row.fingerprint, response: row.response };
}
}
throw new Error(`The idempotency record for "${key}" kept changing; try again`);
}
async complete(
key: string,
owner: string,
response: IdempotencyStoredPayload,
ttl: number,
): Promise<boolean> {
const now = Date.now();
const done = await this.db
.update(idempotencyKeys)
.set({ owner: null, response, expiresAt: now + ttl })
.where(this.ownedBy(key, owner, now))
.returning({ keyHash });
return done.length === 1;
}
async release(key: string, owner: string): Promise<boolean> {
const released = await this.db
.delete(idempotencyKeys)
.where(this.ownedBy(key, owner, Date.now()))
.returning({ keyHash });
return released.length === 1;
}
async extend(key: string, owner: string, lockTtl: number): Promise<boolean> {
const now = Date.now();
const extended = await this.db
.update(idempotencyKeys)
.set({ expiresAt: now + lockTtl })
.where(this.ownedBy(key, owner, now))
.returning({ keyHash });
return extended.length === 1;
}
/**
* Deletes up to `limit` expired rows, for a scheduled job. Expired rows are
* already ignored, and reused by the next request with the same key.
*/
async prune(limit = 1000): Promise<number> {
const now = Date.now();
const expired = this.db
.select({ keyHash })
.from(idempotencyKeys)
.where(lte(expiresAt, now))
.limit(limit);
const deleted = await this.db
.delete(idempotencyKeys)
// Checks the expiry again: a row taken over since the subquery read it stays.
.where(and(inArray(keyHash, expired), lte(expiresAt, now)))
.returning({ keyHash });
return deleted.length;
}
/** The caller's lock, if it still holds it and it hasn't expired. */
private ownedBy(key: string, owner: string, now: number) {
return and(eq(keyHash, hashOf(key)), eq(ownerColumn, owner), gt(expiresAt, now));
}
}
function hashOf(key: string) {
return createHash('sha256').update(key).digest('hex');
}
数据库类型为 PgDatabase,这是 Drizzle 的 node-postgres 和 PGlite 驱动共有的类型,便于后面的测试使用 PGlite。每个方法是一条 PostgreSQL 原子执行的语句,不需要事务:
acquire()用一条INSERT ... ON CONFLICT DO UPDATE ... WHERE expires_at <= now插入锁或接管过期记录。多个调用者竞争同一键时,PostgreSQL 只让一个写入,并针对获胜者的记录重新检查其他调用的WHERE;其他调用得不到返回行,随后读取获胜者的锁。如果记录期间被释放或过期消失,acquire()再循环尝试。complete()、extend()、release()在写入语句的WHERE中检查所有者和未过期条件。它们统计returning()返回行数,各 Drizzle 驱动行为一致。- 过期行立即被忽略,同一键的下一次请求会复用该行。其他过期行可通过定时任务调用
prune()清理,例如使用@nestjs/schedule的@Interval()。删除语句本身再次检查过期条件,避免删除期间被重试接管的记录。 过期时间来自各实例的Date.now(),因此需要通过 NTP 等机制同步时钟。每个幂等请求占用连接池执行一至两条语句,处理器运行时还会每隔lockTtl / 3执行一次UPDATE。
存储是普通 provider。用 DrizzleModule 注册数据库,参见 Drizzle,并把存储加入 AppModule 的 providers:
// 文件:app.module
import { Module } from '@nestjs/common';
import { DrizzleModule } from '@nestjs/drizzle';
import { IdempotencyModule } from '@nestjs/idempotency';
import { drizzle } from 'drizzle-orm/node-postgres';
import { DrizzleIdempotencyStore } from './idempotency/drizzle-idempotency.store.js';
import type { User } from './orders/order.js';
import { OrdersModule } from './orders/orders.module.js';
@Module({
imports: [
DrizzleModule.forRootAsync({
useFactory: () => ({ drizzle, connection: process.env.DATABASE_URL! }),
}),
IdempotencyModule.forRoot({
scope: (req: { user?: User }) => req.user?.id,
}),
OrdersModule,
],
// Registers itself with IdempotencyStorage when Nest creates it.
providers: [DrizzleIdempotencyStore],
})
export class AppModule {}
提示:示例从
process.env读取DATABASE_URL和启用加密后的IDEMPOTENCY_KEYS。实际应用可通过@nestjs/config加载环境,使用校验 schema 在变量缺失时阻止启动,并向这些工厂注入ConfigService。
启动时模块会记录正在使用的存储:IdempotencyStorage: DrizzleIdempotencyStore,随后停止接受注册:
- 应用只有一个存储。第二次调用
registerSource()会抛错,并列出两个类名。 - 应在单例 provider 的构造函数中注册。
IdempotencyModule初始化后注册表锁定;生命周期钩子、按请求创建的 provider 或延迟加载模块中的后续注册会抛错,而不是被静默忽略。 - 未注册存储时,模块使用内存。生产环境
NODE_ENV=production会拒绝这种方式,除非设置allowInMemoryStorage: true,明确接受每次重启或部署都丢失记录。 所有存储都应运行同一契约测试套件。@nestjs/idempotency/testing将其导出为适配任意测试运行器的普通用例,覆盖各方法、所有者检查、过期,以及启用concurrent后先读后写实现容易出错的竞态。原示例在 PGlite 和通过pg连接池访问的 PostgreSQL 上运行;后者让 32 个调用者通过最多 16 个连接竞争每个键:
// 文件:test/step-3.drizzle-store.spec
import { PGlite } from '@electric-sql/pglite';
import { IdempotencyStorage } from '@nestjs/idempotency';
import { idempotencyStoreContract } from '@nestjs/idempotency/testing';
import { drizzle as drizzlePglite } from 'drizzle-orm/pglite';
import { migrate as migratePglite } from 'drizzle-orm/pglite/migrator';
import { DrizzleIdempotencyStore } from '../src/idempotency/drizzle-idempotency.store.js';
const migrationsFolder = new URL('../drizzle', import.meta.url).pathname;
// ...
describe('DrizzleIdempotencyStore on PGlite', () => {
let pglite: PGlite;
let store: DrizzleIdempotencyStore;
beforeAll(async () => {
pglite = new PGlite();
const db = drizzlePglite(pglite);
await migratePglite(db, { migrationsFolder }); // npx drizzle-kit migrate
store = new DrizzleIdempotencyStore(db, new IdempotencyStorage());
});
afterAll(() => pglite.close());
describe('the IdempotencyStore contract', () => {
// The store reads Date.now(): fake Date alone, so the driver keeps its timers.
beforeEach(() => vi.useFakeTimers({ toFake: ['Date'] }));
afterEach(() => vi.useRealTimers());
// One store and one table for every case: each case uses keys of its own.
const cases = idempotencyStoreContract(() => store, {
advanceTime: (ms) => vi.setSystemTime(Date.now() + ms),
concurrent: true,
});
for (const c of cases) {
it(c.name, c.run);
}
});
// ...
});
使用 TypeORM
下面是使用相同表的 TypeORM 实现。应以它替换前面的 Drizzle 实现,两者不要同时注册。用实体表示表,每个列明确声明类型,让不依赖装饰器元数据加载实体的 TypeORM CLI 也能看到与应用相同的 schema:
// 文件:typeorm/idempotency-key.entity
import { Column, Entity, Index, PrimaryColumn } from 'typeorm';
// The idempotency_keys table of the Drizzle schema in src/database/schema.ts, as an entity.
// Every column states its type, so the entity loads the same with or without emitted
// decorator metadata (the TypeORM CLI runs it through tsx, which emits none).
/** Reads a `bigint` column as a number: TypeORM returns it as a string. */
const epochMs = {
to: (value: number) => value,
from: (value: string) => Number(value),
};
/** One row per idempotency record: an in-flight lock, or a completed response. */
@Entity('idempotency_keys')
@Index('idempotency_keys_expires_at_idx', ['expiresAt'])
export class IdempotencyKeyEntity {
// SHA-256 of the record key: keys can be longer than an index entry allows.
@PrimaryColumn({ type: 'text', name: 'key_hash' })
keyHash: string;
@Column({ type: 'text' })
key: string;
@Column({ type: 'text' })
fingerprint: string;
// The attempt holding the lock; null once the record is completed.
@Column({ type: 'text', nullable: true })
owner: string | null;
// null while in flight. json, not jsonb: jsonb rejects "\u0000" in strings. An
// IdempotencyStoredPayload, typed `object`: its `unknown` body doesn't fit TypeORM's insert types.
@Column({ type: 'json', nullable: true })
response: object | null;
// Epoch milliseconds.
@Column({ type: 'bigint', name: 'expires_at', transformer: epochMs })
expiresAt: number;
}
TypeORM 把 bigint 返回为字符串,所以使用 transformer 将 expiresAt 转回数字。响应列类型是 object,而非 IdempotencyStoredPayload,因为后者的 unknown 响应体不符合 TypeORM 插入类型;读取时再转换。TypeORM CLI 与 TypeOrmModule 在 data source 文件中共享配置:
// 文件:typeorm/data-source
import { DataSource, type DataSourceOptions } from 'typeorm';
import { IdempotencyKeyEntity } from './idempotency-key.entity.js';
import { Idempotency1790321611502 } from './migrations/1790321611502-Idempotency.js';
/** What the application's TypeOrmModule and the TypeORM CLI share. */
export const dataSourceOptions = {
type: 'postgres',
url: process.env.DATABASE_URL,
entities: [IdempotencyKeyEntity],
migrations: [Idempotency1790321611502],
} satisfies DataSourceOptions;
// The TypeORM CLI's data source: `migration:generate` compares the entities with this database.
export default new DataSource(dataSourceOptions);
从实体生成迁移、应用迁移,再检查实体与数据库是否一致。CLI 需要 TypeScript loader 来加载 data source,此处使用 tsx:
$ npx tsx ./node_modules/typeorm/cli.js migration:generate src/typeorm/migrations/Idempotency -d src/typeorm/data-source.ts --pretty
Migration /store/src/typeorm/migrations/1790321611502-Idempotency.ts has been generated successfully.
$ npx tsx ./node_modules/typeorm/cli.js migration:run -d src/typeorm/data-source.ts
$ npx tsx ./node_modules/typeorm/cli.js migration:generate src/typeorm/migrations/Check -d src/typeorm/data-source.ts --check
No changes in database schema were found
生成的 1790321611502-Idempotency.ts 包含与 Drizzle 迁移相同的 CREATE TABLE、CREATE INDEX 语句,并在 data source 的 migrations 中列出。存储遵循相同规则:
// 文件:typeorm/typeorm-idempotency.store
import { Injectable } from '@nestjs/common';
import {
IdempotencyStorage,
type IdempotencyAcquireResult,
type IdempotencyStore,
type IdempotencyStoredPayload,
} from '@nestjs/idempotency';
import { createHash } from 'node:crypto';
import { DataSource, MoreThan } from 'typeorm';
import { IdempotencyKeyEntity } from './idempotency-key.entity.js';
/**
* Keeps idempotency records in PostgreSQL, through TypeORM: the Drizzle
* store, on the same table. Each method is a single statement, so the
* guarantees hold with any number of API instances.
*/
@Injectable()
export class TypeOrmIdempotencyStore implements IdempotencyStore {
constructor(
private readonly dataSource: DataSource,
storage: IdempotencyStorage,
) {
storage.registerSource(this);
}
async acquire(
key: string,
owner: string,
fingerprint: string,
lockTtl: number,
): Promise<IdempotencyAcquireResult> {
const hash = hashOf(key);
// A second round only happens when the row was released or expired
// between the two statements below.
for (let attempt = 0; attempt < 5; attempt++) {
const now = Date.now();
// Insert the lock, or take over an expired row. Of many concurrent
// callers, PostgreSQL lets one write, and re-checks the overwrite
// condition for the others against the winner's row: they get no row back.
const { raw } = await this.dataSource
.createQueryBuilder()
.insert()
.into(IdempotencyKeyEntity)
.values({ keyHash: hash, key, fingerprint, owner, response: null, expiresAt: now + lockTtl })
.orUpdate(['fingerprint', 'owner', 'response', 'expires_at'], ['key_hash'], {
overwriteCondition: { where: 'idempotency_keys.expires_at <= :now', parameters: { now } },
})
.returning(['keyHash'])
.updateEntity(false)
.execute();
if ((raw as unknown[]).length === 1) {
return { state: 'acquired' };
}
// Someone else holds the key: report their lock or record.
const row = await this.dataSource.manager.findOneBy(IdempotencyKeyEntity, { keyHash: hash });
if (row && row.expiresAt > now) {
return row.response === null
? { state: 'in-flight', fingerprint: row.fingerprint }
: { state: 'completed', fingerprint: row.fingerprint, response: row.response as IdempotencyStoredPayload };
}
}
throw new Error(`The idempotency record for "${key}" kept changing; try again`);
}
async complete(
key: string,
owner: string,
response: IdempotencyStoredPayload,
ttl: number,
): Promise<boolean> {
const now = Date.now();
const { affected } = await this.dataSource.manager.update(IdempotencyKeyEntity, ownedBy(key, owner, now), {
owner: null,
response,
expiresAt: now + ttl,
});
return affected === 1;
}
async release(key: string, owner: string): Promise<boolean> {
const { affected } = await this.dataSource.manager.delete(IdempotencyKeyEntity, ownedBy(key, owner, Date.now()));
return affected === 1;
}
async extend(key: string, owner: string, lockTtl: number): Promise<boolean> {
const now = Date.now();
const { affected } = await this.dataSource.manager.update(IdempotencyKeyEntity, ownedBy(key, owner, now), {
expiresAt: now + lockTtl,
});
return affected === 1;
}
/**
* Deletes up to `limit` expired rows, for a scheduled job. Expired rows are
* already ignored, and reused by the next request with the same key.
*/
async prune(limit = 1000): Promise<number> {
const expired = this.dataSource.manager
.createQueryBuilder(IdempotencyKeyEntity, 'expired')
.select('expired.keyHash')
.where('expired.expiresAt <= :now')
.limit(limit);
const { affected } = await this.dataSource
.createQueryBuilder()
.delete()
.from(IdempotencyKeyEntity)
.where(`key_hash IN (${expired.getQuery()})`)
// Checks the expiry again: a row taken over since the subquery read it stays.
.andWhere('expires_at <= :now', { now: Date.now() })
.execute();
return affected ?? 0;
}
}
/** The caller's lock, if it still holds it and it hasn't expired. */
function ownedBy(key: string, owner: string, now: number) {
return { keyHash: hashOf(key), owner, expiresAt: MoreThan(now) };
}
function hashOf(key: string) {
return createHash('sha256').update(key).digest('hex');
}
acquire()仍通过插入查询构建器执行INSERT ... ON CONFLICT DO UPDATE ... WHERE expires_at <= now。orUpdate()列出接管时覆盖的列,overwriteCondition提供WHERE,returning()只向成功写入者返回行。complete()、extend()、release()使用实体管理器的update()、delete(),条件包含所有者和MoreThan(now),读取affected。不要使用先读记录的save(),否则接管可能插入读写之间。prune()使用 select 查询构建器生成选择过期行的子查询,将其 SQL 放入DELETE的WHERE,同时再次检查过期时间。 存储注入TypeOrmModule提供的DataSource。使用 data source 的配置注册该模块,替换DrizzleModule;存储也替换DrizzleIdempotencyStore:
// 文件:typeorm/app.module
import { Module } from '@nestjs/common';
import { IdempotencyModule } from '@nestjs/idempotency';
import { TypeOrmModule } from '@nestjs/typeorm';
import type { User } from '../orders/order.js';
import { OrdersModule } from '../orders/orders.module.js';
import { dataSourceOptions } from './data-source.js';
import { TypeOrmIdempotencyStore } from './typeorm-idempotency.store.js';
@Module({
imports: [
TypeOrmModule.forRootAsync({
// The entity and migration the CLI uses; migrations run on deploy (`migration:run`).
useFactory: () => ({ ...dataSourceOptions, url: process.env.DATABASE_URL }),
}),
IdempotencyModule.forRoot({
scope: (req: { user?: User }) => req.user?.id,
}),
OrdersModule,
],
// Registers itself with IdempotencyStorage when Nest creates it.
providers: [TypeOrmIdempotencyStore],
})
export class AppModule {}
契约测试套件无需修改。TypeORM 通过 pg 连接,因此测试需要 PostgreSQL,缺少服务器时跳过。另加一个检查实体与迁移是否一致的用例,相当于 migration:generate --check:
// 文件:test/typeorm-idempotency.store.spec
import { IdempotencyStorage } from '@nestjs/idempotency';
import { idempotencyStoreContract } from '@nestjs/idempotency/testing';
import { DataSource } from 'typeorm';
import { startPostgres } from '../../../test-support/postgres.js';
import { dataSourceOptions } from '../src/typeorm/data-source.js';
import { TypeOrmIdempotencyStore } from '../src/typeorm/typeorm-idempotency.store.js';
// ...
const { postgres, reason } = await startPostgres();
afterAll(() => postgres?.stop());
describe.skipIf(!postgres)(`TypeOrmIdempotencyStore on PostgreSQL${postgres ? '' : ` (skipped: ${reason})`}`, () => {
let dataSource: DataSource;
let url: string;
let store: TypeOrmIdempotencyStore;
beforeAll(async () => {
url = await postgres!.createDatabase('idempotency_typeorm_store');
dataSource = await new DataSource({ ...dataSourceOptions, url, poolSize: 16 }).initialize();
await dataSource.runMigrations(); // npx typeorm migration:run
store = new TypeOrmIdempotencyStore(dataSource, new IdempotencyStorage());
});
afterAll(() => dataSource?.destroy());
describe('the IdempotencyStore contract', () => {
// The store reads Date.now(): fake Date alone, so the driver keeps its timers.
beforeEach(() => vi.useFakeTimers({ toFake: ['Date'] }));
afterEach(() => vi.useRealTimers());
// One store and one table for every case: each case uses keys of its own.
// 32 callers per race, on up to 16 connections.
const cases = idempotencyStoreContract(() => store, {
advanceTime: (ms) => vi.setSystemTime(Date.now() + ms),
concurrent: { callers: 32 },
});
for (const c of cases) {
it(c.name, c.run);
}
});
it('has a migration that matches the entity (what `migration:generate --check` checks)', async () => {
const { upQueries } = await dataSource.driver.createSchemaBuilder().log();
expect(upQueries.map((query) => query.query)).toEqual([]);
});
// ...
});
startPostgres() 是测试辅助函数,用于在 PostgreSQL 服务器上创建临时数据库,没有服务器时跳过套件。原示例还在同一数据库上运行两个 TypeORM AppModule 实例,一个扣款时另一个返回 409。
提示:若使用 Prisma,教程后的“存储契约”说明了如何表达这些语句;也应使用相同测试套件验证。
改用 Redis 存储键
如果 API 已使用 Redis,可以把记录放在那里。记录是 Redis 可自动过期的小哈希,无需清理任务;客户端每秒重试同一付款这样的热点键,只需访问内存,避免主数据库行锁。代价是额外运行、保护和备份一个存储,并必须保证 Redis 不会提前驱逐记录,详见生产检查部分。 每个方法使用一个 Lua 脚本,由 Redis 原子执行:
// 文件:idempotency/redis-idempotency.store
import { Inject, Injectable } from '@nestjs/common';
import {
IdempotencyStorage,
type IdempotencyAcquireResult,
type IdempotencyStore,
type IdempotencyStoredPayload,
} from '@nestjs/idempotency';
import { REDIS } from '../redis/redis.module.js';
/** The one ioredis method this store needs. */
export interface RedisClient {
eval(script: string, numKeys: number, ...args: (string | number)[]): Promise<unknown>;
}
// Creates the lock if no record exists. Otherwise returns the record, in the
// same round trip.
const ACQUIRE = `
if redis.call('EXISTS', KEYS[1]) == 0 then
redis.call('HSET', KEYS[1], 'state', 'in-flight', 'fp', ARGV[1], 'owner', ARGV[2])
redis.call('PEXPIRE', KEYS[1], ARGV[3])
return {1}
end
local r = redis.call('HMGET', KEYS[1], 'state', 'fp', 'resp')
return {0, r[1], r[2], r[3]}
`;
// Turns our lock into a completed record, unless another attempt took it over.
const COMPLETE = `
if redis.call('HGET', KEYS[1], 'owner') ~= ARGV[1] then return 0 end
redis.call('HSET', KEYS[1], 'state', 'completed', 'resp', ARGV[2])
redis.call('HDEL', KEYS[1], 'owner')
redis.call('PEXPIRE', KEYS[1], ARGV[3])
return 1
`;
// Deletes the lock, but only if it is still ours.
const RELEASE = `
if redis.call('HGET', KEYS[1], 'owner') == ARGV[1] then
return redis.call('DEL', KEYS[1])
end
return 0
`;
// Renews our lock while the handler runs, but only if it is still ours.
const EXTEND = `
if redis.call('HGET', KEYS[1], 'owner') ~= ARGV[1] then return 0 end
redis.call('PEXPIRE', KEYS[1], ARGV[2])
return 1
`;
type AcquireReply =
| [1]
| [0, 'in-flight' | 'completed', string, string | null];
@Injectable()
export class RedisIdempotencyStore implements IdempotencyStore {
constructor(
@Inject(REDIS) private readonly redis: RedisClient,
storage: IdempotencyStorage,
) {
storage.registerSource(this);
}
async acquire(
key: string,
owner: string,
fingerprint: string,
lockTtl: number,
): Promise<IdempotencyAcquireResult> {
const reply = (await this.redis.eval(
ACQUIRE, 1, this.key(key), fingerprint, owner, lockTtl,
)) as AcquireReply;
if (reply[0] === 1) {
return { state: 'acquired' };
}
const [, state, storedFingerprint, response] = reply;
return state === 'completed'
? { state, fingerprint: storedFingerprint, response: JSON.parse(response!) }
: { state, fingerprint: storedFingerprint };
}
async complete(
key: string,
owner: string,
response: IdempotencyStoredPayload,
ttl: number,
): Promise<boolean> {
const done = await this.redis.eval(
COMPLETE, 1, this.key(key), owner, JSON.stringify(response), ttl,
);
return done === 1;
}
async release(key: string, owner: string): Promise<boolean> {
const released = await this.redis.eval(RELEASE, 1, this.key(key), owner);
return released === 1;
}
async extend(key: string, owner: string, lockTtl: number): Promise<boolean> {
const extended = await this.redis.eval(EXTEND, 1, this.key(key), owner, lockTtl);
return extended === 1;
}
private key(key: string) {
return `idem:${key}`;
}
}
每条记录是键 idem:<scope>:<key> 下的哈希,包含 state、fp(指纹)、owner、resp(JSON 响应)。执行中的锁在最后续租后的 lockTtl 到期,已完成记录在 ttl 后到期,因此无需清理任务。每个脚本只访问一个键,也适用于 Redis Cluster。Redis 原子运行脚本,但后续命令失败时,不会撤销脚本之前的写入。 例如 PEXPIRE 拒绝非整数毫秒。如果它在 HSET 之后失败,会留下永不过期的锁,因此拦截器向存储传递整数毫秒。 Nest 注入 Redis 客户端,构造函数负责注册存储,与 PostgreSQL 实现相同。应用只有一个存储,所以在 AppModule 中用 RedisIdempotencyStore 替换前一个:
// 文件:app.module
import { Module } from '@nestjs/common';
import { IdempotencyModule } from '@nestjs/idempotency';
import { RedisIdempotencyStore } from './idempotency/redis-idempotency.store.js';
import type { User } from './orders/order.js';
import { OrdersModule } from './orders/orders.module.js';
import { RedisModule } from './redis/redis.module.js';
@Module({
imports: [
RedisModule,
IdempotencyModule.forRoot({
scope: (req: { user?: User }) => req.user?.id,
}),
OrdersModule,
],
providers: [RedisIdempotencyStore],
})
export class AppModule {}
RedisModule 表示应用自己的全局 Redis 模块,以 REDIS token 提供 ioredis 客户端,例如 useFactory: () => new Redis(process.env.REDIS_URL!);在包完成执行中请求后,由模块的 onApplicationShutdown() 调用 quit() 关闭。RedisClient 只要求 eval(),ioredis Redis 实例可直接满足。 启动日志变为 IdempotencyStorage: RedisIdempotencyStore,注册规则仍是只有一个存储,并在构造函数中注册。 原示例项目也为此存储运行契约套件:
// 文件:test/step-4.redis-store.spec
import { idempotencyStoreContract } from '@nestjs/idempotency/testing';
// ...
describe('RedisIdempotencyStore, the IdempotencyStore contract (on FakeRedis)', () => {
let redis: FakeRedis;
const cases = idempotencyStoreContract(
() => {
redis = new FakeRedis();
return new RedisIdempotencyStore(redis, new IdempotencyStorage());
},
// Redis expires keys on its own clock: the fake's.
{ advanceTime: (ms) => redis.advance(ms), concurrent: true },
);
for (const c of cases) {
it(c.name, c.run);
}
});
提示:原教程示例项目没有 Redis 服务器,测试使用进程内替身,以 JavaScript 实现相同四个脚本。它覆盖存储逻辑及响应读取,但不覆盖 Lua 本身。实际依赖该存储前,应对真实 Redis(例如容器中的实例)运行契约测试。不传
advanceTime时,套件等待真实过期,约需十秒。 后续模块代码仍展示DrizzleIdempotencyStore。使用 Redis 时以RedisIdempotencyStore替换即可;加密、配送服务和调优参数相同。启用加密后,密封数据存于resp字段,而不是数据库的response列。
加密已存储的收据
保存的响应包含完整收据:支付服务商的收据 ID、金额和银行卡末四位。若明文存入 idempotency_keys,数据库的每份备份和转储都会包含这些数据,因此应加密静态记录。 加密键是秘密,应在应用启动时读取,而不是导入 app.module.ts 时读取,因此改用 forRootAsync()。存储仍是自行注册的 AppModule provider,工厂返回配置:
// 文件:app.module
IdempotencyModule.forRootAsync({
useFactory: () => ({
scope: (req: { user?: User }) => req.user?.id,
encryption: {
// Newest first: seal with the first key, open with any of them.
keys: process.env.IDEMPOTENCY_KEYS!.split(','),
},
}),
}),
使用 openssl rand -base64 32 生成密钥,在秘密管理器中保存为 IDEMPOTENCY_KEYS。启用加密后,response 列保存密封格式 v1.<keyId>.<iv>.<ciphertext>.<tag>,而不是明文收据:
- 记录使用 AES-256-GCM 加密,每条记录使用随机 IV。
- 指纹保持明文,因为存储需要比较它,而且它本身是哈希。
- 记录键与数据一起参与认证,把记录复制到另一用户的键下会导致解密失败,而不是在那里重放。
- 加密启用期间拒绝明文记录,防止持有表写权限的人植入响应。
- 字符串密钥至少 32 个字符,并通过 HKDF-SHA256 扩展。这适用于上面生成的随机秘密,但不是密码哈希,不应使用人自行选择的密码。密钥太短或首尾含空格(例如
IDEMPOTENCY_KEYS=new, old的第二个键以空格开头),会在启动时抛出指明问题键的错误。 应在端点保存第一条记录前开启加密。开关会双向改变记录格式:切换前保存的明文,或关闭加密后仍存在的密封记录,在过期前都无法读取,重放时返回下述 500。
要在不破坏重放的情况下轮换密钥:
- 生成新密钥,部署
IDEMPOTENCY_KEYS=old-key,new-key。实例仍用旧键加密,但已经能解密新键的数据。 - 所有实例采用该配置后,再部署
IDEMPOTENCY_KEYS=new-key,old-key。新记录使用新键,旧记录仍可读取;滚动部署期间尚未重启的实例也已能读取新记录。 - 等待
ttl(下文配置为 48 小时)过去,旧键加密的记录全部过期后,部署IDEMPOTENCY_KEYS=new-key。 过早删除密钥会让对应记录的重放以拒绝方式失败:客户端收到 500IDEMPOTENCY_RECORD_UNREADABLE,系统记录错误,处理器不会重新运行。该记录说明已经扣过款,再运行可能重复扣款。记录到期后该错误才消失。
通过 GraphQL 付款
网店网站通过 GraphQL mutation 付款。在 AppModule 注册 GraphQLModule:
// 文件:app.module
GraphQLModule.forRoot<ApolloDriverConfig>({
driver: ApolloDriver,
autoSchemaFile: true,
}),
在订单 resolver 的查询旁添加 mutation。GraphQL 的键既可来自与 REST 相同的 Idempotency-Key 请求头,也可来自参数。参数把键公开在 schema 中,便于代码生成器和审核者发现,也适合通常按连接而非按操作设置请求头的 GraphQL 客户端。 @Idempotent() 优先读取 mutation 的 idempotencyKey 参数,否则读取请求头,因此 resolver 只需声明参数:
// 文件:orders/orders.resolver
@Mutation(() => PaymentReceipt)
@Idempotent({ required: true })
payOrder(
@Args('orderId', { type: () => ID }) orderId: string,
@Args('input') input: PayOrderDto,
@Args('idempotencyKey') _idempotencyKey: string,
@CurrentUser() user: User,
) {
return this.ordersService.pay(orderId, user, input);
}
PayOrderDto 是 REST 使用的同一个类,增加 @InputType() 和 @Field();PaymentReceipt 是包含收据字段的 @ObjectType()。resolver 使用同一个 AuthGuard,从 GraphQL context 读取请求;模块的 scope 也收到同一个请求,因此键仍按用户隔离。参数是非空 String,缺少它会被 GraphQL 校验直接拒绝:
mutation PayOrder($orderId: ID!, $paymentMethod: String!, $key: String!) {
payOrder(orderId: $orderId, input: { paymentMethod: $paymentMethod }, idempotencyKey: $key) {
receiptId
amount
cardLast4
paidAt
}
}
相同变量的重试返回保存的 payOrder 结果,并保留 resolver 返回的类型,例如 Date 仍作为 Date 重放。记录键附加字段路径,例如 usr_alice:<key>:payOrder,因为一次操作可以包含多个 mutation。每个字段独立保存和重放,别名也是路径的一部分,因此重试必须使用相同文档。 拒绝以字段错误表示,而非 HTTP 状态码;响应仍为 200,因为同一操作中的其他字段可能成功。状态和重试延迟位于 extensions,以下省略 locations:
{
"errors": [
{
"message": "A request with this idempotency key is still being processed.",
"path": ["payOrder"],
"extensions": {
"code": "IDEMPOTENCY_KEY_IN_USE",
"httpStatus": 409,
"retryAfter": 2
}
}
],
"data": null
}
同理,GraphQL 重放不会设置 Idempotent-Replayed 响应头。
对 order.paid 事件去重
付款后 OrdersService 发出 order.paid,配送服务安排快递。实际系统通常至少投递一次:发布者超时后重试,或 RabbitMQ、Kafka、NATS JetStream 重新投递未确认消息,都可能让同一事件到达两次,造成顾客收到两份包裹。 去重需要同一事件的每次投递都使用相同 ID,因此应从数据派生,而不是每次发送时重新生成。订单只付款一次,所以 OrdersService 使用 order-paid: 加订单 ID。为事件处理器 加上 @Idempotent(),并让 keyFrom 指向这个 ID:
// 文件:shipping/shipping.controller
import { Controller } from '@nestjs/common';
import { EventPattern, Payload } from '@nestjs/microservices';
import { Idempotent } from '@nestjs/idempotency';
import { ShipmentsService } from './shipments.service.js';
export interface OrderPaidEvent {
eventId: string;
orderId: string;
}
@Controller()
export class ShippingController {
constructor(private readonly shipmentsService: ShipmentsService) {}
@EventPattern('order.paid')
@Idempotent({ required: true, keyFrom: { payload: 'eventId' } })
async onOrderPaid(@Payload() event: OrderPaidEvent) {
await this.shipmentsService.createFor(event.orderId);
}
}
微服务处理器默认先读取 payload 的 idempotencyKey,再读取由 header 指定名称的传输头。这里用 keyFrom 改读 eventId,并用 required: true 拒绝没有 ID 的事件。 配送服务是独立应用,有自己的数据库和 IdempotencyModule,注册相同的存储类:
// 文件:shipping/shipping.module
import { Module } from '@nestjs/common';
import { DrizzleModule } from '@nestjs/drizzle';
import { IdempotencyModule } from '@nestjs/idempotency';
import { drizzle } from 'drizzle-orm/node-postgres';
import { DrizzleIdempotencyStore } from '../idempotency/drizzle-idempotency.store.js';
import { ShipmentsService } from './shipments.service.js';
import { ShippingController } from './shipping.controller.js';
@Module({
imports: [
DrizzleModule.forRootAsync({
useFactory: () => ({ drizzle, connection: process.env.DATABASE_URL! }),
}),
IdempotencyModule.forRoot({
ttl: '7d', // brokers can redeliver long after the first delivery
}),
],
controllers: [ShippingController],
providers: [ShipmentsService, DrizzleIdempotencyStore],
})
export class ShippingModule {}
// 文件:shipping/main
import { NestFactory } from '@nestjs/core';
import { Transport, type MicroserviceOptions } from '@nestjs/microservices';
import { ShippingModule } from './shipping.module.js';
async function bootstrap() {
const app = await NestFactory.createMicroservice<MicroserviceOptions>(ShippingModule, {
transport: Transport.TCP,
options: { host: '0.0.0.0', port: 4001 },
});
await app.listen();
}
bootstrap();
与 API 相比,有几处不同:
-
不使用
scope,因为事件 ID 在整个系统中唯一。 -
ttl为 7 天,因为消息可能在首次处理很久之后才从死信队列返回。 -
记录归属于处理器,记录键以类名和方法名结尾,例如
order-paid%3Aord_1001:ShippingController.onOrderPaid。Nest 运行事件的所有已注册处理器,因此另一个发送物流单号邮件的order.paid处理器拥有自己的记录,也能处理每个事件一次。 重复事件按两种情况处理: -
首次投递已经完成:跳过重复事件。保存的结果被“重放”,但事件没有回复,因此不发生其他操作。
-
首次投递仍在运行:返回
IDEMPOTENCY_KEY_IN_USE,Nest 记录错误。RabbitMQ、Kafka 等需要确认消息的传输中,不应把该次投递确认成已处理,应让传输稍后重投,届时会跳过。 这样只预约一次快递。HTTP 付款重试本身也不会生成重复事件,因为重放请求不会进入OrdersService。
调整 TTL 和失败处理
默认值是起点。API 的最终配置如下:
// 文件:app.module
IdempotencyModule.forRootAsync({
useFactory: () => ({
scope: (req: { user?: User }) => req.user?.id,
encryption: {
// Newest first: seal with the first key, open with any of them.
keys: process.env.IDEMPOTENCY_KEYS!.split(','),
},
// Longer than the app keeps retrying a queued payment (24 hours).
ttl: '48h',
// A crashed instance blocks a key for at most 45 seconds. While a
// charge runs, its lock is renewed every 15 seconds.
lockTtl: '45s',
// How long a client waits after a 409 before it asks again.
retryAfter: '2s',
}),
}),
所有时长接受毫秒数,或 '45s'、'48h'、'7d' 等字符串。无法解析的值会阻止启动,并指出对应选项。
ttl(默认'24h')决定已完成结果可以重放多久,必须长于客户端最长重试时间。应用将离线付款排队最多一天,因此这里设置为48小时。TTL过后记录被遗忘,重试会再次运行处理器。在本例中,订单状态检查会得到409“已付款”而非再次扣款,但顾客无法取得原收据。lockTtl(默认'60s')决定执行处理器的实例崩溃后,执行中的锁还能存活多久。正常运行时每隔lockTtl / 3续租;这里是每 15 秒一次,因此扣款超过 45 秒也能保有锁,重试得到 409。停止续租后,锁过期,重试才接管。更短的值可更快释放崩溃后的键,但增加续租频率。 无法访问存储以续租时,拦截器记录丢失锁的警告。retryAfter(默认'1s')决定 409 的Retry-After,单位为整数秒。扣款通常需要一至两秒,因此这里要求客户端等待两秒。@Idempotent()接受required、ttl、lockTtl、retryAfter、keyFrom、scope、fingerprint、storeIf,覆盖该处理器对应的模块选项。fingerprint决定重试必须重复的内容,例如排除客户端时间戳后的请求体;方法和 URL 始终参与。装饰器也可用于控制器或 resolver 类,覆盖会修改状态的处理器;GET 路由和 GraphQL 查询默认跳过,除非自身也有@Idempotent()。 处理器级选项与类级选项合并,因此处理器只设置lockTtl时,类上的required: true仍有效。scope适用于所有处理器,但收到的内容不同:HTTP、GraphQL 为请求,微服务为消息 payload。同一模块服务多个上下文时,可以传对象,分别在http、graphql、rpc键下提供函数,各函数按对应上下文类型定义;没有提供函数的上下文没有 scope。
并非所有结果都保存。默认保存状态码低于 500 的结果: | 结果 | 默认行为 | 原因 | | ————————————————————————— | ———————————————————– | ———————————————————- | | 2xx | 保存并重放 | 付款成功 | | 4xx HttpException:校验400、拒付402、404、已付款409 | 保存,每次重放再次抛出 | 确定性结果:相同请求得到相同答案 | | 其他带4xx status 的错误,例如403 AuthorizationError | 作为相应状态的Nest异常保存并重放 | 所属包会把它转换为该状态 | | 5xx或未知错误,例如服务商503、数据库不可用 | 不保存,释放键,允许重试重新执行 | 多数故障发生在副作用之前 | | StreamableFile、流、未启用 passthrough 的 @Res() | 不保存,释放键并记录警告 | 无法重放 | 重放错误会重建为带有已保存状态和响应体的普通 HttpException,再次经过异常过滤器。若过滤器用 instanceof 检查自定义异常类,重放时看到的是基类。 storeIf 可以改变保存规则,接收状态码,以及处理器抛出的错误。例如 storeIf: () => true 也保存5xx,与 Stripe 的做法相同。默认规则有一个空隙:副作用之后的5xx。如果服务商扣款成功,随后 markPaid() 失败,键被释放,重试可能再次扣款。可以用以下三种方式弥补:
- 向支付服务商传入由用户ID和幂等键共同派生的键。多数支付服务商有自己的幂等键支持,可以对扣款去重。只用幂等键不够,因为两个用户可能发送相同键。
- 使用
storeIf保存该处理器的5xx,让应用在5xx后开始新的尝试。若只保存扣款后的失败,可抛出专门异常并检查:storeIf: (status, error) => status < 500 || error instanceof ChargeCapturedError。 - 使用事务发件箱,在同一个事务中写入付款和事件。
尝试运行
启动配送服务及 API,配置 DATABASE_URL(已执行迁移)、IDEMPOTENCY_KEYS 和支付服务商沙盒。沙盒约半秒返回。以下为原教程示例输出,仅保留相关响应头。
不携带键的付款被拒绝:
$ curl -i localhost:3000/orders/ord_1001/pay \
-H 'Authorization: Bearer alice-token' \
--json '{"paymentMethod":"pm_card_visa"}'
HTTP/1.1 400 Bad Request
{"statusCode":400,"error":"Bad Request","code":"IDEMPOTENCY_KEY_REQUIRED","message":"An idempotency key is required for this operation."}
首次携带键的请求会扣款:
$ KEY=$(uuidgen)
$ curl -i localhost:3000/orders/ord_1001/pay \
-H 'Authorization: Bearer alice-token' \
-H "Idempotency-Key: $KEY" \
--json '{"paymentMethod":"pm_card_visa"}'
HTTP/1.1 201 Created
Content-Type: application/json; charset=utf-8
{"orderId":"ord_1001","receiptId":"ch_0001","amount":4098,"cardLast4":"4242","paidAt":"2026-09-22T11:09:25.813Z"}
相同键的重试得到相同收据,包括相同的 paidAt,且不再调用支付服务商:
$ curl -i localhost:3000/orders/ord_1001/pay \
-H 'Authorization: Bearer alice-token' \
-H "Idempotency-Key: $KEY" \
--json '{"paymentMethod":"pm_card_visa"}'
HTTP/1.1 201 Created
Idempotent-Replayed: true
Content-Type: application/json; charset=utf-8
{"orderId":"ord_1001","receiptId":"ch_0001","amount":4098,"cardLast4":"4242","paidAt":"2026-09-22T11:09:25.813Z"}
相同键搭配不同银行卡属于客户端错误:
$ curl -i localhost:3000/orders/ord_1001/pay \
-H 'Authorization: Bearer alice-token' \
-H "Idempotency-Key: $KEY" \
--json '{"paymentMethod":"pm_card_mastercard"}'
HTTP/1.1 422 Unprocessable Entity
{"statusCode":422,"error":"Unprocessable Entity","code":"IDEMPOTENCY_KEY_REUSED","message":"This idempotency key was already used for a different request."}
最后,对另一个订单同时发送两个使用新键的请求。一个完成扣款,另一个在支付服务商尚未完成时得到409和 Retry-After: 2:
$ KEY=$(uuidgen)
$ for i in 1 2; do
curl -s -o /dev/null -w '%{http_code} %header{retry-after}\n' \
localhost:3000/orders/ord_1002/pay \
-H 'Authorization: Bearer alice-token' \
-H "Idempotency-Key: $KEY" \
--json '{"paymentMethod":"pm_card_visa"}' &
done; wait
409 2
201
提示:
--json需要 curl 7.82 或更新版本;-w中的%header需要 curl 7.84 或更新版本。
测试
端到端测试保留应用实际配置,但以普通 InMemoryIdempotencyStore 实例覆盖 DrizzleIdempotencyStore provider。普通实例不会自行注册,所以模块回退到内存存储;Nest 不会创建 DrizzleIdempotencyStore,也就没有数据库查询。node-postgres 连接池在首次查询时连接,因此 DrizzleModule 也不会打开连接。 加密键仍来自 IDEMPOTENCY_KEYS,应在测试环境中设置,例如在 vitest.config.mts 使用 test.env。 为支付服务商 mock 添加短暂延迟,让两个同时发送的请求重叠:
// 文件:test/pay-order.e2e-spec
import { ValidationPipe, type INestApplication } from '@nestjs/common';
import { InMemoryIdempotencyStore } from '@nestjs/idempotency';
import { Test } from '@nestjs/testing';
import { setTimeout } from 'node:timers/promises';
import { EMPTY } from 'rxjs';
import request from 'supertest';
import type { MockInstance } from 'vitest';
import { AppModule } from '../src/app.module.js';
import { DrizzleIdempotencyStore } from '../src/idempotency/drizzle-idempotency.store.js';
import { SHIPPING_SERVICE } from '../src/orders/orders.service.js';
import { PaymentProviderClient } from '../src/payments/payment-provider.client.js';
describe('POST /orders/:id/pay', () => {
let app: INestApplication;
let charge: MockInstance<PaymentProviderClient['charge']>;
beforeEach(async () => {
const moduleRef = await Test.createTestingModule({ imports: [AppModule] })
.overrideProvider(DrizzleIdempotencyStore)
.useValue(new InMemoryIdempotencyStore())
.overrideProvider(SHIPPING_SERVICE)
.useValue({ emit: () => EMPTY })
.compile();
app = moduleRef.createNestApplication();
app.useGlobalPipes(new ValidationPipe({ whitelist: true }));
await app.listen(0, '127.0.0.1');
charge = vi.spyOn(app.get(PaymentProviderClient), 'charge').mockImplementation(async (input) => {
await setTimeout(100); // slow enough for duplicates to overlap
return {
id: 'ch_test',
amount: input.amount,
cardLast4: '4242',
reference: input.reference,
createdAt: new Date().toISOString(),
};
});
});
afterEach(() => app.close());
const pay = (key: string) =>
request(app.getHttpServer())
.post('/orders/ord_1001/pay')
.set('Authorization', 'Bearer alice-token')
.set('Idempotency-Key', key)
.send({ paymentMethod: 'pm_card_visa' });
it('charges once when the app retries', async () => {
const first = await pay('key-1').expect(201);
const retry = await pay('key-1').expect(201).expect('Idempotent-Replayed', 'true');
expect(retry.body).toEqual(first.body);
expect(charge).toHaveBeenCalledTimes(1);
});
it('rejects a duplicate that arrives while the first is in flight', async () => {
const [a, b] = await Promise.all([pay('key-2'), pay('key-2')]);
expect([a.status, b.status].sort()).toEqual([201, 409]);
expect(charge).toHaveBeenCalledTimes(1);
});
});
初始化完成后,app.get(IdempotencyStorage).source 返回正在使用的存储,此处为内存存储。peek(key) 返回 usr_alice:key-1 等记录键下的原始记录(state、fingerprint、response),便于断言拒付结果已存储等情况。自定义存储应另用前述契约套件对真实服务器测试。 若端到端测试仍保留 Drizzle 存储,可用 overrideProvider(getDrizzleToken()).useValue(db) 替换 DrizzleModule 注册的数据库;getDrizzleToken 来自 @nestjs/drizzle,db 是已经执行应用迁移的 PGlite Drizzle 数据库。原示例以此测试加密、GraphQL、事件及调优设置。 测试设置参见端到端测试。
生产环境检查
- 即使单实例也注册共享数据库或 Redis 存储。内存记录会在重启和部署时丢失,之后的重试会再次执行处理器。启用
concurrent运行契约套件。生产环境未注册存储时应用拒绝启动,除非设置allowInMemoryStorage。Redis 设置maxmemory-policy noeviction。 提前驱逐记录可能导致重复扣款;每条幂等记录都有TTL,因此volatile-*策略同样会驱逐它们。 - 对所有转移资金、或重复创建会被顾客察觉的资源的端点,设置
required: true。 - 所有面向客户端的端点都设置
scope,避免一位用户的键重放另一位用户的响应。已登录用户调用未设置scope的处理器时,模块会记录警告。 - 让幂等拦截器处于其他全局拦截器外层。全局拦截器按注册顺序运行,
AppModule的APP_INTERCEPTORproviders 比其导入模块中的先注册,因此处于幂等拦截器外层并看到每次重放。模块会在启动时警告;若其中有ClassSerializerInterceptor,则拒绝启动,因为重放交给它的是普通对象,可能发送本应由@Exclude()隐藏的字段。 这些拦截器应在main.ts通过app.useGlobalInterceptors()注册,或放在IdempotencyModule之后导入的模块。 - 数据库存储通过迁移建表,按额外语句量配置连接池,同步实例时钟,并定期调用
prune()。 ttl应大于客户端重试窗口。因为处理期间续租,lockTtl只需覆盖崩溃后的等待时间。- 监控事件。可注入的
IdempotencyEvents通过events$报告每次重放、带code的拒绝和锁丢失;同样的事件发布到nestjs:idempotency:replayed、nestjs:idempotency:rejected、nestjs:idempotency:lock-lost诊断通道,供追踪工具使用。 对lock-lost报警:它表示结果未进入存储,原因可能是续租失败(存储不可达或事件循环阻塞超过lockTtl),或保存结果失败,重试可能再次运行处理器。 - 响应包含个人或付款数据时,在第一条记录保存前启用
encryption。将密钥保存在秘密管理器,按上述三步轮换,在ttl过去前不要移除旧键。 - 向支付服务商传递由用户ID和幂等键派生的键,覆盖扣款后失败的情形。
- 多服务共享Redis时,各服务使用独立前缀(例如实现中的
idem:)或数据库。记录键只在单服务内唯一,两个服务可能有同名处理器类,客户端也可能向两者发送同一个键。 - 保留服务层状态检查,最好使用条件数据库更新。幂等键去重的是同一次尝试的重试;两次点击创建两个键,就是两次尝试。
- 允许浏览器发送请求头。跨域客户端在 CORS
allowedHeaders中加入Idempotency-Key,在exposedHeaders中加入Idempotent-Replayed和Retry-After。 - 保持幂等响应较小,因为会保存整个响应体,包目前尚无大小限制。
存储契约
教程已给出 Drizzle、TypeORM、Redis 实现。这里概括所有 IdempotencyStore 必须满足的要求,供其他ORM实现参考;包 README 的“Implementing a store”详细说明各方法。 注册。 存储是普通单例provider,在构造函数中对注入的 IdempotencyStorage 调用 registerSource(this)。注册表立即检查结构;除非传入 replace: true,否则拒绝第二次注册;在 IdempotencyModule 初始化时锁定并记录存储名称。未注册时使用 InMemoryIdempotencyStore,实例之间不能共享,重启会丢失记录。 NODE_ENV=production 时这种情况会启动失败,除非设置 allowInMemoryStorage: true,接受每次重启或部署丢失记录。 方法。 四个方法都必须原子执行。先读再单独写会让两个请求同时执行处理器,产生本包要避免的重复扣款。方法不接收业务事务;幂等记录独立写入,不属于业务事务。 | 方法 | 行为 | 如何保证原子性 | | — | — | — | | acquire(key, owner, fingerprint, lockTtl) | 对空闲或过期键加锁,或返回现有记录 | 一条插入语句同时接管过期行,竞争者中恰好一个获胜,其他读取其记录 | | complete(key, owner, response, ttl) | 将所有者未过期的锁转为完成记录 | 比较并交换:写语句的WHERE同时检查所有者与过期时间 | | release(key, owner) | 删除所有者未过期的锁,允许重试 | 与complete相同的比较并交换 | | extend(key, owner, lockTtl) | 续租所有者未过期的锁 | 与complete相同的比较并交换 | Prisma实现。 以Drizzle和TypeORM实现及相同表为参照:
acquire()使用一条INSERT ... ON CONFLICT (key_hash) DO UPDATE ... WHERE expires_at <= now RETURNING。Prismaupsert()的更新不支持此条件,应使用$queryRaw。未返回行时,与其他实现一样读取现有记录。complete()、release()、extend()把所有者和过期时间放入WHERE并统计影响行数:updateMany()或deleteMany()的where包含keyHash、owner,以及expiresAt的gt: now条件,检查count === 1。不要先findUnique()再update()。- 响应列 在任何ORM中都使用
json,不是jsonb;后者拒绝字符串中的\x00,而响应体可能包含它。 测试。 使用@nestjs/idempotency/testing套件。idempotencyStoreContract()返回适配任意运行器的用例;concurrent: true让16个调用者竞争每个键,应在带连接池的PostgreSQL上运行,确保真实重叠。advanceTime推进存储时钟用于过期测试;省略则等待真实时间。
参考
模块选项
IdempotencyModule.forRoot() 接受下列选项。forRootAsync() 的 isGlobal、imports 放在顶层,与 useFactory、useClass、useExisting 并列,工厂返回其他选项。除 header、replayHeaders、encryption、allowInMemoryStorage、isGlobal 外,其余也可在类或处理器的 @Idempotent() 中设置,按字段覆盖模块配置。 | 选项 | 类型 | 默认值 | 说明 | | — | — | — | — | | required | boolean | false | 缺少键时以 IDEMPOTENCY_KEY_REQUIRED 拒绝 | | scope | 函数、对象或 false | 无 | 键命名空间,如用户ID;对象按http/graphql/rpc分别提供函数,false表示明确共享命名空间 | | ttl | Duration | '24h' | 已完成结果可重放的时长 | | lockTtl | Duration | '60s' | 停止续租后锁的存活时间;每lockTtl/3续租 | | retryAfter | Duration | '1s' | IDEMPOTENCY_KEY_IN_USE 的等待时间,向上取整到秒 | | keyFrom | 来源,或按上下文分别指定来源 | 见下表 | header、arg(GraphQL)、payload(微服务),或接受ExecutionContext的函数 | | fingerprint | (payload, context) => unknown | 整个payload | 重试必须重复的内容;方法及URL、GraphQL字段或消息模式,以及scope始终参与 | | storeIf | (status, error) => boolean | status < 500 | 哪些结果保存并重放;其他结果释放键 | | header | string | 'Idempotency-Key' | HTTP请求头,也供GraphQL读取,并作为微服务传输头名称 | | replayHeaders | string[] | 无 | 除Location、Content-Type、Content-Language、Content-Location、ETag、Last-Modified外额外重放的响应头;Set-Cookie及平台为每个响应写入的头会在启动时被拒绝 | | encryption | 包含keys的对象 | 关闭 | AES-256-GCM加密;keys列出32字节Buffer或至少32字符的字符串,新键在前 | | allowInMemoryStorage | boolean | false | 允许生产应用在未注册存储时启动 | | isGlobal | boolean | true | 全局注册模块 | 所有时长支持毫秒数,或 '45s'、'48h'、'7d' 等字符串。
键的默认来源:
| 上下文 | 默认来源 |
|---|---|
| HTTP | Idempotency-Key 请求头,由header选项指定 |
| GraphQL | mutation的idempotencyKey参数,其次是请求头 |
| 微服务 | payload的idempotencyKey属性,其次是header指定的传输头 |
键为1至255个可打印ASCII字符。payload中的数字按其数字字符串处理。
事件
IdempotencyEvents.events$ 发出每个事件,也将其发布到诊断通道。每个payload都有 type、context('http'、'graphql'或'rpc')、handler(如OrdersController.pay)、key,以及存在时的scope。 | 类型 | 通道 | Payload类型 | 其他字段 | 时机 | | — | — | — | — | — | | replayed | nestjs:idempotency:replayed | IdempotencyReplayedEvent | status | 重试获得已保存结果,或跳过重复事件 | | rejected | nestjs:idempotency:rejected | IdempotencyRejectedEvent | code、status | 处理器执行前被拒绝;KEY_REQUIRED和KEY_INVALID没有key | | lock-lost | nestjs:idempotency:lock-lost | IdempotencyLockLostEvent | phase:extend、complete或release | 结果未进入存储,重试可能再次运行处理器 | IdempotencyEvent 是这三种payload类型的联合类型。
错误
拒绝携带以下错误代码,类型为 IdempotencyErrorCode。HTTP放在响应体;GraphQL放在字段错误的 extensions,包含 httpStatus,409另有 retryAfter;微服务的 RpcException payload包含 status、code、statusCode、message,409另有 retryAfter。 | 代码 | 状态 | 原始消息 | 原因 | | — | — | — | — | | IDEMPOTENCY_KEY_REQUIRED | 400 | An idempotency key is required for this operation. | required为true但没有提供键 | | IDEMPOTENCY_KEY_INVALID | 400 | The idempotency key must be 1 to 255 printable ASCII characters. | 键过长、含非允许字符,或不是字符串/数字 | | IDEMPOTENCY_KEY_IN_USE | 409,带Retry-After | A request with this idempotency key is still being processed. | 首次请求仍运行 | | IDEMPOTENCY_KEY_REUSED | 422 | This idempotency key was already used for a different request. | 相同键搭配不同指纹 | | IDEMPOTENCY_RECORD_UNREADABLE | 500 | The stored result for this idempotency key could not be read. | 记录无法解密,例如其密钥已移除 |
来源:Idempotency keys,NestJS 文档,Kamil Myśliwiec 与贡献者。本文为中文翻译;保留完整实现、测试和原教程示例输出,将源码的文件名展示标记改为注释。
Copyright (c) 2017-present Kamil Myśliwiec http://kamilmysliwiec.com。采用 MIT 许可。
特此免费授予任何取得本软件及相关文档副本的人不受限制地处理本软件的权利,包括使用、复制、修改、合并、发布、分发、再许可和/或出售副本,并允许获得本软件的人这样做,但须在所有副本或实质部分中保留上述版权声明和本许可声明。
本软件按原样提供,不作任何明示或默示保证,包括但不限于适销性、特定用途适用性和不侵权。无论依据合同、侵权或其他理由,作者或版权持有人均不对因本软件、使用本软件或其他相关交易而产生的索赔、损害或其他责任负责。











暂无评论内容