最近在优化我们团队的智能客服助手时遇到了一个典型的高并发瓶颈当在线用户数激增TPS每秒事务处理数超过500时系统响应变得极其缓慢用户体验直线下降。经过一番深入的排查和重构我们成功将平均响应时间从令人抓狂的2秒优化到了流畅的200毫秒以内。这个过程就像给一辆老爷车换上了赛车引擎感触颇多特地记录下来希望能给遇到类似问题的朋友一些参考。1. 问题定位传统轮询架构的“阿喀琉斯之踵”我们的旧系统采用的是最经典的HTTP短轮询Polling架构。客服端每隔几秒就向服务器发起一次请求询问是否有新的用户消息或状态更新。在用户量少的时候这没什么问题。但当并发上来后问题就暴露无遗。在模拟500 TPS的压测下我们观察到了以下现象CPU占用率居高不下大量的请求创建和销毁带来了巨大的上下文切换开销CPU使用率长期维持在80%以上其中大部分时间花在了处理网络I/O和线程调度上。响应时间P95/P99飙升平均响应时间虽然还能看但P9595%的请求响应时间超过了1.5秒P99更是达到了3秒以上。这意味着每100个请求中就有1个用户需要等待超过3秒这完全不可接受。数据库连接池被打满每个请求几乎都要查询数据库获取会话状态频繁的短连接导致数据库连接池迅速耗尽形成恶性循环。问题的核心在于轮询是一种“拉Pull”模型无论是否有数据更新客户端都会频繁发起请求造成了大量的无效查询和资源浪费。2. 技术选型长轮询、WebSocket与SSE的“三国演义”要解决实时性问题我们必须从“拉”模型转向“推Push”模型或更高效的“拉”模型。我们重点对比了三种主流方案特性维度短轮询 (Polling)长轮询 (Long-Polling)WebSocket服务端事件 (Server-Sent Events, SSE)通信协议HTTPHTTPWebSocket (基于TCP)HTTP通信方向客户端轮询客户端发起服务端挂起至有数据全双工服务端向客户端单向推送协议开销高频繁建立/断开HTTP中连接挂起时间较长低一次握手持久连接低持久连接纯文本协议断连处理简单下次轮询重试复杂需处理超时和重连复杂需心跳保活和重连机制简单浏览器自动重连数据格式任意任意二进制或文本仅文本UTF-8浏览器兼容完美完美良好IE10良好除IE服务端压力极高高低连接数即用户数低我们的选择 对于智能客服场景消息主要是服务端向客服坐席推送如用户新消息、会话转移通知。同时我们需要考虑开发复杂度和运维成本。WebSocket功能最强大但需要额外的协议管理和心跳维护。SSE是纯HTTP协议天然支持断线重连开发更简单且完全满足我们“服务端推送”的核心需求。因此我们最终选择了SSE作为主要的推送通道辅以WebSocket用于需要双向实时交互的特定功能如协同编辑。3. 核心重构基于Spring Reactor的异步非阻塞架构选定SSE后我们利用Spring WebFlux基于Project Reactor重构了消息推送网关实现了真正的异步非阻塞。3.1 带背压控制的请求处理流水线背压Backpressure是响应式编程中流量控制的核心概念防止快速的生产者压垮慢速的消费者。以下是核心处理链Service public class MessagePushService { private final Sinks.ManyServerSentEvent messageSink Sinks.many().multicast().onBackpressureBuffer(1000); // 客服坐席连接端点 GetMapping(path “/stream” produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxServerSentEvent streamMessages(RequestParam String agentId) { return messageSink.asFlux() .filter(event - event.agentId().equals(agentId)) // 1. 按坐席ID过滤 .doOnSubscribe(sub - log.info(“Agent {} subscribed” agentId)) .onBackpressureBuffer(50 // 2. 设置背压缓冲区 BufferOverflowStrategy.DROP_OLDEST) .delayElements(Duration.ofMillis(10) Schedulers.boundedElastic()) // 3. 平滑推送控制速率 .doOnError(ex - log.error(“Stream error for agent {}” agentId ex)); } // 内部服务调用推送消息 public void pushMessageToAgent(String agentId String event String data) { ServerSentEvent sse ServerSentEvent.builder(data) .event(event) .id(UUID.randomUUID().toString()) .build(); // 4. 尝试发射失败则记录如缓冲区满 Sinks.EmitResult result messageSink.tryEmitNext(sse); if (result.isFailure()) { log.warn(“Message dropped for agent {} due to {}” agentId result); // 可在此处将消息转入持久化队列等待后续重试 } } }3.2 分布式会话状态的Redis分片策略客服会话状态如当前服务的用户、会话上下文需要被集群中的所有节点访问。我们使用Redis进行存储并设计了分片策略以提升性能和容量。分片键使用会话ID进行分片确保同一个会话的所有操作都落在同一个Redis实例上避免分布式事务。数据结构使用Hash存储会话详情Set存储坐席的当前会话列表String存储一些快速访问的计数器。过期与持久化设置合理的TTL如30分钟配合RDBAOF持久化防止内存泄漏和数据丢失。# application.yml 片段 spring: redis: cluster: nodes: redis-node1:6379 redis-node2:6379 redis-node3:6379 max-redirects: 3 lettuce: pool: max-active: 16 max-idle: 83.3 熔断降级与超时重试微服务间调用引入Resilience4j进行容错处理防止某个依赖服务故障导致推送服务雪崩。Configuration public class CircuitBreakerConfig { Bean public CircuitBreakerConfigCustomizer customerServiceBreaker() { return CircuitBreakerConfigCustomizer.of(“customerService” builder - builder .slidingWindowSize(10) // 基于最近10次调用计算失败率 .failureRateThreshold(50) // 失败率阈值50% .waitDurationInOpenState(Duration.ofSeconds(30)) // 熔断开启30秒后进入半开 .permittedNumberOfCallsInHalfOpenState(3) // 半开状态下允许的调用数 ); } } Service public class CustomerServiceClient { CircuitBreaker(name “customerService” fallbackMethod “getCustomerInfoFallback”) Retry(name “customerService” fallbackMethod “getCustomerInfoFallback”) TimeLimiter(name “customerService”) public CompletableFutureCustomerInfo getCustomerInfoAsync(String customerId) { // 异步调用客户服务 return webClient.get().uri(“/customers/{id}” customerId).retrieve().bodyToMono(CustomerInfo.class).toFuture(); } private CompletableFutureCustomerInfo getCustomerInfoFallback(String customerId Exception ex) { log.warn(“Fallback triggered for customer {}” customerId ex); // 返回兜底数据如缓存中的旧数据或默认信息 return CompletableFuture.completedFuture(CustomerInfo.defaultInfo(customerId)); } }4. 性能测试数据说话效果显著我们使用JMeter对新架构进行了压测模拟1000个坐席同时在线消息产生速率稳定在600 TPS。压测报告关键指标对比优化后平均响应时间从 ~2000ms 降至~180msP99响应时间从 3000ms 降至~350ms错误率从 5% (主要是超时) 降至 0.1%系统吞吐量提升超过300%服务器资源CPU使用率从80%降至40%左右内存使用更加平稳。不同云实例规格的性价比对比 我们测试了两种常见的云服务器规格2核4G通用型优化前只能勉强支撑200 TPS优化后能稳定处理600 TPS性价比提升约2倍。4核8G计算优化型优化前处理500 TPS就出现瓶颈优化后能轻松应对1200 TPS资源利用率提升显著。结论是架构优化带来的性能收益远大于单纯升级硬件配置。5. 避坑指南那些年我们踩过的“坑”5.1 消息幂等处理的常见误区SSE或WebSocket都可能因为网络抖动导致客户端重复收到消息。实现幂等性不能只依赖前端。误区仅在客户端用lastEventId去重。正确做法服务端在发送消息时生成全局唯一的ID如Snowflake ID并将其作为SSE事件的id字段。关键业务处理如更改会话状态时需要先检查该ID是否已处理过可利用Redis的SETNX命令实现简易的分布式锁或幂等校验。5.2 会话上下文的内存泄漏防范在异步环境中如果将会话上下文如MapString Session缓存在服务内存中很容易因为连接断开后未清理而导致泄漏。解决方案使用WeakHashMap或更好的方式——完全将会话状态外置到Redis中。在Netty或WebFlux中可以利用Channel或WebSession的关闭监听器在连接断开时触发清理逻辑。// 在连接断开时清理资源示例 connection.onClose().subscribe(null null () - { redisTemplate.delete(“session:” sessionId); log.info(“Session {} resources cleaned up” sessionId); });5.3 灰度发布时的流量调度策略直接全量上线新推送网关风险极高。我们结合服务网格如Istio或网关如Spring Cloud Gateway实现了灰度发布。策略基于坐席ID的哈希值或特定Header如X-Version将流量路由到不同版本的推送服务实例。监控密切对比灰度组和基线组的核心指标P99延迟、错误率一旦灰度组出现异常立即切回流量。6. 总结与思考这次架构优化让我们深刻体会到面对高并发实时场景从同步阻塞转向异步非阻塞从轮询转向服务端推送几乎是必然的选择。通过结合Spring Reactor、SSE和Redis我们构建了一个高效、可伸缩的智能客服消息推送系统。最后抛出一个我们在设计后期反复权衡的开放性问题如何平衡实时性与最终一致性在客服系统中坐席看到的消息顺序必须和用户发送的顺序严格一致强一致性但这在分布式、异步推送的场景下成本极高。我们最终采用了“会话内顺序一致性”的折中方案保证单个会话窗口内的消息顺序绝对正确而跨会话或全局的状态如坐席空闲数则接受秒级的延迟最终一致性。这需要在产品体验和技术复杂度之间做出精准的权衡。你的场景中这个平衡点又在哪里呢