HITL 人工介入的并发难题:两个人同时审批一张证书怎么办
HITL 人工介入的并发难题两个人同时审批一张证书怎么办这是 LangGraph 生产化系列的第二篇。上一篇讲了混合检索、Supervisor 熔断和缓存三防这篇讲一个更隐蔽的领域人工介入HITL的状态管理。完整开源github.com/muyiyang09/…⭐先还原一个真实场景教练上传国职证书 → AI 服务 OCR 核验 风险评估 → 挂起等管理员最终确认。AI 建议是通过但决定权在人。管理员 A 正在核对编号此时管理员 B 也打开了这张审核单。A 点了通过。三秒后B 也点了通过。系统会发生什么如果你答不出来说明你的 HITL 还停留在 Demo 层。这篇文章拆解这背后的三层问题状态存哪里、状态活多久、状态谁说了算。先看官方 Demo 的 HITL 长什么样LangGraph 的 HITL 机制本身很优雅interrupt()暂停图执行状态交给 Checkpointer 持久化人工决定通过Command(resume...)注入后恢复from langgraph.types import interrupt, Command async def hitl_checkpoint(state) - dict: Node 5HITL人工最终确认 组装最终结果。 if settings.hitl_enabled: decision interrupt({ prompt: f证书审核人工确认风险等级{state.get(risk_level)}, fields: state.get(fields), verifications: state.get(verifications), suggestion: state.get(suggestion), }) ...# 管理员提交决定后恢复执行 state_out await CERT_REVIEW_GRAPH.ainvoke( Command(resume{action: payload.action}), config{configurable: {thread_id: thread_id}}, )Demo 里interruptMemorySaver十行代码就能跑。但生产环境的问题是审核单是长生命周期对象——AI 判断只要 2 秒人做决定要 5 分钟、5 小时甚至第二天上班再处理。这中间你的服务要滚动更新、要扩容副本、Redis key 会过期。于是三层问题依次浮出水面。坑一MemorySaver 存 HITL 状态等于把审核单写在便利贴上这是最容易被轻视的坑因为它在开发环境永远不会暴露。MemorySaver 把图状态存在进程内存里。单进程跑 Demo 没问题但生产部署的第一个动作就是多副本 滚动更新这时两个致命场景出现了滚动更新管理员 A 点了通过请求落在副本 1。恰好此时发布新版本副本 1 被杀。管理员 B 刷新页面——审核单消失了因为状态随着副本 1 的进程一起蒸发了。多副本状态分裂AI 审核发起的 interrupt 状态存在副本 1管理员的 resume 请求被负载均衡打到副本 2——副本 2 的 MemorySaver 里根本没有这个 thread 的状态。我们生产配置是SERVICE_ENVprod时强制切换 RedisSaver并且构建过程是层层降级的def build_checkpointer(): 构建顺序Redis → RedisDBCheckpointer 包装 → MemorySaver 回退。 每层都有 fail-open 降级最终保证返回一个可用 Checkpointer。 if settings.checkpointer_backend redis: try: from langgraph.checkpoint.redis import AsyncRedisSaver ttl { default_ttl: settings.checkpoint_ttl_minutes, # 单位是「分钟」 refresh_on_read: True, # 读取时续期活跃 thread 不会被误清 } saver AsyncRedisSaver(redis_urlsettings.redis_url, ttlttl) # DB 灾备包装防止 Redis TTL 过期导致会话中断 if settings.checkpoint_db_fallback: saver RedisDBCheckpointer(saver) return saver except Exception as exc: logger.warning([Checkpoint] RedisSaver 构建失败回退 MemorySaver%s, exc) from langgraph.checkpoint.memory import MemorySaver return MemorySaver()注意两个细节refresh_on_readTrue活跃会话每次读取都续期TTL 只清真的死了的 thread构建失败回退 MemorySaver 而不是拒绝启动Checkpointer 故障不能拖垮整个 AI 服务——宁可退化为单副本可用不可全局 500。一句话总结HITL 状态的生命周期以天计而你的进程生命周期以分钟计。生命周期不匹配的状态不能存在进程里。坑二Redis TTL 到期时管理员正好打开审核单换了 RedisSaver滚动更新的问题解决了。但 TTL 引入了新问题Redis key 过期的瞬间管理员正好点开一张 6 小时前挂起的审核单。resume 请求到达Checkpointer 从 Redis 读 thread 状态——miss。此时要么报审核单不存在要么更糟整条审核流程从头重跑一遍重新 OCR、重新核验、重新 interrupt管理员看到一张复活的审核单。我们给了 Checkpointer 加了一层DB 灾备包装RedisDBCheckpointer每次状态写入时双写一份 JSON 到 MySQLRedis miss 时从 DB 读回、回填 Redis、继续会话。核心读路径长这样async def aget_tuple(self, config): result await self._inner.aget_tuple(config) # 先读 Redis if result is not None: return result # Redis 命中直接返回 thread_id self._thread_id(config) # singleflight并发读只放一个进 DB lock self._locks.setdefault(thread_id, asyncio.Lock()) async with lock: # 空值缓存防穿透确认 DB 无数据后60s 内不再查 last_miss self._empty_cache.get(thread_id) if last_miss and time.time() - last_miss _EMPTY_CACHE_TTL: return None self._empty_cache.pop(thread_id, None) data await session_store.get_state(thread_id) # DB 兜底 if data is None: self._empty_cache[thread_id] time.time() return None # 回填 Redis后续请求走快速路径 await self._inner.aput(config, data[checkpoint], data.get(metadata) or {}, {}) return CheckpointTuple(configconfig, checkpointdata[checkpoint], ...)你会发现这就是第一篇讲的缓存三防原样复用——雪崩靠refresh_on_read读时续期击穿靠 singleflight穿透靠空值缓存。同一个问题域换了皮从教练数据缓存换成会话状态灾备解法是可以整体搬运的。还有一个容易被忽略的对称性设计DB 写失败只告警、不抛异常Redis 已写成功主流程不能被灾备层拖死。灾备层的定位是锦上添花它的故障等级必须低于主链路。一句话总结给 Checkpointer 配灾备不是不信任 Redis是不信任人做决定的速度和 TTL 的交集。坑三两个人同时审批——状态机是第一道防线现在回答开头的问题A 点了通过之后B 再点会发生什么如果没有任何防护B 的 resume 请求会带着Command(resume...)去调ainvoke。此时 thread 的状态已经是终态图已执行到 END——LangGraph 不会替你拦截它会把整张图重新执行一遍。B 的通过变成了一次全新的审核重新 OCR、重新核验、重新 interrupt 挂起。审核单凭空多出一个幽灵流程审计日志里出现两条互相独立的审核记录。我们的防线是一个独立的HITL 状态机和 Checkpointer 解耦# hitl_state.py审核单生命周期 pending → approved / rejected / cancelled终态 # 用 Redis key hitl:{thread_id}:status 记录TTL 24h 与 Checkpointer 对齐 _TTL 86400 _PENDING pending _TERMINAL {approved, rejected, cancelled} def is_terminal(status: str | None) - bool: 是否已是终态不能再 resume。 return status in _TERMINALresume 端点在恢复图执行之前先做冲突检测# 冲突检测已处理过 / 已取消 → 拒绝重复 resume status await hitl_state.get_status(thread_id) if hitl_state.is_terminal(status): raise ConflictError(该审核单已处理过) if status is None: raise NotFoundError(审核单不存在或已过期) state_out await CERT_REVIEW_GRAPH.ainvoke( Command(resume{action: payload.action}), ...) ... await hitl_state.set_status(thread_id, payload.action) # 执行成功后才写终态三个设计决定值得展开为什么不用 Checkpointer 的状态判断而要单独的 keyCheckpointer 存的是图执行状态节点、消息、pending writes这张审核单在业务上是否已结案是业务语义塞进图状态里会让图状态承担业务职责两个抽象互相污染。分开之后审核单列表页可以直接扫hitl:*:status这类 key不必反序列化 checkpoint。为什么执行成功后才写终态而不是先占位先写终态再恢复执行一旦ainvoke中途失败审核单会永久卡在已通过但实际什么都没发生。后写终态的代价是检查-执行-写回之间存在竞态窗口——两个管理员真正同时提交时可能双双通过检测。高并发场景下这里应该升级为Redis Lua 原子 CASGET 比对 SET在脚本内一次完成我们把它标记为 resume QPS 上量后的第一个加固点——当前审核操作是人工低频行为先留竞态窗口、后加原子性是性价比排序。TTL 为什么要和 Checkpointer 对齐24h如果状态机 TTL 更短会出现状态机认为审核单已过期可重开但 Checkpointer 里图还挂着的分裂如果更长会出现审核单显示已通过但 resume 已被拦截的另一种分裂。同一生命周期的状态TTL 必须对齐否则两边会各自腐烂。一句话总结并发审批的本质不是锁问题是业务终态和图执行状态没分家。分了家状态机管幂等Checkpointer 管恢复各司其职。一个贯穿三层的哲学fail-open把三层的失败策略排在一起会发现一个统一的取向层故障场景策略Checkpointer 构建RedisSaver 初始化失败回退 MemorySaver 告警DB 灾备读MySQL 查询失败返回 None当无状态 告警DB 灾备写MySQL 写入失败仅告警主流程继续HITL 状态读写Redis 不可用降级为无状态冲突检测失效每一层的降级都被明确地选择过而不是异常抛到哪里算哪里。我们的原则是审核流程的可用性 冲突检测的完备性——单副本开发环境 Redis 挂了最多是失去防重复 resume 能力审核本身还能跑反之如果这些层选择 fail-closed任何依赖故障就拒绝服务AI 服务会被最外围的一个 Redis 打死。当然 fail-open 有边界涉及资金和合规的动作不允许 fail-open。证书审核给了 suggestion 但决定权在人所以状态检测失效可接受假如哪天加了自动吊销教练资质那一层必须 fail-closed。fail-open 不是默认值是逐层审批的结果。写在最后HITL 的完整生产化清单其实就五件事✅ Checkpointer 用 RedisSaver多副本共享 崩溃恢复refresh_on_read防误清✅ DB 灾备层双写 miss 回填singleflight / 空值缓存一并带上✅ 独立 HITL 状态机管业务终态与图状态解耦TTL 对齐✅ resume 前冲突检测终态拒绝重复恢复⏳ resume 上量后状态翻转升级为 Lua 原子 CAS已标记未实装。第 5 条我特意保留在清单里——诚实标注知道但还没做的加固点比假装系统无懈可击更接近生产级。如果你在准备 AI 工程化方向的面试这块的完整设计推演在仓库文档里多 Agent 实现与 HITL 落地#08上线检查清单#1040 道 Agent 高频面试题#09如果这篇文章帮你绕开了 HITL 的至少一个深坑请到仓库右上角点个 ⭐ Star——这是我持续更新系列的最大动力。 下篇预告《改 prompt 之前你敢不看分数吗给 LangGraph 加一套离线 Eval 基建》——聊聊 Eval 和单元测试的区别、三类 metric 的实现以及离线 95% 通过率 ≠ 线上有效的诚实边界。项目地址github.com/muyiyang09/…