本文是「成都硅基边界」零代码构建平台技术复盘系列,继续 Java 侧业务后端(Silicon-Service)。上一篇聊了 JWT + Redis 双层鉴权,这一篇聊长连接:为什么 WebSocket 一旦上了多实例部署就全是坑,以及我们怎么用"网关粘性路由 + Netty 原生服务"把这两个坑填平。
文中代码均为示意代码(根据设计改写、非项目真实源码),按生产级 Java 规范书写。
一、背景:长连接在分布式下会遇到什么
HTTP 请求是"短平快"的:来一个请求、做完就走,随便打到哪个实例都行——所以负载均衡可以肆无忌惮地轮询。
WebSocket 完全不是这样:连接一旦建立就绑死在某个实例上,后续所有消息都必须由这个实例收发。于是多实例部署下立刻冒出两个难题:
- 难题一:连接该往哪个实例路由? 默认轮询负载均衡会把同一个客户端的重连请求分到不同实例,导致"连接在实例 A、消息发到实例 B"——推送永远到不了。
- 难题二:流式消息为什么会乱序? LLM 应用的回答是流式分片推送的(打字机效果)。多个线程同时往同一个连接写数据,到了 Netty 的 EventLoop 里就可能重排——"结束标记"比最后一段正文先到,前端直接错乱。
我们最终落地的方案是:网关注入粘性路由(解决难题一)+ Netty 原生独立服务(解决难题二)。整体链路如下:
flowchart LR
C[客户端<br/>网页客服 / 管理端对话] -->|ws://.../ws/chat<br/>?token=xxx| GW[网关 :8080]
GW --> SF[WsStickyRoutingFilter<br/>order=10050]
SF -->|按客户端 IP 一致性哈希<br/>替换 lb:// 为具体实例| N1[agent-service A<br/>Netty :27051]
SF -.->|实例下线则重分配| N2[agent-service B<br/>Netty :27051]
N1 --> H[NioWebSocketHandler<br/>连接表 / 账号绑定]
H --> SE[单线程 STREAM_EXECUTOR<br/>保证流式分片保序]
SE -->|TextWebSocketFrame| C
二、难题一:网关粘性路由
2.1 为什么不能用默认负载均衡
Spring Cloud Gateway 默认会把 lb://agent-service/** 交给 ReactiveLoadBalancerClientFilter 做轮询。对 WebSocket 来说这等于"每次都随机换一个实例",连接和后续消息必然对不上。
我们要做的是:在负载均衡之前,把 lb:// 直接替换成具体实例地址,让它跳过默认 LB。
2.2 关键在于过滤器顺序
网关的过滤器靠 Ordered 决定执行顺序,这里存在一个精确的"夹缝":
RouteToRequestUrlFilter order = 10000 ← 把路由解析成 lb://agent-service/ws/chat
↓ ★ 我们的粘性路由插在这里(10050)★
ReactiveLoadBalancerClientFilter order = 10150 ← 默认负载均衡(要被跳过)
只有卡在这两者之间,才能拿到已经解析成 lb:// 的 URL,并且在默认 LB 生效前把它改掉。实现如下:
/**
* WebSocket 粘性路由过滤器:按客户端标识把 WS 请求固定转发到同一实例。
* <p>执行时机在 RouteToRequestUrlFilter(10000) 之后、ReactiveLoadBalancerClientFilter(10150)
* 之前,通过替换 GATEWAY_REQUEST_URL_ATTR 跳过默认负载均衡。
*/
@Component
@RequiredArgsConstructor
@Slf4j
public class WsStickyRoutingFilter implements GlobalFilter, Ordered {
private static final String TARGET_SERVICE = "agent-service";
private static final String WS_PATH_PREFIX = "/ws";
/** agent-service 独立 Netty WebSocket 端口(与 HTTP 端口不同,需替换注册端口) */
private static final int WS_PORT = 27051;
private final ReactiveDiscoveryClient discoveryClient;
/** 客户端标识 -> 实例地址(host:port) 的粘性映射缓存 */
private final Map<String, String> stickyMap = new ConcurrentHashMap<>();
@Override
public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
URI url = exchange.getAttribute(GATEWAY_REQUEST_URL_ATTR);
// 只处理 lb://agent-service/ws/** 这一类请求
if (url == null || !"lb".equals(url.getScheme())
|| !TARGET_SERVICE.equals(url.getHost())
|| !url.getPath().startsWith(WS_PATH_PREFIX)) {
return chain.filter(exchange);
}
String clientKey = resolveClientKey(exchange);
return discoveryClient.getInstances(TARGET_SERVICE)
.collectList()
.flatMap(instances -> {
if (instances.isEmpty()) {
log.error("粘性路由:无可用实例,回落默认负载均衡");
return chain.filter(exchange);
}
String address = resolveInstance(clientKey, instances);
URI directUri = URI.create(
url.toString().replace("lb://" + TARGET_SERVICE, "http://" + address));
exchange.getAttributes().put(GATEWAY_REQUEST_URL_ATTR, directUri);
return chain.filter(exchange);
});
}
/**
* 解析目标实例:优先复用缓存,缓存实例已下线则重新分配。
*/
private String resolveInstance(String clientKey, List<ServiceInstance> instances) {
String cached = stickyMap.get(clientKey);
if (cached != null) {
boolean alive = instances.stream()
.anyMatch(i -> (i.getHost() + ":" + WS_PORT).equals(cached));
if (alive) {
return cached;
}
stickyMap.remove(clientKey); // 实例下线,重新分配
log.info("粘性路由:实例已下线,重新分配 client={} old={}", clientKey, cached);
}
ServiceInstance selected = instances.get(Math.abs(clientKey.hashCode()) % instances.size());
String address = selected.getHost() + ":" + WS_PORT;
stickyMap.put(clientKey, address);
log.info("粘性路由:新分配 client={} -> {}", clientKey, address);
return address;
}
/** 取客户端标识:优先代理头,逐级回落到真实远端地址 */
private String resolveClientKey(ServerWebExchange exchange) {
HttpHeaders headers = exchange.getRequest().getHeaders();
String forwarded = headers.getFirst("X-Forwarded-For");
if (StringUtils.isNotBlank(forwarded)) {
return forwarded.split(",")[0].trim();
}
String realIp = headers.getFirst("X-Real-IP");
if (StringUtils.isNotBlank(realIp)) {
return realIp;
}
InetSocketAddress remote = exchange.getRequest().getRemoteAddress();
return remote != null ? remote.getAddress().getHostAddress() : "unknown";
}
@Override
public int getOrder() {
return 10050; // 必须介于 10000 与 10150 之间
}
}
2.3 两个必须知道的局限
hashCode % 实例数在扩缩容时会大面积重映射:实例数从 3 变 4,几乎所有客户端的映射都会变,老实例上的连接被迫全部重连。规模上来后应换成一致性哈希环 + 虚拟节点,把重映射控制在 1/N;- 同 NAT 出口的多用户会挤到同一实例:公司内网几百个用户可能共用一个出口 IP,按 IP 哈希会把它们全压到一台机器上。更稳的做法是用"用户/会话 ID"而不是 IP 做键(鉴权后即可拿到)。
三、难题二:流式消息为什么必须保序
3.1 根因在 EventLoop 的 MPSC 队列
Netty 的 Channel.writeAndFlush() 是线程安全的——你可以从任意线程调用。但"线程安全"不等于"全序保证":多个线程并发提交写任务时,任务进入 EventLoop 的 MPSC(多生产者单消费者)队列,入队顺序取决于线程调度的先后,并不能保证与你业务代码的调用顺序一致。
对流式回答来说这是致命的:
业务线程顺序: 分片1 → 分片2 → 分片3 → 结束(end=true)
EventLoop 实际:分片1 → 分片3 → 结束 → 分片2
前端看到: "你好," + "世界" + [结束] + ",欢迎"
"结束标记"提前到达,前端会把回答提前截断,用户看到半句话就停了。
3.2 解法:单线程串行化写入
我们给流式消息单独开了一个单线程执行器,所有流式分片都提交到这一个线程里按序发送:
/**
* 流式消息专用单线程执行器。
* <p>保证 streamMsg 严格按调用顺序串行写入 Netty Channel,
* 避免多线程 writeAndFlush 提交到 EventLoop MPSC 队列时发生乱序
* (表现为 end=true 的分片先于正文分片到达前端,前端提前终止渲染)。
*/
private static final ExecutorService STREAM_EXECUTOR = Executors.newSingleThreadExecutor(r -> {
Thread t = new Thread(r, "ws-stream-sender");
t.setDaemon(true); // 守护线程,不阻塞应用退出
return t;
});
/**
* 推送流式分片(按序)。
*
* @param customerId 目标客户(可能对应多个连接:多标签页)
* @param msgId 流式消息 ID(前端据此归并同一段回答)
* @param text 本片文本
* @param end 是否为最后一片
*/
public static boolean pushStreamChunk(Long customerId, String msgId, String text, boolean end) {
try {
STREAM_EXECUTOR.execute(() -> {
StreamMessage payload = new StreamMessage(msgId, text, end);
broadcastToCustomer(customerId, ServerMessage.of(MessageType.streamMsg, payload));
});
return true;
} catch (Exception e) {
log.error("流式消息推送失败 customerId={} msgId={}", customerId, msgId, e);
return false;
}
}
这个设计的本质是"用一次线程切换换取顺序保证"。有人会担心单线程成为瓶颈——对文本流式分片来说完全不必:分片体积小(几十字节到几 KB),发送是内存写 + 网络缓冲,单线程每秒可处理数万条;而顺序正确性对此场景是硬需求。如果未来单线程真的成为瓶颈,正确方向是"每连接一个有序队列"而不是"多线程共享连接"。
四、Netty 服务搭建的四个工程细节
4.1 异步启动,不拖慢应用启动
Netty 服务放在独立的初始化组件里,通过线程池异步启动,避免 bind() 阻塞 Spring 启动流程:
@Component
@RequiredArgsConstructor
@Slf4j
public class WebSocketServerBootstrap implements InitializingBean {
private static final int WS_PORT = 27051;
private final WebSocketChannelInitializer channelInitializer;
@Resource(name = "taskExecutor")
private ThreadPoolTaskExecutor taskExecutor;
@Override
public void afterPropertiesSet() {
taskExecutor.execute(this::start); // 异步启动,不阻塞 Spring 容器
}
private void start() {
// bossGroup 只负责接受连接,1 个线程足够;workerGroup 负责 IO 读写,默认 CPU * 2
EventLoopGroup bossGroup = new NioEventLoopGroup(1);
EventLoopGroup workerGroup = new NioEventLoopGroup();
try {
ServerBootstrap bootstrap = new ServerBootstrap()
.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.handler(new LoggingHandler(LogLevel.INFO)) // 入/出站事件日志,便于排查
.childHandler(channelInitializer);
Channel channel = bootstrap.bind(WS_PORT).sync().channel();
log.info("WebSocket 服务启动成功,端口={}", WS_PORT);
channel.closeFuture().sync();
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // 保留中断标记
log.error("WebSocket 启动被中断", e);
} catch (Exception e) {
log.error("WebSocket 启动异常", e);
} finally {
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully();
log.info("WebSocket 服务已关闭");
}
}
}
4.2 独立端口是刻意的取舍
WebSocket 跑在 27051,与服务的 HTTP 端口完全分离。好处是互不干扰(HTTP 接口扩容不影响长连接)、可以独立调参;代价是必须处理端口映射与地址发现——这正是上一节粘性路由里要把注册端口替换成 WS_PORT 的原因。这是"独立部署"必然付出的复杂度,属于知情选择。
4.3 Pipeline 顺序不能错
@Override
protected void initChannel(SocketChannel ch) {
ch.pipeline()
.addLast(new HttpServerCodec()) // HTTP 编解码(WS 握手基于 HTTP)
.addLast(new ChunkedWriteHandler()) // 支持大块数据分片写
.addLast(new HttpObjectAggregator(65536)) // 聚合 HTTP 分片为完整请求(握手需要)
.addLast(new WebSocketServerProtocolHandler("/ws/chat", true)) // 完成 101 协议升级
.addLast(webSocketHandler); // 业务处理器
}
顺序的逻辑链条很清楚:WebSocket 握手本质是一次 HTTP 请求——先用 HTTP 编解码器,再聚合分片请求,才轮到协议升级处理器,最后才是业务 Handler。WebSocketServerProtocolHandler 的第二个参数是"是否处理 ping/pong 帧",交给它处理可以省掉手写心跳应答。
4.4 单例 Handler 与线程安全
业务 Handler 用 @ChannelHandler.Sharable 声明为单例(这样才好统计在线数),配套要求是:所有共享状态必须线程安全,且不能用 @ChannelHandler 的实例字段存连接状态。
@Component
@ChannelHandler.Sharable // 单例:便于全局在线统计
public class NioWebSocketHandler extends SimpleChannelInboundHandler<TextWebSocketFrame> {
/** 连接ID -> Channel */
private static final Map<String, Channel> CHANNEL_MAP = new ConcurrentHashMap<>();
/** 连接ID -> 绑定的账号/客户ID */
private static final Map<String, Long> PRINCIPAL_MAP = new ConcurrentHashMap<>();
/** 用 @Lazy 打破"启动期 Handler 依赖 Service、Service 又依赖容器"的循环 */
@Lazy @Resource private ChannelConfigService channelConfigService;
...
}
五、连接生命周期:绑定、鉴权与清理
5.1 握手期完成鉴权与绑定
长连接只在握手时有一次"带凭证"的机会(业务消息里不该再传 token),所以鉴权必须在这里做完——正好复用上一篇的 JWT + Redis 双层鉴权:
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
if (evt instanceof WebSocketServerProtocolHandler.HandshakeComplete handshake) {
MultiValueMap<String, String> params =
UriComponentsBuilder.fromUriString(handshake.requestUri()).build().getQueryParams();
String token = params.getFirst("token");
if (StringUtils.isNotBlank(token)) {
// 异步线程没有请求上下文,把鉴权头显式放进线程上下文,供 Feign 透传
Map<String, String> headers = Map.of(
HttpHeaders.AUTHORIZATION, URLUtil.decode(token),
PLATFORM_HEADER, PlatformEnum.TENANT.name());
ThreadHeaderHolder.setHeaders(headers);
try {
Result<LoginVO> result = authFeignClient.verifyToken(); // 走 JWT + Redis 双层校验
if (result != null && result.getSuccess()) {
PRINCIPAL_MAP.put(ctx.channel().id().asLongText(),
result.getData().getAccount().getId());
}
} finally {
ThreadHeaderHolder.clear(); // 用完即清,避免污染 Netty 线程
}
}
}
ctx.fireUserEventTriggered(evt);
}
两个要点:① 未通过鉴权的连接不绑定身份,后续业务消息会被直接丢弃;② ThreadHeaderHolder 用完立刻清理——Netty 的 EventLoop 线程同样是长期复用的,ThreadLocal 残留会导致下一个连接读到别人的身份(与上一篇服务层拦截器的"双向清理"是同一类风险)。
5.2 防伪校验:别信客户端报的账号
客户端消息里带 accountId,但绝不能直接采信。对需要身份一致的通道,必须与握手期绑定的身份比对:
@Override
protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame frame) throws Exception {
String key = ctx.channel().id().asLongText();
ClientMessage message = objectMapper.readValue(frame.text(), ClientMessage.class);
message.setChannelKey(key);
// 管理端对话通道:消息里的 accountId 必须与连接绑定的身份一致,否则丢弃
if (ChannelType.chat == message.getChannel()
&& !Objects.equals(PRINCIPAL_MAP.get(key), message.getAccountId())) {
log.warn("身份不匹配,丢弃消息 key={} claimed={}", key, message.getAccountId());
return;
}
switch (message.getType()) {
case create -> handleCreate(ctx, message);
case heartbeat -> sendHeartbeatAck(ctx, message);
case msg -> handleChatMessage(message);
case clear -> handleClear(message);
default -> log.warn("未处理的消息类型: {}", message.getType());
}
}
5.3 三处清理,一处都不能漏
长连接最容易出的生产事故是内存泄漏——连接断了但映射表里的记录没删,跑几天内存就满了。所以"加入 / 移除 / 异常"三个入口都要成对清理:
@Override
public void handlerAdded(ChannelHandlerContext ctx) {
CHANNEL_MAP.put(ctx.channel().id().asLongText(), ctx.channel());
log.info("客户端接入 {}", ctx.channel().remoteAddress());
}
@Override
public void handlerRemoved(ChannelHandlerContext ctx) {
String key = ctx.channel().id().asLongText();
ctx.close();
CHANNEL_MAP.remove(key); // 两张表必须同步清理
PRINCIPAL_MAP.remove(key);
log.info("客户端断开 {}", ctx.channel().remoteAddress());
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
String key = ctx.channel().id().asLongText();
log.error("连接异常 {}", ctx.channel().remoteAddress(), cause);
CHANNEL_MAP.remove(key); // 异常路径也要清,否则一定泄漏
PRINCIPAL_MAP.remove(key);
ctx.close();
}
另外注意 @ChannelHandler.Sharable 单例是全局共享的,所以 handlerAdded / handlerRemoved 会并发触发——映射表必须用 ConcurrentHashMap(这点我们踩过:早期用 HashMap 曾出现并发写导致条目丢失)。
六、消息协议:一个枚举搞定双向通信
客户端与服务端用同一套消息信封,靠 type 分发:
/** 双向消息类型 */
public enum MessageType {
create, // 建立会话(携带客户端配置)
heartbeat, // 心跳保活 / 心跳应答
msg, // 普通消息(含用户发送与机器人回复)
historyList, // 历史消息列表
streamMsg, // 流式分片(打字机效果)
conversation, // 会话列表变更(管理端)
clear, // 清除会话
}
/** 客户端 -> 服务端 */
@Data
public class ClientMessage {
private String channelKey; // 连接标识(服务端回填,客户端无需传)
private ChannelType channel; // web=网页客服,chat=管理端对话
private MessageType type;
private Long accountId;
private JsonNode data; // 按 type 多态解析
}
设计要点:data 用 JsonNode 承载多态载荷,按 type 再转换成具体 DTO——避免为一个枚举值建一个大而全的字段集合;channelKey 由服务端回填,客户端伪造无效(服务端以连接 ID 为准)。
七、踩坑与经验清单
- 粘性路由必须卡在 10000~10150 之间:早于
RouteToRequestUrlFilter拿不到lb://URL,晚于ReactiveLoadBalancerClientFilter就已经被负载均衡改写过了; - WS 独立端口要同步改地址发现:服务注册的是 HTTP 端口,粘性路由必须替换成
WS_PORT,否则转发到 HTTP 端口上握手必失败; - 多线程写同一连接必然可能乱序:只要有"流式 + 顺序敏感"的组合,就必须串行化写入(单线程执行器是最小代价解法);
- ThreadLocal 在 EventLoop 线程上同样会残留:Netty 线程也是长期复用的,用完必须
clear(); - 客户端上报的身份一律不可信:必须与握手期绑定身份比对(越权风险);
- 清理要覆盖三条路径:连接加入、正常断开、异常断开,两处映射表同步清理;
- @Sharable 单例下所有静态状态都要线程安全:容器必须是并发容器;
- 异步启动 + 优雅关闭成对出现:启动不阻塞 Spring,关闭时
shutdownGracefully()释放线程组。
八、结语
WebSocket 在多实例部署下的复杂度,本质来自两个"反直觉":连接是有状态的(所以不能随便负载均衡)、写入是并发安全的但不保证全序(所以流式必须串行化)。
我们的应对可以浓缩成两句话:
网关负责"把连接粘对地方"——按客户端标识做粘性路由,跳过默认负载均衡;
服务负责"把消息按序送达"——Netty 原生独立服务 + 单线程写入队列保序。
再叠加握手期鉴权、身份防伪、三路清理这三件事,长连接才真正能扛住生产环境的并发与长时间运行。