创建订单的事务已经提交,客户端没有收到响应。客户端继续发送相同请求,服务端既要返回已有订单,也要防止同一个键被拿来购买另一件商品。这个问题可以缩到一个可以观察的边界:同一业务身份在数据库里竞争同一条唯一记录,获胜者把订单和可重放结果放进同一事务。
本轮实验使用独立 PostgreSQL 18.6 数据库,两个独立的 psql 连接。脚本先断言数据库名、用户和主版本,再创建只属于本次运行的表。提交、回滚、等待超时与清理竞争都在真实数据库中执行;响应丢失通过调用端舍弃结果后重试来模拟,没有注入网络断开,也没有连接支付服务。
为一笔订单确定稳定身份
假设员工为部门采购一件设备,金额一百元,客户端为这次确认动作生成键 K。用户点一次确认后产生的重试都沿用 K;用户修改金额再确认,则发起新的业务意图。若网络库每次重试都生成新键,服务端的唯一约束看见的是不同请求,无法替客户端判断它们是不是同一件事。
服务端把登录身份、租户和操作类型组合成可信作用域,作用域与 K 一起作为主键。实验把三个固定值拼成一个字符串,便于查看竞争过程。实际接口宜使用独立列,明确租户之间、调用方之间是否共享业务身份;客户端提交的用户编号不能代替鉴权结果。
请求指纹只包含影响业务结果的参数。金额按最小货币单位规范化,币种、商品、收货对象、接口语义版本都要进入指纹。追踪编号和发送时间不参与,字段顺序也不应改变指纹。服务端补默认值之后再计算,才能避免省略默认币种与显式填写默认币种被当成不同意图。
规范化本身值得保留测试样例。整数一百和字符串“一百”是否允许进入同一接口,应由输入校验先决定。数组的顺序是否影响含义,也取决于商品清单的定义。直接给原始 JSON 字节做摘要,虽然实现容易,却把序列化差异变成了业务冲突。
实验用固定的两个指纹字符串检查同键异参分支,没有实现摘要算法或 HTTP 参数解析。读者可以据此观察数据库协议,但不能把这项输出理解为已经测试了规范化安全性。线上还需要键长度上限、请求体大小限制,以及控制恶意调用方占满幂等表的配额。
RFC 9110 对幂等方法的定义关注重复请求的预期服务端效果。本文为创建订单的 POST 另订应用协议,并选择同键异参返回冲突;它没有把自定义请求头或某个状态码解释成 HTTP 自动提供的保证。
把领取执行权和业务提交放在一起
实验建四张表:幂等记录、订单、待发送事件、模拟下游接收表。幂等记录以作用域和键唯一,保存指纹、响应和过期标记;订单以作用域和业务编号唯一;事件与接收表也有稳定的业务身份。脚本给表名增加本次运行编号,重复执行不会清空其他实验的数据。
同步协议开始一笔事务,先执行冲突时不插入的 INSERT,并要求 RETURNING 返回键。返回一行表示本事务拿到了执行权,于是写订单、写事件、更新响应,最后提交。返回零行表示已有同键记录,当前请求进入读取和重放分支,不能继续创建订单。
BEGIN;
INSERT INTO idem.records_example(scope, key, hash, body, expired)
VALUES ('tenant1:user1:create', 'K', 'hash100', NULL, false)
ON CONFLICT (scope, key) DO NOTHING
RETURNING key;
-- 应用检查 RETURNING:仅新插入者写业务与响应。
COMMIT;
这段 SQL 展示分支入口,表名是阅读用名称,完整可执行 DDL 和双方会话由文末脚本提供。同步路径不单独提交一个处理中记录。业务失败时,订单、事件和幂等记录一起回滚;业务成功时,它们一起可见。只有发生了单独提交,才需要另外设计处理中记录的超时接管协议。
没有把业务变更和幂等结果放进同一事务的反例很直接:先创建订单并提交,再写幂等结果,进程在两次提交之间退出。重试找不到键,会再次创建订单。反过来先提交幂等占位,再执行业务,进程退出后又会留下没有业务结果的占位。单纯补一个状态字段解决不了这两个窗口。
事务里也不要等待用户确认或调用远端支付。占着数据库锁等待远端结果,会让同键请求排队,还会把外部超时与本地事务失败混在一起。本文的同步事务只处理同一数据库里的订单与发送意图,外部发送在提交之后完成。
PostgreSQL 18 的 INSERT 文档说明了冲突处理与 RETURNING 的语义。生产表增加其他唯一约束后,建议显式指定幂等主键作为冲突目标,避免把与幂等键无关的约束冲突也解释成正常重放。实验记录表只有这一项唯一身份,因此省略冲突目标不改变本次测试的含义。
两个连接争同一个键
会话甲开始事务,插入键 commit,写入订单、事件和响应,但暂不提交。会话乙开始事务,插入同一个键。控制脚本不靠固定睡眠猜测竞争已经发生,而是让甲查询乙的会话状态;观察到等待事件类型为 Lock 后,才允许甲提交或回滚。
甲提交的场景里,乙的 INSERT 等待结束,返回零行,SQLSTATE 为成功。乙接着发出新的 SELECT,读到状态二百零一和金额一百的已存响应。最后双方只留下一个订单。乙参与了真实唯一约束竞争,没有拿 Node 内存锁假装数据库并发。
甲回滚的场景使用另一个键 rollback。甲先写的订单和事件随事务撤销。乙的 INSERT 从等待中恢复,返回键,表示乙成为新插入者;脚本才让乙写订单和事件,保存响应并提交。最终仍然只有一笔业务,没有把甲曾经插入过当成已成功。
| 场景 | 甲的结束动作 | 乙的观察 | 后续动作 |
|---|---|---|---|
| 先提交 | COMMIT | INSERT 不返回键 | 新语句读已有响应 |
| 先回滚 | ROLLBACK | INSERT 返回键 | 乙创建订单并提交 |
| 等待超时 | 乙先超时,甲之后提交 | SQLSTATE 55P03 | 乙回滚,再以原键重试 |
这里固定使用读已提交隔离级别。PostgreSQL 18 隔离级别文档说明,不插入的决定可能来自当前语句快照还看不见的竞争事务。因此,不能把“插入或从已有记录读取”的两个分支塞进同一语句,再假定输掉竞争的请求一定能读到对方刚提交的行。
下一条语句会取得新的读视图,这是本文保留两步协议的原因。提高到可重复读或可串行化后,工程师要重新处理整个事务的重试,不能照搬这段等待行为。本文没有运行其他隔离级别,也没有把读已提交下的三个结果外推为所有关系数据库的保证。
等待超时以后,先处理当前事务
第三个场景让乙设置二百五十毫秒的锁等待上限。甲仍然持有键,乙竞争插入时得到锁等待超时。脚本读到 SQLSTATE 55P03,随后执行回滚,恢复会话设置。等甲提交后,乙重新开始事务,沿用原键,读取已有响应。
这个结果区分了两个对象:乙没有获得执行权,甲的业务却可能完成。接口若把乙的超时翻译成“订单创建失败,可以换键再试”,客户端就可能制造第二笔业务。比较合适的契约是提示当前结果尚待确认,并允许用同一键查询或重试;具体采用哪种响应码由接口约定。
锁等待超时与语句超时覆盖的范围不同。前者只计算获取锁的等待,后者限制语句执行时间;二者都不能替代请求整体截止时间。连接池等待、输入校验和响应发送还会花时间。本文只配置锁等待上限来制造竞争故障,没有测量端到端截止时间。PostgreSQL 18 会话参数给出了两个参数的边界。
失败事务中的后续业务语句不能继续当成正常执行。应用捕获数据库异常后,应结束该事务,再决定是否从头重试。连接池归还连接之前也要完成清理,否则下一位借用者可能接到失败事务或遗留的会话参数。实验显式回滚并重置锁等待参数,方便读者看到这一步。
判断重试分支应使用数据库错误码,而不是匹配英文错误消息。锁等待超时和业务唯一冲突走不同路径,本轮分别观察到 55P03 与 23505;它们的定义可查 PostgreSQL 18 错误码表。应用还要把无法确定提交结果的连接中断单列,不能因为连接断开就声称提交没有发生。
响应重放、异参冲突与业务唯一性
提交场景结束后,调用端舍弃第一次结果,再次发送相同键与指纹。新事务插入失败,读取已有响应,订单数不增加。这模拟了客户端对提交结果失去认知后的恢复步骤。它没有模拟 TCP 故障、连接池失联或服务进程被终止,核查记录按这个范围标注。
同一个键换成另一个指纹时,实验读取比较结果为假。真实接口应在返回旧响应之前拒绝这次请求,否则用户购买两百元商品却拿到一百元订单的响应,会误以为新请求已经成功。只比较键而不比较参数,保护了记录数量,却破坏了业务意图。
业务唯一编号还负责跨键去重。实验用新键 newkey 再创建已经存在的业务编号,订单主键触发唯一冲突,整个事务回滚,新键记录也没有留下。它说明传输幂等键和领域身份各自承担一层约束:客户端换了键,领域约束仍能拦住同一采购申请重复成单。
但业务唯一冲突不能无条件转换成成功。服务端需要读取该业务编号的现有订单,核对归属和参数,决定是重放、冲突还是人工处理。文末实验只证明数据库拒绝重复写入,并没有实现这段跨键返回策略。若两个不同采购申请碰巧共享了错误的外部编号,把冲突当成成功会掩盖上游数据问题。
读取旧响应时也要执行当前授权。员工离开采购部门后,保留中的键不应让他继续拿到订单明细。保存结果对象 ID、读取时重新鉴权,通常比长期保存完整个人信息响应更容易处理撤权;代价是重放时还要访问业务表,并明确对象后续变化会不会改变返回内容。
输入校验失败是否占用键,也属于协议的一部分。本例在开始事务前拒绝格式错误,因此修正格式后可以重新发送;已经形成确定业务拒绝的请求,则需要决定是否保存并重放拒绝结果。例如采购额度不足若只回滚,用户补足额度后用原键可能成功;如果业务把首次拒绝视为这次申请的终态,就应把拒绝和指纹一起保存。两种行为面向不同产品预期,数据库约束不会替团队作选择。文末实验没有模拟额度检查,正式接口应把选定行为加入验收样例。
清理响应时保留键的墓碑
只给幂等记录配置过期删除,会出现另一种竞争:乙的 INSERT 因已有记录而返回零行,清理任务删除那一行,乙下一条 SELECT 却读不到。若乙此时重试插入并执行业务,就可能在业务仍然存在时重新执行。领域唯一约束能够拦一部分重复操作,但不能替清理协议定义重放行为。
本例选择保留键、指纹和过期标记,过期只清空响应正文。重放事务读取记录时使用共享行锁,检查指纹和过期状态,把响应复制到进程内后提交。清理事务更新同一行时等待共享锁释放;先完成的重放可以使用已经读出的正文,之后到来的请求看到过期标记。
本轮清理竞争确实运行了两个连接。甲持有共享行锁并读取响应,乙尝试清空正文并设置过期;甲观察到乙在等待锁,随后提交。乙完成更新后,再用原键插入仍然冲突,查询返回“已过期、正文为空、指纹仍在”。键没有重新开放,清理者也没有在甲读结果的中途删掉身份记录。
这里使用共享锁,而非只保护键的较弱锁,因为清理更新的是普通列。PostgreSQL 18 行锁文档说明共享锁会与更新冲突,而只保护键的锁不会阻止所有非键更新。这是协议中一个容易被名字误导的细节,不能认为带有“共享”二字就足够。
墓碑期间,接口可以对相同指纹返回“结果已超过重放窗口,请查询业务对象”,对不同指纹仍返回冲突。本文没有提供一个自动恢复原响应的保证,也不让过期请求重新执行业务。正式表还需要创建时间和重放截止时间;实验用布尔标志手动触发清理,验证锁竞争,没有测试定时扫描器。
这种方案会持续保留身份,存储成本与历史键数量一起增长。团队可以把短期响应和长期身份拆到不同表,但领取执行权仍要检查那份长期身份。若最终必须删除墓碑,就要在契约中定义键何时可以复用,或使用不会与新请求碰撞的命名空间。不可重复的订单仍由持久业务编号约束。
清理任务处理多行时要采用一致的顺序和受控批量,避免长事务占住大量记录。重放路径也应缩短持锁时间,别在共享锁期间调用授权远端服务或向慢客户端发送正文。授权可以先做,进入事务后只进行必要的状态核对;更严格的撤权边界仍需要和权威授权服务约定。
提交订单后再发送外部事件
订单通知采用事务 outbox:订单和“需要发送订单事件”的记录一起提交。发送者读事件,调用下游,再记录发送结果。若下游完成后发送者崩溃,本地仍可能看见未发送状态,恢复后会再次投递。稳定事件身份让接收者有机会去重,发送成功标志本身消除不了这个窗口。
本轮将三条事件向本地接收表插入两遍,接收表按业务身份唯一,最后仍有三行。订单、事件、接收表和幂等记录的总数分别为三、三、三、三。这验证的是本地重复投递夹具,接收表没有执行真实扣款或发送短信,因此不能声称测到了外部副作用恰好一次。
真实接收者也要把去重身份和自己的业务变更放进同一事务。若先登记收到事件再扣款,登记后崩溃会漏扣;若先扣款再登记,崩溃后重投会重复扣。对方只有外部接口且不支持稳定键或结果查询时,本地 outbox 无法替它实现原子性。
AWS 的事务 outbox 说明也要求处理重复消息。本文采用这条设计边界:本地确认订单与发送意图一致,远端是否完成另行核对。外部超时后保留结果未知状态,按照对方的幂等窗口和查询能力恢复,不自动假定它什么都没做。
同一业务流若要求事件有序,还需要业务版本或序号。只有去重 ID 能消除重复,不能阻止较旧事件晚到覆盖较新状态。接收者可以拒绝过期版本,或者按对象串行处理;本文的三条订单相互独立,没有测试乱序传播,不能借用本地行数证明顺序一致。
复现范围与接入步骤
文末脚本通过两个 docker exec 连接独立测试容器,不读取应用连接串。先按后面的准备命令新建一次性 PostgreSQL 18 容器,再把完整代码放入 idem-lab.mjs 运行;代码在 idem 中创建带运行编号的新表,不删除旧表,也不停止容器。
node idem-lab.mjs
文末代码包含建表、会话屏障和断言,执行时会打印实际版本、双方 SQL 与返回结果。终端的标准输出和错误输出来自不同管道,错误文字可能落在下一条命令的展示块中;状态码由同一 psql 会话在失败语句之后读取,判断依据不依赖文字出现顺序。
接入应用时,先把新插入、同参重放、异参冲突和已过期四条分支写成入口协议,再补连接中断后的结果查询。连接中断发生在提交前还是提交后,应用未必知道;恢复流程需要沿用原业务身份向权威数据库核对,不能通过重新生成键来解除不确定性。
本轮已经验证数据库竞争与墓碑清理的核心时序。摘要规范化、真实 HTTP 状态码、进程中断和远端副作用还没有端到端证据。这份实现适合作为接入协议的参考实验;上线验收仍应在真实入口记录故障点、原键与最终业务记录,尤其要覆盖“提交完成但调用方无法确认”的那一次中断。
独立测试库准备与完整脚本
下面命令只用于新建本地一次性测试容器。网络关闭且不映射端口,客户端通过 docker exec 进入容器连接。若已经有本例专用容器,跳过创建命令;不要把名字改成业务数据库容器。镜像按主版本十八选择,脚本还会打印实际补丁版本,本文记录来自十八点六。
docker run -d --name logz-editorial-pg-20260908 --network none \
-e POSTGRES_USER=editorial -e POSTGRES_DB=editorial_depth \
-e POSTGRES_HOST_AUTH_METHOD=trust postgres:18
until docker exec logz-editorial-pg-20260908 \
pg_isready -U editorial -d editorial_depth; do sleep 1; done
node idem-lab.mjs
下面是 idem-lab.mjs 的全部内容。它只依赖本机 Node 与 Docker 命令,psql 由 PostgreSQL 容器提供。本节准备命令供读者复现;本轮实测使用已准备好的专用容器,没有通过这些命令重新创建服务。
import {spawn} from 'node:child_process';
import assert from 'node:assert/strict';
const run=Date.now().toString(),table='idem.records_'+run,orders='idem.orders_'+run,events='idem.events_'+run,sink='idem.sink_'+run;
class Session{constructor(name){this.name=name;this.buf='';this.full='';this.seq=0;this.p=spawn('docker',['exec','-i','logz-editorial-pg-20260908','psql','-X','-qAt','-U','editorial','-d','editorial_depth'],{stdio:['pipe','pipe','pipe']});for(const stream of [this.p.stdout,this.p.stderr])stream.on('data',d=>{this.buf+=d;this.full+=d;this.check?.();});this.p.on('error',e=>this.reject?.(e));}async sql(s){assert(!this.check,'one outstanding command per connection');const mark='END_'+this.name+'_'+(++this.seq);console.log('\n['+this.name+' SQL]\n'+s);return await new Promise((resolve,reject)=>{this.reject=reject;const timer=setTimeout(()=>{this.check=null;reject(Error('session timeout '+this.name));},15000);this.check=()=>{if(this.buf.includes(mark)){clearTimeout(timer);const out=this.buf.slice(0,this.buf.indexOf(mark));this.buf=this.buf.slice(this.buf.indexOf(mark)+mark.length).replace(/^\r?\n/,'');this.check=null;console.log('['+this.name+' OUT]\n'+out.trim());resolve(out.trim());}};this.p.stdin.write(s+'\n\\echo '+mark+'\n');this.check();});}close(){this.p.stdin.end('\\q\n');}}
const A=new Session('A'),B=new Session('B');
const guard=`DO $$ BEGIN IF current_database()<>'editorial_depth' OR current_user<>'editorial' OR current_setting('server_version_num')::int/10000<>18 THEN RAISE EXCEPTION 'wrong database/user/version'; END IF; END $$;`;
const insert=(key)=>`INSERT INTO ${table} VALUES ('tenant1:user1:create','${key}','hash100',NULL,false) ON CONFLICT DO NOTHING RETURNING key;`;
const business=(key,biz=key)=>`INSERT INTO ${orders} VALUES ('tenant1:user1:create','${biz}',100); INSERT INTO ${events} VALUES ('tenant1:user1:create','${biz}'); UPDATE ${table} SET body='{"status":201,"amount":100}' WHERE key='${key}';`;
async function blocked(){for(let i=0;i<100;i++){const out=await A.sql(`SELECT EXISTS(SELECT 1 FROM pg_stat_activity WHERE application_name='idem_B_${run}' AND wait_event_type='Lock');`);if(out==='t')return;await new Promise(r=>setTimeout(r,20));}throw Error('B never entered lock wait');}
try{
for(const s of [A,B])await s.sql(`\\set ON_ERROR_STOP on\n${guard}\nSET application_name='idem_${s.name}_${run}'; SET default_transaction_isolation='read committed'; SELECT current_database(),current_user,current_setting('server_version');`);
await A.sql(`CREATE SCHEMA IF NOT EXISTS idem; CREATE TABLE ${table}(scope text,key text,hash text NOT NULL,body jsonb,expired boolean NOT NULL,PRIMARY KEY(scope,key)); CREATE TABLE ${orders}(scope text,biz text,amount int,PRIMARY KEY(scope,biz)); CREATE TABLE ${events}(scope text,biz text,PRIMARY KEY(scope,biz)); CREATE TABLE ${sink}(scope text,biz text,PRIMARY KEY(scope,biz));`);
for(const mode of ['commit','rollback','timeout']){
await A.sql(`BEGIN; ${insert(mode)} ${business(mode)}`);
if(mode==='timeout')await B.sql("SET lock_timeout='250ms'; \\set ON_ERROR_STOP off");
const pending=B.sql(`BEGIN; ${insert(mode)}\n\\echo SQLSTATE=:SQLSTATE`);await blocked();
if(mode==='timeout'){const out=await pending;assert(out.includes('55P03'));await B.sql("ROLLBACK; SET lock_timeout=0;\n\\set ON_ERROR_STOP on");await A.sql('COMMIT;');await B.sql(`BEGIN; ${insert(mode)} SELECT body FROM ${table} WHERE key='${mode}'; COMMIT;`);}
else{await A.sql(mode==='commit'?'COMMIT;':'ROLLBACK;');const out=await pending;if(mode==='rollback'){assert(out.includes('rollback'));await B.sql(business(mode));}else{assert(!out.split('\n').includes('commit'));}const replay=await B.sql(`SELECT body FROM ${table} WHERE key='${mode}'; COMMIT;`);assert(replay.includes('201'));}
console.log('CASE '+mode+' PASS');}
const retry=await B.sql(`BEGIN; ${insert('commit')} SELECT hash='hash100',body FROM ${table} WHERE key='commit'; COMMIT;`);assert(retry.includes('t|'));console.log('CASE lost-response replay PASS (client discards committed result; no network fault injected)');
const mismatch=await B.sql(`SELECT hash='hash200' FROM ${table} WHERE key='commit';`);assert.equal(mismatch,'f');console.log('CASE same-key different-hash detected PASS');
await A.sql(`BEGIN; SELECT body FROM ${table} WHERE key='commit' FOR SHARE;`);
const clean=B.sql(`BEGIN; UPDATE ${table} SET body=NULL,expired=true WHERE key='commit'; COMMIT;`);await blocked();await A.sql('COMMIT;');await clean;
const tombstone=await A.sql(`BEGIN; ${insert('commit')} SELECT expired,body IS NULL,hash FROM ${table} WHERE key='commit'; COMMIT;`);assert(tombstone.includes('t|t|hash100'));console.log('CASE cleanup waits for replay; tombstone retained PASS');
await B.sql('\\set ON_ERROR_STOP off');const duplicate=await B.sql(`BEGIN; ${insert('newkey')} INSERT INTO ${orders} VALUES ('tenant1:user1:create','commit',100);\n\\echo SQLSTATE=:SQLSTATE\nROLLBACK;`);assert(duplicate.includes('23505'));await B.sql('\\set ON_ERROR_STOP on');
await A.sql(`INSERT INTO ${sink} SELECT * FROM ${events} ON CONFLICT DO NOTHING; INSERT INTO ${sink} SELECT * FROM ${events} ON CONFLICT DO NOTHING;`);
const counts=await A.sql(`SELECT (SELECT count(*) FROM ${orders}),(SELECT count(*) FROM ${events}),(SELECT count(*) FROM ${sink}),(SELECT count(*) FROM ${table});`);assert.equal(counts,'3|3|3|3');console.log('FINAL PASS orders=3 outbox=3 sink=3 records=3; sink is local duplicate-delivery simulation, not external service');
console.log('TABLES '+JSON.stringify({table,orders,events,sink}));
}finally{A.close();B.close();}











