如何用PostgreSQL锁实现应用并发控制,避免非幂等API重复调用?
问题
我的应用某模块需调用外部非幂等的sendMessage API,希望通过PostgreSQL锁实现并发控制,采用常规事务+行级锁而非advisory locks。核心需求如下:
sendMessageAPI调用需恰好执行一次(除非数据中心级灾难);- API响应需存入数据库(除非灾难);
- 每个用户同一时间仅能执行一次
sendMessage调用。
数据库使用SERIALIZABLE隔离级别保障数据完整性,当前简化方案为创建messages表:
CREATE TABLE messages ( userId TEXT NOT NULL, request_payload JSONB NOT NULL, response_payload JSONB, PRIMARY KEY (userId) ); CREATE INDEX pending_msgs ON messages(response_payload IS NOT NULL);
由独立进程插入待发送请求,行处理器通过以下逻辑获取行锁后调用API:
const txn = conn.transaction(); # `BEGIN` under the hood const [{userId, request_payload}] = await txn.execute(``` SELECT * FROM messages WHERE response_payload IS NULL FOR UPDATE SKIP LOCKED LIMIT 1 ```); const response = await http.post(`/sendMessage/${userId}`, payload); await txn.execute('UPDATE messages SET response_payload = $1 WHERE userId = $2', response, userId); await txn.commit();
现有疑问:
SELECT FOR UPDATE是否会触发serialization_failure,进而因重试导致重复调用API?- 未使用
SELECT FOR UPDATE的进程是否可能并发修改该行,导致提交时出现serialization_failure(此时已调用API)? - 文档提到REPEATABLE READ或SERIALIZABLE事务中,若锁定行自事务启动后已变更会抛出错误,但测试中发现对锁定行的更新会阻塞而非触发序列化失败。当前方案是否足够?是否有更严格的PostgreSQL锁机制,确保事务期间无人能修改该行,若有其他进程尝试则直接让对方事务失败?
回答
当前方案的安全性分析
你的核心顾虑是序列化失败导致API重复调用,但当前基于SELECT FOR UPDATE SKIP LOCKED的方案在SERIALIZABLE隔离级别下,能有效规避这类风险,原因如下:
- 行级排他锁的阻塞特性:当事务通过
SELECT FOR UPDATE锁定某行后,其他任何尝试修改该行的操作(无论是否使用SELECT FOR UPDATE)都会被阻塞,直到当前事务提交或回滚。这和你测试中看到的“阻塞而非序列化失败”一致——PostgreSQL会优先尝试获取锁,仅在锁等待超时或死锁发生时才抛出错误,不会直接触发序列化失败。 - SERIALIZABLE级别的兜底保障:极端场景下(如事务启动后该行被其他事务修改并提交),当前事务提交时会触发
serialization_failure,但这种情况在SELECT FOR UPDATE的锁机制下几乎不可能发生,因为锁会直接阻止其他事务修改该行。
彻底避免API重复调用的优化方案
如果要进一步消除极端场景下的风险,可以调整逻辑:
- 在UPDATE语句中增加条件校验:确保仅当
response_payload仍为NULL时才执行更新,即使出现异常情况也能避免无效操作:
UPDATE messages SET response_payload = $1 WHERE userId = $2 AND response_payload IS NULL
执行后检查UPDATE的行影响数,若为0则说明该行已被其他事务处理,直接回滚当前事务,不重试API调用。
- 缩短锁持有时间:API调用属于事务外的耗时操作,长时间持有锁会降低并发度。可以新增
status字段拆分逻辑:
CREATE TABLE messages ( userId TEXT NOT NULL, request_payload JSONB NOT NULL, response_payload JSONB, status TEXT NOT NULL DEFAULT 'pending', PRIMARY KEY (userId) ); CREATE INDEX pending_msgs ON messages(status = 'pending');
处理逻辑改为:
// 第一步:锁定行并标记为处理中,快速提交事务释放锁 const txn1 = conn.transaction(); const [{userId, request_payload}] = await txn1.execute(``` SELECT * FROM messages WHERE status = 'pending' FOR UPDATE SKIP LOCKED LIMIT 1 ```); await txn1.execute('UPDATE messages SET status = $1 WHERE userId = $2', 'processing', userId); await txn1.commit(); // 第二步:事务外调用API,无锁持有 const response = await http.post(`/sendMessage/${userId}`, request_payload); // 第三步:更新响应和状态 const txn2 = conn.transaction(); await txn2.execute(` UPDATE messages SET response_payload = $1, status = 'completed' WHERE userId = $2 `, response, userId); await txn2.commit();
这种方式把锁持有时间压缩到最短,即使第一步出现序列化失败,此时还未调用API,重试不会产生重复请求。
关于“更严格的锁机制”
PostgreSQL的FOR UPDATE行级排他锁已经是最严格的行级锁之一,会阻止其他事务对该行的修改、删除及获取其他排他锁。如果希望其他尝试修改该行的事务直接失败而非阻塞,可以设置锁等待超时:
SET lock_timeout = '100ms';
其他事务尝试获取锁时,若超过设定时间仍未拿到,会直接抛出lock_timeout错误,而非持续阻塞。超时时间需根据业务场景调整,避免正常锁等待被误判。
内容的提问来源于stack exchange,提问作者Jake Biesinger
相关产品推荐
相关产品推荐

