基于QingLedger项目三级锁并发控制机制详解
- 目录
- 一.项目背景
- 二.为什么是"三级锁",而不是一把锁
- 三.核心概念
- 四. 第 1 级Session Lock(会话级锁)
- 1 锁放在哪
- 2 抢锁:一个原子 UPDATE 搞定"三种场景"
- 3 续租与释放
- 五.第 2 级:Request Lock(请求级锁)
- 1 锁放在哪
- 2 幂等锚点:唯一索引兜底
- 3 状态迁移:每个动作都带令牌校验
- 4 请求级"抢锁"的三种路径
- 六.第 3 级:ToolExecution Lock(工具执行级锁)
- 1 锁放在哪
- 2 幂等获取执行记录
- 3 认领已存在记录的四分支
- 4 工具执行的状态迁移
- 5 数据库层的幂等双保险
- 七、三级锁的完整编排:sendMessage 全流程
- 八、心跳与所有权校验:防止"长时间任务被误抢"
- 九、为什么不会死锁
- 十、数据库层的最终兜底:唯一索引 + 事务
- 十一、为什么不用 Redis 分布式锁 / Redisson
- 十二、409 冲突处理策略:为什么不阻塞等待
- 十三、总结:一套锁的三层分工
目录
一.项目背景
QingLedger 是一个"AI 记账"应用:用户用自然语言跟 AI 聊天(“中午吃粉花了 35”),AI 通过大语言模型(LLM)理解意图,调用记账工具,自动把账记到用户的账本里。地址:https://github.com/wsdjxhw/QingLedger
这个场景天然带有几个棘手的并发问题:
- LLM 调用很慢(通常 2~10 秒),用户在多设备上同时发消息,很容易出现"同一时刻两条消息打到同一个会话"。
- LLM 会多次调用工具(一轮
list_categories查分类、一轮create_transaction记账),这些工具调用是有副作用的写操作(往账本里插交易),必须保证只执行一次,否则会记出两笔账。 - 用户可能重复提交(前端网络重试、双击发送),同一个请求不能重复执行。
一句话:要在多设备并发 + LLM 慢调用 + 有副作用的工具调用这三重约束下,保证一个会话同一时刻只有一个请求在处理,且每个工具副作用只落地一次。
这套机制最终被设计成了三级锁。
二.为什么是"三级锁",而不是一把锁
先看一个最朴素的方案:给会话加一把"全局锁",同一时刻只允许一个请求处理。
问题在哪?
- 锁太粗:如果 Session 级锁把整个 LLM 调用(含多轮工具调用)都抱住,那一次请求可能要 10 秒,期间同一个会话的其他请求全部被挡在外面。用户想同时发起两条消息?没门。
- 锁太细:如果只在"单次工具执行"上加锁,那同一会话的两条消息可能交错执行,LLM 上下文就乱了,账也会记串。
所以需要分粒度管理:让锁的粒度匹配"需要保护的资源边界"。
- 会话是一次对话的完整上下文,必须串行 → 会话级锁
- 一次请求是一次完整的 LLM 处理流程,它有自己的状态(处理中/成功/失败)→ 请求级锁
- 一次工具调用是一个幂等副作用,粒度最细,甚至同一请求内的多个工具调用可以并发 → 工具执行级锁
这就是三级锁的来源:
App请求
│
▼
┌──────────────┐
│ Session Lock │ ← 会话级:确保同一时刻只有一个请求在处理
│ (行级UPDATE) │ 粒度最粗,保护整个会话
└──────┬───────┘
│
┌──────▼───────┐
│ Request Lock │ ← 请求级:状态机 processing→success/error
│ (状态机+lease)│ 支持重试失败、认领超时
└──────┬───────┘
│
┌──────▼──────────┐
│ ToolExec Lock │ ← 工具级:单次工具执行的幂等
│ (幂等键+lease) │ 粒度最细
└─────────────────┘
简单来说,会话级锁保证一次只能运行一条请求,请求级锁保证了同一个请求只运行一次,工具级锁保证工具的执行结果不会重复落库
三.核心概念
传统意义的"等待锁"(如 synchronized、Redis SETNX 自旋),但是这三把锁是"租约"(Lease)式的抢锁——抢不到就直接返回 HTTP 409,让客户端别等了。
每个锁都建立在两个核心概念上:
- 租约令牌(lease token):一把随机 UUID,谁持有令牌谁就拥有锁的所有权。后续所有操作(心跳、释放、状态迁移)都要带上令牌做校验,防止别人把你手里的锁抢走后你的操作还生效(这就是经典的"误删别人锁"问题)。
- 心跳(heartbeat):持有者需要定期刷新心跳时间戳。超过阈值(项目里统一是5 分钟)没有心跳,就认为持有者已死,允许其他人抢锁。
这其实是在数据库行上实现的一个分布式锁:用一行UPDATE ... WHERE ...的原子 CAS完成"检查 + 抢占",用 lease token 做所有权校验,用心跳时间戳做过期检测。
原子 CAS(Compare-And-Swap,比较并交换) 是并发编程中最核心的底层原语(Primitive)。你可以把它理解为“无锁编程”的基石——它是一条CPU级别的原子指令,能在多线程环境下,在不使用互斥锁(Synchronized/Lock)的情况下,保证数据修改的线程安全性。
它的工作机制(“读-改-写”三步合一步)
CAS 涉及三个操作数:
内存地址 V(你要改的变量在哪)
期望值 A(你认为这个变量现在应该是什么值)
新值 B(你想把它改成什么值)
执行逻辑(极其简单):
CPU 执行时,瞬间判断:如果 V 当前的值等于 A,就把 B 写入 V;如果不等于,则什么都不做。无论是否修改,都返回 V 原来的值。
为什么叫“原子”?(它是如何保证安全的)
这是理解的关键。在普通的编程中,先比较后交换是两步操作,多线程下必然出现“竞态条件”。而 CAS 是硬件层面支持的:
在 x86 架构下,它对应 CMPXCHG 指令。
CPU 在执行这条指令时,会锁住总线(或缓存行),使得在执行比较和交换的极短时间内,其他线程/核心无法访问该内存地址。
因此,“比较”和“交换”在物理上是一气呵成的,中间不可能被线程切换打断,这就是“原子”的本质。
下面逐级展开。
四. 第 1 级Session Lock(会话级锁)
1 锁放在哪
会话锁的三个字段直接挂在chat_session表上(ChatSession.java:32-39,注释明确写着"分布式锁"):
/** 当前正在执行的请求 ID(分布式锁) */privateLongcurrentRequestId;/** 当前请求的租约令牌(分布式锁) */privateStringcurrentRequestLeaseToken;/** 当前请求的心跳时间(用于检测过期锁) */privateLocalDateTimecurrentRequestHeartbeatAt;一次"会话级锁"就是:把current_request_id、current_request_lease_token、current_request_heartbeat_at三个字段填成当前请求的值。
2 抢锁:一个原子 UPDATE 搞定"三种场景"
抢锁的核心 SQL 在ChatSessionMapper.java:13-27:
UPDATEchat_sessionSETcurrent_request_id=#{requestId},current_request_lease_token=#{leaseToken},current_request_heartbeat_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMPWHEREid=#{sessionId}AND(current_request_idISNULL--① 无锁,直接抢OR(current_request_id=#{requestId}ANDcurrent_request_lease_token=#{leaseToken})--② 重入(同一个请求)ORcurrent_request_heartbeat_atISNULL--③ 心跳为空,可抢ORcurrent_request_heartbeat_at<DATE_SUB(CURRENT_TIMESTAMP,INTERVAL5MINUTE)--③ 心跳超时,可抢)这行 UPDATE 就是一把分布式锁的完整实现,它利用 MySQL 行锁的原子性,把"判断能否抢 + 抢到并写入"合成了一个不可分割的操作。并发时只有一个人能抢到(影响行数 = 1),其他人WHERE不匹配,影响行数 = 0,抢锁失败。
WHERE里的四个条件覆盖了三种场景:
| 条件 | 场景 | 说明 |
|---|---|---|
current_request_id IS NULL | 无锁 | 会话空闲,直接抢 |
current_request_id = requestId AND lease_token = leaseToken | 重入 | 同一个请求再次进入(比如重试自己),幂等放行 |
| 心跳为空 / 心跳超过 5 分钟 | 超时抢锁 | 持有者疑似宕机,其他人可以接管 |
注意:调用方根据 SQL返回的影响行数判断是否抢到:
int返回值 > 0 表示成功(ChatServiceImpl.java:549-551)。
3 续租与释放
续租(心跳)ChatSessionMapper.java:33-43:刷新心跳时间,但必须校验current_request_id和lease_token都匹配,否则说明锁已经易主,不能续:
UPDATEchat_sessionSETcurrent_request_heartbeat_at=CURRENT_TIMESTAMPWHEREid=#{sessionId}ANDcurrent_request_id=#{requestId}ANDcurrent_request_lease_token=#{leaseToken}释放ChatSessionMapper.java:46-58:同样校验令牌,把三个锁字段清空。只有锁的当前持有者能释放,防止"抢了别人锁还帮别人释放"。
UPDATEchat_sessionSETcurrent_request_id=NULL,current_request_lease_token=NULL,current_request_heartbeat_at=NULLWHEREid=#{sessionId}ANDcurrent_request_id=#{requestId}ANDcurrent_request_lease_token=#{leaseToken}五.第 2 级:Request Lock(请求级锁)
1 锁放在哪
请求记录存在chat_request表,它的两个字段构成请求级锁(ChatRequest.java:23-27):
/** 客户端请求幂等 ID(session 内唯一) */privateStringclientRequestId;/** 处理租约令牌(分布式锁) */privateStringleaseToken;请求级锁的精髓是状态机:每条请求有processing / success / error三种状态,所有状态迁移都是带令牌校验的原子 UPDATE。
2 幂等锚点:唯一索引兜底
client_request_id是客户端每次发消息时生成的一个幂等 ID。同一会话内它必须唯一,这是由数据库唯一索引强约束的(V1.0.6__upgrade_chat_agent_schema.sql:27):
UNIQUEKEYuk_session_request(session_id,client_request_id)应用层逻辑在ChatServiceImpl.findOrCreateChatRequest(ChatServiceImpl.java:332-364):
- 先按
session_id + client_request_id查,存在直接复用(返回claimed=false,走抢锁逻辑); - 不存在则插入新记录,并返回
claimed=true; - 如果插入撞了唯一键(并发下两个请求同时插),捕获异常后回落查询,把别人先插入的那条拿出来用。
这就是"应用层幂等键 + 数据库唯一索引"的双保险。
3 状态迁移:每个动作都带令牌校验
请求级锁的核心操作都在ChatRequestMapper.java:
| 操作 | 方法 | 行号 | 语义 |
|---|---|---|---|
| 按幂等键查 | selectBySessionAndClientRequest | 16-18 | 查询已有请求 |
| 认领超时请求 | claimProcessingRequest | 21-30 | status='processing'且updated_at < 5分钟前→ 换新令牌 |
| 重试失败请求 | retryErroredRequest | 33-46 | status='error'→ 重置为processing,换新令牌,清空旧结果 |
| 回退误认领 | restoreClaimedProcessingRequest | 49-60 | 更新了请求锁修改了数据但是抢会话锁失败后,把误认领的令牌还原,并且还原数据 |
| 刷新心跳 | touchProcessingRequest | 63-71 | 校验lease_token后刷新updated_at |
| 标记成功 | markSuccess | 74-89 | processing+ 令牌匹配 →success,写入回复 |
| 标记失败 | markError | 92-102 | processing+ 令牌匹配 →error,写入错误信息 |
以markSuccess为例(ChatRequestMapper.java:74-89),状态迁移必须同时满足"状态正确 + 令牌匹配",否则影响行数为 0,调用方会抛 409"请求已被其他执行流接管":
UPDATEchat_requestSETstatus='success',reply=#{reply}, ...WHEREid=#{id}ANDstatus='processing'ANDlease_token=#{leaseToken}4 请求级"抢锁"的三种路径
请求级锁的抢锁逻辑在sendMessage里(ChatServiceImpl.java:122-152),根据请求当前状态分三条路径:
路径 1:status = success → 不走任何锁,直接返回缓存结果 路径 2:status = error → 先抢会话锁,成功则把请求重置为 processing 重试 路径 3:status = processing(处理中) ├─ 未超时(updated_at 在 5 分钟内)→ 抛 409"请求处理中,请勿重复提交" └─ 已超时 → 抢锁(claimProcessingRequest)+ 抢会话锁超时判断在tryClaimStaleProcessing(ChatServiceImpl.java:429-435):
privatebooleantryClaimStaleProcessing(ChatRequestrequest,StringnewLeaseToken){LocalDateTimeupdatedAt=request.getUpdatedAt();if(updatedAt!=null&&updatedAt.isAfter(LocalDateTime.now().minusSeconds(REQUEST_STALE_SECONDS))){returnfalse;// 300 秒内没超时,不能抢}returnchatRequestMapper.claimProcessingRequest(request.getId(),newLeaseToken)>0;}其中REQUEST_STALE_SECONDS = 300(ChatServiceImpl.java:36),即 5 分钟。
细节:路径 3 抢到请求锁之后,如果抢会话锁失败,要把请求的令牌和updated_at回退到原来的值——因为你的抢锁动作只是"临时借用",既然会话锁没抢到,就必须把请求还原,不能让别人的请求以为被抢了。这段逻辑在restoreClaimedProcessingRequest(ChatServiceImpl.java:443-453)。
六.第 3 级:ToolExecution Lock(工具执行级锁)
1 锁放在哪
工具执行记录存在chat_tool_execution表,核心字段(ChatToolExecution.java:23-30):
/** 参数 SHA-256 哈希(用于幂等判断) */privateStringargumentsHash;/** 幂等键:requestId + toolName + argumentsHash */privateStringidempotentKey;/** 处理租约令牌(分布式锁) */privateStringleaseToken;工具级锁的核心是"幂等键":requestId + 工具名 + 参数哈希。同一个请求、调用同一个工具、参数也一模一样,就认为是同一次工具执行,绝不能执行两次。
这里有个巧妙的点——参数要先规范化再哈希(normalizeJson,BookingAgent.java:737-747):把 JSON 的 key 按字典序排序(TreeMap)再序列化。这样 LLM 返回{"amount":"35","type":"expense"}和{"type":"expense","amount":"35"}会算出同一个哈希,从而正确识别为同一次调用。
2 幂等获取执行记录
acquireExecution(BookingAgent.java:316-343)是工具级锁的入口:
- 先按
idempotentKey查(selectByIdempotentKey); - 不存在就插一条新记录(状态
processing,生成新令牌); - 插入撞唯一键(并发下两个人同时插,
uk_idempotent_key/uk_request_tool_args兜底)→ 回落查询,拿到别人那条,再走认领逻辑。
3 认领已存在记录的四分支
claimExistingExecution(BookingAgent.java:355-380)对已存在的记录分四支处理:
| 状态 | 处理 |
|---|---|
success | 直接复用,返回缓存结果(工具已执行过,不重复执行) |
error | retryErroredExecution重置为processing,重新执行 |
processing且已超时(5 分钟无心跳) | claimStaleExecution抢锁,换新令牌 |
processing且未超时 | 返回 null,调用方视为"操作冲突",放弃本轮 |
4 工具执行的状态迁移
工具级状态迁移与请求级几乎一样,只是多存了result_json和transaction_id(ChatToolExecutionMapper.java):
retryErroredExecution(18-29):error → processingtouchProcessingExecution(32-40):心跳续租claimStaleExecution(43-52):超时抢锁markSuccess(55-68):processing → success,写入结果和交易 IDmarkError(71-82):processing → error,写入错误结果
注意这里有个硬约束:create_transaction工具在transactionTemplate.execute(...)事务内完成"创建交易 + 标记工具成功",二者要么都成功要么都回滚(BookingAgent.java:256-293)。这就保证了:不会出现"交易建了,但工具记录没标记成功"的中间状态。
5 数据库层的幂等双保险
工具级幂等在数据库层有两条唯一索引(V1.0.6__upgrade_chat_agent_schema.sql:66-67)做最终兜底:
UNIQUEKEYuk_request_tool_args(request_id,tool_name,arguments_hash),UNIQUEKEYuk_idempotent_key(idempotent_key)Redis 仅作为热点缓存,加速命中;数据库记录才是资金写操作的最终防重依据。
七、三级锁的完整编排:sendMessage 全流程
把三级锁串起来的是sendMessage(ChatServiceImpl.java:109-225)。完整时序如下:
客户端发消息 (sessionId, content, clientRequestId) │ ├─ 1. findOrCreateChatRequest 查/建请求记录(唯一索引兜底幂等) │ ├─ 2. status == success ? ──是──→ 直接返回缓存结果,不碰锁 │ ├─ 3. status == error ? ──是──→ 抢 Session Lock → retryErroredRequest │ ├─ 4. status == processing ? ──┬─ 超时(5分钟)? → 抢 Request Lock → 抢 Session Lock │ └─ 未超时 → 抛 409"请求处理中,请勿重复提交" │ ├─ 5. 抢到 Session Lock 后: │ tryRecoverCompletedRequest → 结果已落库?直接返回(防重复执行 LLM) │ tryRecoverToolSideEffect → 交易已建但回复没生成?恢复交易结果 │ ├─ 6. 调用 bookingAgent.run(...) ← 进入 Agent 执行(第 2、3 级锁在此展开) │ ├─ LLM 多轮循环(最多 3 轮) │ │ 每轮:assertRequestOwnership → LLM 调用(心跳保活) → assertRequestOwnership │ │ 有 tool_calls 就进入 executeTool(工具级幂等) │ └─ 生成最终回复 │ ├─ 7. markRequestSuccess 请求状态机 → success,带令牌校验 │ └─ 8. finally: releaseSessionExecution 释放会话锁关键点回顾:
- 第 5 步的"恢复"是防重复执行的关键:如果 LLM 已经跑完、结果已落库,但客户端因为网络重发了一次,
tryRecoverCompletedRequest会通过消息的dedupe_key直接找到已持久化的回复并返回,不再重新调用 LLM。这一步避免了"白白花一次 LLM 费用 + 重复建账"。 - 第 8 步用 finally 释放会话锁,保证无论成功、失败还是抛异常,会话锁最终都会被释放,不会把会话锁死。
八、心跳与所有权校验:防止"长时间任务被误抢"
LLM 调用 + 工具执行是长时间任务(一次可能几秒到十几秒,多轮可能更长)。如果只是"抢锁时不看心跳、持有后也不续租",那别的请求就会在 5 分钟超时后把它抢走,导致两个执行流并发操作同一个会话。
所以项目做了一个定时心跳守护executeWithLeaseHeartbeat(BookingAgent.java:853-889):
- 用
ScheduledExecutorService启动一个每 60 秒(LEASE_HEARTBEAT_SECONDS = 60,BookingAgent.java:71)运行一次的心跳线程; - 每次心跳调用
assertRequestOwnership(校验请求 + 会话两层所有权)和touchProcessingExecution(刷新工具执行心跳); - 心跳失败(说明所有权丢了)会记录
heartbeatFailure,任务结束后抛 409; - 任务完成后
future.cancel(true)+shutdownNow()关掉线程。
所有权校验assertRequestOwnership(BookingAgent.java:892-898)同时校验两层锁,任一失败即抛 409:
privatevoidassertRequestOwnership(LongsessionId,LongrequestId,StringleaseToken){intrequestRows=chatRequestMapper.touchProcessingRequest(requestId,leaseToken);intsessionRows=chatSessionMapper.touchSessionExecution(sessionId,requestId,leaseToken);if(requestRows<=0||sessionRows<=0){thrownewBusinessException(409,"请求已被其他执行流接管");}}这等于给"长时间执行"加了一层主动保护:只要我在跑,就一直续租;一旦续不上,立刻停下来报错,绝不跟别人并行改数据。
九、为什么不会死锁
三个锁的获取顺序固定:Session → Request → ToolExecution。
获取顺序(必须一致): Session Lock → Request Lock → ToolExecution Lock为什么这个顺序不会死锁?因为死锁的必要条件是"多个线程以不同顺序获取多把锁,互相等对方持有的锁"。只要所有线程都按同一个方向、同一个顺序拿锁,就不会出现循环等待。释放顺序相反(先释放最内层,最后释放会话锁)。
而且注意:这把锁不是阻塞等待式的,抢不到直接 409 返回,天然杜绝了"无限期等锁"导致的死锁——冲突的一方直接放弃,根本不会停在锁上等。
十、数据库层的最终兜底:唯一索引 + 事务
三级锁再完善,也扛不住"极端并发下的竞态窗口"(比如抢锁逻辑里"先查后插"的两步间隙)。所以项目在数据库层做了最终兜底:
| 表 | 唯一索引 | 作用 |
|---|---|---|
chat_request | uk_session_request (session_id, client_request_id) | 请求级幂等强约束 |
chat_tool_execution | uk_request_tool_args (request_id, tool_name, arguments_hash) | 工具级幂等强约束 |
chat_tool_execution | uk_idempotent_key (idempotent_key) | 幂等键唯一 |
chat_message | uk_message_dedupe_key (dedupe_key) | 消息去重,防重复插入 |
应用层对这些冲突的处理方式非常统一:捕获DuplicateKeyException/DataIntegrityViolationException,当成"别人已经做过了",回落查询拿到既有数据,而不是当作系统错误。例如:
- 消息插入:
insertMessageIgnoreDuplicate捕获DuplicateKeyException后直接忽略(ChatServiceImpl.java:617-623、BookingAgent.java:839-845) - 账本并发重复添加成员:
addMemberToLedger捕获DataIntegrityViolationException(LedgerServiceImpl.java:491-502)
这种"先乐观插入,撞了唯一键就回落查询"的模式,是这套系统的并发安全底座。
十一、为什么不用 Redis 分布式锁 / Redisson
项目里没有引入 Redisson,也没有用SET NX自旋锁,原因有三:
- 幂等兜底必须落在数据库:这是记账场景,最终要防的是"同一笔交易建两次"。这个约束只能靠
chat_tool_execution的唯一索引 + 事务来强保证,Redis 做不到持久化级别的兜底。 - MySQL 行级原子 UPDATE 天然够用:会话锁是"一行一个锁",
UPDATE ... WHERE ...本身就是原子的 CAS,不需要额外引入分布式锁中间件。 - 成本与复杂度:三级锁的粒度设计让并发冲突率极低(同一会话同一时刻基本只有一个请求),MySQL 方案在数据一致性上是更稳的锚点。Redis 在这里只承担了热点读缓存、验证码限流等轻量职责。
一句话:这个项目选择了"数据库作为唯一事实来源",而不是用 Redis 做锁再想方设法跟数据库保持一致。
十二、409 冲突处理策略:为什么不阻塞等待
三级锁在"抢不到"时统一返回HTTP 409 Conflict,而不是让请求排队等待。这是有意为之的设计:
- LLM 调用动辄 2~10 秒,让客户端阻塞等待的体验很差;
- 返回 409 后,前端可以展示"上一消息正在处理中",然后通过轮询消息列表来感知处理完成,再决定是否允许用户发新消息;
- 实现简单,不需要队列、回调、长轮询。
冲突时的报错信息很明确(ChatServiceImpl.java:127、143):
"当前会话仍有消息处理中,请稍后再试" "请求处理中,请勿重复提交" "请求已被其他执行流接管"十三、总结:一套锁的三层分工
把三级锁用一句话概括:同一个会话不能并行,同一个请求不能重跑,同一个工具不能重复记账。
| 层级 | 表 | 锁字段 | 粒度 | 保护对象 |
|---|---|---|---|---|
| 第 1 级 Session | chat_session | current_request_id / _lease_token / _heartbeat_at | 粗 | 同一会话串行 |
| 第 2 级 Request | chat_request | client_request_id(幂等键) + lease_token | 中 | 同一请求不重跑(状态机) |
| 第 3 级 ToolExec | chat_tool_execution | idempotent_key(幂等键) + lease_token | 细 | 同一工具副作用只落地一次 |
底层是同一套"分布式锁三件套"在不同粒度上的复用:
- 行级原子
UPDATE ... WHERE ...(CAS)实现抢锁; - lease token 校验所有权,防止误释放、误迁移;
- 心跳 + 5 分钟超时实现崩溃自动接管。
再加上两个不可忽视的底座:
- 数据库唯一索引做并发竞态窗口的最终兜底;
@Transactional事务保证"工具执行记录 + 交易创建"的原子性。
这套设计把分布式系统里的"并发控制"问题,下沉到了数据库这一层原子能力上,用最朴素的 SQL 解决了最麻烦的并发一致性问题。