多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

WMS集成AI Agent生产部署:Semaphore与Flux.defer并发安全实战

WMS集成AI Agent生产部署:Semaphore与Flux.defer并发安全实战 1. 从一次凌晨三点的压测事故说起WMS 仓储系统集成 AI Agent 这个系列做到第 8 讲前面七讲把 Agent 的对话链路、工具调用、SSE 流式推送、会话上下文管理都跑通了本地开发环境里一切顺滑。但真正把服务推到生产、挂上并发压测的那一刻问题才集中爆发——这也是我把终篇留给生产部署与并发安全的原因。那天晚上我用 JMeter 起了 200 个并发线程去打 SSE 接口前 30 秒一切正常第 40 秒开始大量请求卡死日志里刷屏的是Semaphore获取超时紧接着是Flux的onErrorDropped最后整个 Agent 服务线程池被打满连健康检查接口都返回 503。我盯着屏幕从凌晨一点查到四点最后定位到的根因非常反直觉Semaphore 被放在了Flux.defer外面。这个坑的隐蔽性在于本地单机、低并发、短连接场景下它完全不暴露只有生产环境的高并发 SSE 长连接 慢下游三者叠加才会触发。所以这篇终篇我不打算再讲 Agent 怎么调工具而是把生产部署这一关的并发安全讲透Semaphore 到底该放在哪、Flux.defer 的求值时机为什么关键、SSE 长连接下的信号量该怎么设计、压测时怎么快速定位这类问题。如果你正在把任何带流式输出的 Spring Boot 服务往生产推这篇内容应该能帮你省掉一个通宵。关键词里出现的 Semaphore、Flux.defer、Reactor、Spring Boot、SSE正好构成了这条链路的五个关键节点我会沿着信号量语义 → Reactor 求值时机 → SSE 连接生命周期 → 压测定位 → 生产参数这条线一路讲下去。2. Semaphore 与 Flux.defer 的求值时机到底谁先谁后2.1 先搞清楚 Semaphore 在响应式链路里扮演什么角色在 WMS 集成 AI Agent 的场景里Semaphore 的用途很明确限制同时打到下游大模型或工具服务的并发请求数。因为 Agent 一次对话可能触发多个工具调用每个工具调用又可能是一次 HTTP 请求如果不加限制200 个用户同时提问就会瞬间产生几百个下游请求把下游打挂或者触发限流。传统阻塞式代码里Semaphore 的用法很直观private final Semaphore semaphore new Semaphore(20); public String callDownstream(String input) throws InterruptedException { semaphore.acquire(); try { return doHttpCall(input); } finally { semaphore.release(); } }acquire()阻塞当前线程拿到许可才继续release()在 finally 里保证归还。这套逻辑在同步世界里天经地义。但到了 Reactor 里问题来了Reactor 的线程模型是事件驱动的你不能阻塞线程去等信号量。一旦在Flux的某个操作符里调用阻塞的acquire()就会占用 event loop 线程整个流的吞吐直接崩掉。所以响应式场景下必须用非阻塞的信号量获取也就是tryAcquire()配合重试或者用 Reactor 自带的Sinks、RateLimiter之类的工具。我当时的实现是这样的错误版本private final Semaphore semaphore new Semaphore(20); public FluxString streamAgent(String question) { return Flux.defer(() - { if (!semaphore.tryAcquire()) { return Flux.error(new RuntimeException(并发超限)); } return agentService.stream(question) .doFinally(signal - semaphore.release()); }); }看起来没问题对吧Flux.defer里做tryAcquire流结束时release。但压测时就是这里出的问题。2.2 Flux.defer 的真正语义每次订阅才求值Flux.defer的核心语义是它接收的是一个 Supplier只有在每次订阅发生时才会调用这个 Supplier 去创建真正的 Flux。这是它和Flux.just、Flux.fromIterable这类创建时即求值的操作符最本质的区别。用生活化的类比Flux.just(x)像是你提前把菜做好放桌上谁来吃都是这盘Flux.defer(() - cook())像是每次有客人来你才现做每个客人吃到的都是新做的。这个区别在信号量场景下是致命的。看下面两种写法// 写法 ASemaphore 在 defer 外面 public FluxString streamA(String q) { semaphore.tryAcquire(); // 订阅前就执行了 return Flux.defer(() - agentService.stream(q) .doFinally(s - semaphore.release())); } // 写法 BSemaphore 在 defer 里面 public FluxString streamB(String q) { return Flux.defer(() - { if (!semaphore.tryAcquire()) { return Flux.error(new RuntimeException(并发超限)); } return agentService.stream(q) .doFinally(s - semaphore.release()); }); }写法 A 的问题在于streamA方法被调用的那一刻也就是 Flux 被组装的那一刻tryAcquire就执行了。但此时流还没被订阅甚至可能永远不会被订阅。如果调用方组装了 Flux 但因为某种原因没订阅这个许可就永远泄漏了。更糟的是如果同一个 Flux 被多次订阅比如retry、repeat、或者被多个下游共享tryAcquire只会执行一次但release会执行多次信号量计数直接错乱。写法 B 才是正确的每次订阅都重新tryAcquire每次流结束都release一一对应。2.3 我踩的那个坑defer 里套 defer但真实事故比这更绕。我当时的代码其实是这样的public FluxString streamAgent(String question) { return Flux.defer(() - { Semaphore localSem new Semaphore(20); // 每次订阅都 new 一个 return Flux.defer(() - { if (!localSem.tryAcquire()) { return Flux.error(new RuntimeException(并发超限)); } return agentService.stream(question) .doFinally(s - localSem.release()); }); }); }外层defer每次订阅都 new 一个新的 Semaphore内层defer用这个新 Semaphore 做限流。看起来每次订阅独立限流实际上限流完全失效——因为每个订阅都有自己的 Semaphore20 个许可对单个订阅来说根本用不完200 个并发订阅就是 200 个独立的 Semaphore下游照样被打爆。这就是压测时前 30 秒正常、后面崩的真相前 30 秒并发还没起来每个订阅独立限流看不出问题并发一上来下游被打挂Agent 服务开始堆积请求SSE 连接迟迟不结束线程池和连接池双双耗尽。提示Flux.defer里 new 出来的对象生命周期是每次订阅一份。如果你希望限流是全局的Semaphore 必须是单例且tryAcquire必须在defer内部执行。3. SSE 长连接把并发问题放大了一个量级3.1 SSE 和普通 HTTP 请求的本质差异普通 REST 接口是请求-响应-关闭连接生命周期以毫秒计。SSE 是请求-持续推送-客户端或服务端关闭连接生命周期以秒甚至分钟计。这个差异对并发模型的影响是数量级的。假设下游 Agent 一次对话平均耗时 8 秒SSE 连接要保持 8 秒。如果并发 200意味着任意时刻都有 200 个活跃连接。而普通接口如果 QPS 是 200、耗时 8 秒那并发也是 200——但普通接口的 200 并发是瞬时的SSE 的 200 并发是持续的。更麻烦的是SSE 连接会占用一个 Servlet 容器线程如果用 Spring MVC 的SseEmitter或者一个 Reactor 的订阅如果用 WebFlux一个 TCP 连接一份会话上下文内存在 WMS 场景里Agent 还要维护对话历史、工具调用状态每个连接的内存开销比普通接口大得多。200 个并发 SSE 连接内存占用可能是普通接口的 10 倍以上。3.2 WebFlux SSE 下的信号量设计如果你用的是 WebFlux我推荐 Agent 服务用 WebFlux因为天然支持背压和流式SSE 的实现大概是这样GetMapping(value /agent/stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxServerSentEventString stream(RequestParam String question) { return agentService.stream(question) .map(chunk - ServerSentEvent.builder(chunk).build()) .timeout(Duration.ofMinutes(5)) .onErrorResume(e - Flux.just( ServerSentEvent.builder([ERROR] e.getMessage()).build())); }这里的并发控制要分两层考虑第一层是连接数限制。不能让无限多的 SSE 连接同时存在否则内存和文件描述符都会爆。这一层可以用全局 Semaphore 控制在Flux.defer里tryAcquire在doFinally里release。第二层是下游并发限制。即使连接数控制住了每个连接内部可能触发多个工具调用下游并发还是要单独限。这一层建议用 Reactor 的flatMap配合concurrency参数或者用 Resilience4j 的RateLimiter。我最终的实现是这样的private final Semaphore connectionSemaphore new Semaphore(200); private final Semaphore downstreamSemaphore new Semaphore(50); GetMapping(value /agent/stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxServerSentEventString stream(RequestParam String question) { return Flux.defer(() - { if (!connectionSemaphore.tryAcquire()) { return Flux.just(ServerSentEvent.builder([BUSY] 服务繁忙请稍后重试).build()); } return agentService.stream(question) .flatMap(chunk - callToolWithLimit(chunk), 4) // 下游并发限制为 4 .map(chunk - ServerSentEvent.builder(chunk).build()) .timeout(Duration.ofMinutes(5)) .doFinally(signal - connectionSemaphore.release()) .onErrorResume(e - Flux.just( ServerSentEvent.builder([ERROR] e.getMessage()).build())); }); } private MonoString callToolWithLimit(String chunk) { return Mono.defer(() - { if (!downstreamSemaphore.tryAcquire()) { return Mono.just([SKIP] 下游繁忙); } return toolService.call(chunk) .doFinally(s - downstreamSemaphore.release()); }); }注意flatMap的第二个参数4这是 Reactor 的并发参数表示同时订阅的上游元素数量。配合downstreamSemaphore可以做到最多 4 个工具调用并发且全局下游并发不超过 50。3.3 连接超时和心跳别让僵尸连接吃掉许可SSE 长连接还有一个隐蔽的坑客户端断开了但服务端不知道。TCP 连接可能因为网络问题变成半开状态服务端还在傻傻地推送Semaphore 许可一直不释放。解决办法有两个一是设置合理的timeout。上面代码里的.timeout(Duration.ofMinutes(5))就是兜底5 分钟没数据就主动结束流触发doFinally释放许可。二是加心跳。SSE 协议本身支持注释行作为心跳FluxServerSentEventString heartbeat Flux.interval(Duration.ofSeconds(15)) .map(i - ServerSentEvent.builder().comment(heartbeat).build()); return Flux.merge(agentStream, heartbeat) .doFinally(s - connectionSemaphore.release());心跳的作用是让客户端和服务端都知道连接还活着同时也能让中间的负载均衡器不因为空闲而断开连接。很多云厂商的 LB 默认 60 秒空闲就断SSE 如果 60 秒没数据就会被切断加个 15 秒心跳就能规避。注意心跳流和业务流 merge 之后doFinally会在两个流都结束时才触发。如果业务流先结束心跳流还在跑许可就不会释放。所以心跳流要用takeUntilOther或者和业务流做正确的组合确保业务结束时心跳也停。4. 压测那一晚从现象到根因的完整排查链路4.1 现象描述三个反常信号压测开始后我观察到三个反常现象信号一JMeter 报告前 30 秒 TPS 稳定在 180 左右第 40 秒开始断崖式下跌到 20 以下且不再恢复。信号二服务端日志里Semaphore相关的异常从第 40 秒开始密集出现但异常信息是并发超限说明tryAcquire返回了 false。信号三jstack抓线程栈发现大量线程卡在reactor.netty的 event loop 上但 CPU 使用率只有 30%不是计算瓶颈。这三个信号组合起来指向一个结论有许可泄漏导致信号量被耗尽后续请求全部被拒。4.2 排查第一步确认是不是许可泄漏许可泄漏的典型特征是信号量计数只减不增。我在代码里加了一行日志每次tryAcquire和release都打印当前availablePermits()log.info(acquire, available{}, connectionSemaphore.availablePermits());压测重跑日志显示availablePermits从 200 一路降到 0然后就不再回升。即使所有 SSE 连接都已经关闭客户端收到 FIN许可也没回来。这就坐实了泄漏。问题是doFinally明明写了release为什么没执行4.3 排查第二步doFinally 到底有没有被调用我在doFinally里也加了日志.doFinally(signal - { log.info(release, signal{}, available{}, signal, connectionSemaphore.availablePermits()); connectionSemaphore.release(); })重跑压测发现release日志的数量远少于acquire日志。也就是说有一部分订阅根本没有走到doFinally。这就奇怪了。Reactor 的doFinally保证在流终止时无论是 complete、error 还是 cancel都会执行。除非……订阅根本没建立或者流被丢弃了。4.4 排查第三步回到 Flux.defer 的求值时机我重新审视代码终于发现了问题。当时的代码结构是这样的public FluxServerSentEventString stream(String question) { FluxServerSentEventString flux Flux.defer(() - { if (!connectionSemaphore.tryAcquire()) { return Flux.just(busyEvent()); } return agentService.stream(question) .map(this::toEvent) .doFinally(s - connectionSemaphore.release()); }); // 这里有个缓存 return cache.putIfAbsent(question, flux); }问题出在最后那行cache.putIfAbsent。这是一个相同问题复用同一个 Flux的优化本意是减少重复计算。但Flux.defer返回的 Flux 被缓存后同一个 Flux 实例被多个订阅者订阅。第一次订阅tryAcquire成功许可 -1流结束release许可 1。看起来正常。但如果第一个订阅还没结束第二个订阅就来了第二个订阅也会触发defer的 Supplier再次tryAcquire许可再 -1。两个订阅共享同一个 Flux 实例但doFinally会执行两次release两次。如果第一个订阅因为某种原因被 cancel 而没有执行doFinally比如客户端断开导致 Reactor 直接 cancel 而不走终止信号许可就泄漏了。更隐蔽的是Flux.defer的 Supplier 在每次订阅时都会执行但如果 Supplier 内部抛异常或者返回的 Flux 在订阅时立即出错doFinally可能不会按预期执行。4.5 根因确认与修复最终确认的根因是Flux 被缓存复用 客户端异常断开导致 cancel 信号未正确传播 doFinally 在某些 cancel 路径下未执行。修复方案有三条第一去掉 Flux 缓存。每个请求独立创建 Flux不共享实例。Agent 对话本来就是有状态的缓存 Flux 本身就是错误设计。第二用using操作符替代手动 tryAcquire/release。Reactor 提供了Flux.using它保证资源在流终止时一定被释放return Flux.using( () - { if (!connectionSemaphore.tryAcquire()) { throw new RuntimeException(并发超限); } return connectionSemaphore; }, sem - agentService.stream(question).map(this::toEvent), sem - sem.release() );using的第三个参数是资源清理函数Reactor 保证它在流终止包括 cancel时一定执行。第三加兜底超时。即使using也有极端情况加一个.timeout()作为最后防线。修复后重跑压测200 并发稳定跑 30 分钟availablePermits始终在 0-200 之间正常波动没有再出现泄漏。5. 生产部署时那些文档不会告诉你的参数5.1 WebFlux 的 event loop 线程数WebFlux 默认的 event loop 线程数是CPU 核数 * 2。在 8 核机器上就是 16 个线程。如果 SSE 连接里有阻塞操作比如同步的 HTTP 调用、数据库查询这 16 个线程很快就会被占满。我的建议是Agent 服务里所有下游调用必须是非阻塞的。如果用 WebClient确保用的是 Reactor Netty 的异步客户端如果用 RestTemplate那就要把阻塞调用调度到boundedElastic线程池return Mono.fromCallable(() - restTemplate.postForObject(url, req, String.class)) .subscribeOn(Schedulers.boundedElastic());boundedElastic的默认线程数是CPU 核数 * 10可以通过reactor.schedulers.defaultBoundedElasticSize调整。但要注意这个线程池是有限的如果阻塞调用太多一样会排队。5.2 SSE 连接数的操作系统级限制Linux 默认的单进程文件描述符限制是 1024。每个 SSE 连接占一个 fd加上其他开销实际能支撑的 SSE 连接可能只有 800 左右。生产环境必须调大ulimit -n 65535同时在 systemd 服务文件里加[Service] LimitNOFILE65535如果是容器环境还要注意容器的 fd 限制和宿主机的限制是两回事都要调。5.3 反向代理的超时配置Nginx 作为反向代理时默认的proxy_read_timeout是 60 秒。SSE 连接如果 60 秒没数据就会被 Nginx 切断。必须调大location /agent/stream { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Connection ; proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s; proxy_send_timeout 3600s; chunked_transfer_encoding on; }proxy_buffering off是关键否则 Nginx 会缓冲 SSE 数据导致客户端收不到实时推送。Connection 是为了清掉默认的close保持长连接。5.4 JVM 参数与 GC 调优SSE 长连接场景下堆内存里会积累大量半成品的响应对象。如果 GC 不及时老年代很快被占满触发 Full GC所有连接都会卡顿。我的配置是-Xms4g -Xmx4g -XX:UseG1GC -XX:MaxGCPauseMillis200 -XX:InitiatingHeapOccupancyPercent45G1 的MaxGCPauseMillis200保证单次 GC 停顿不超过 200ms对 SSE 的实时性影响可控。InitiatingHeapOccupancyPercent45让 G1 更早启动并发标记避免 Full GC。5.5 监控指标哪些必须埋点生产环境必须监控这几个指标否则出问题只能靠猜指标含义告警阈值semaphore.available当前可用许可数持续为 0 超过 1 分钟sse.active.connections活跃 SSE 连接数超过连接上限的 80%reactor.netty.connectionsNetty 连接数超过 fd 限制的 70%agent.downstream.latency下游调用 P99 延迟超过 5 秒jvm.gc.pauseGC 停顿时间P99 超过 500ms这些指标用 Micrometer 暴露到 Prometheus配 Grafana 面板出问题一眼就能定位。6. 几个容易忽略的边界场景6.1 客户端主动断开时的资源清理SSE 客户端断开有两种方式正常关闭发送 FIN和异常断开网络中断。正常关闭时Reactor 会收到 cancel 信号doFinally会执行。异常断开时TCP 连接可能变成半开服务端要等到下一次写数据失败才知道客户端没了。这就是为什么心跳很重要。心跳写失败会触发 error 信号进而触发doFinally。没有心跳的话服务端可能几分钟都不知道客户端已经走了许可一直不释放。6.2 下游超时与信号量的关系如果下游调用超时doFinally会执行许可会释放。但如果下游调用是慢而不超时比如一直挂着不返回许可就一直被占。所以下游调用必须设置超时webClient.post() .uri(url) .bodyValue(req) .retrieve() .bodyToMono(String.class) .timeout(Duration.ofSeconds(30)) .onErrorResume(e - Mono.just([TIMEOUT]));超时时间要根据下游的实际 P99 来定一般设成 P99 的 2-3 倍。6.3 信号量许可数与下游容量的匹配信号量许可数不是拍脑袋定的要根据下游的承载能力来算。假设下游大模型服务单实例能承受 50 QPS你有 2 个实例那总容量是 100 QPS。如果单次 Agent 调用平均耗时 8 秒那并发上限就是100 * 8 800。但为了留余量实际设置应该打 7 折也就是 560 左右。这个计算过程必须做否则要么限流太松打挂下游要么限流太紧浪费容量。6.4 多实例部署时的信号量语义单机 Semaphore 只能限制单实例的并发。如果你部署了 3 个实例每个实例 Semaphore 是 200那全局并发就是 600。这时候要么把单实例 Semaphore 调小200/3 ≈ 66要么用分布式限流Redis Lua 脚本。Agent 服务我建议用单机 Semaphore 实例数控制因为分布式限流会引入 Redis 依赖和网络开销对 SSE 这种长连接场景不划算。只要实例数固定单机 Semaphore 的总和就是全局上限。7. 写在最后几个我实际踩过的坑第一个坑是在Flux.defer里做重操作。defer的 Supplier 每次订阅都会执行如果里面做了数据库查询、远程调用那每次订阅都会触发一次性能直接崩。defer里只应该做轻量的状态判断和对象创建。第二个坑是用Flux.create手动发信号时忘记处理 cancel。Flux.create的 sink 如果不监听onCancel客户端断开后你的推送逻辑还在跑资源不释放。正确做法是Flux.create(sink - { sink.onCancel(() - { // 清理资源 connectionSemaphore.release(); }); // 推送逻辑 });第三个坑是在 SSE 里用block()。有人为了图方便在流里调block()拿结果这在 WebFlux 里会直接抛IllegalStateException因为 event loop 线程不允许阻塞。所有阻塞操作必须调度到boundedElastic。第四个坑是忘记设置produces MediaType.TEXT_EVENT_STREAM_VALUE。少了这个Spring 会按普通 JSON 返回SSE 客户端收不到流式数据表现为请求一直挂着不返回。第五个坑是压测工具本身的问题。JMeter 默认的 HTTP 客户端对 SSE 支持不好可能提前关闭连接。建议用curl -N或者专门的 SSE 压测工具确保压测结果反映的是服务端真实表现。这个系列做到这里就收尾了。从 Agent 对话链路到生产部署8 讲下来最大的体会是响应式编程的坑90% 都在求值时机和资源生命周期这两件事上。Semaphore 放错位置、Flux 被缓存复用、doFinally 没执行本质上都是没搞清楚这段代码什么时候执行、执行几次、什么时候清理。把这三个问题想清楚生产环境的并发安全就稳了一大半。
返回列表