Netty 百万长连接推送网关生产实战
引言
企微推送、电商秒杀通知、IoT 指令下发……这些场景都有一个共同挑战:如何在单台机器上维持几十万甚至百万级长连接,并在下游抖动时保证系统不雪崩。
本篇文章基于 Netty 构建一套生产级推送网关,从连接治理到可观测性全链路展开。
| 指标 | 数值 | 说明 |
|---|---|---|
| 单机连接数 | 1.2M+ | 16C32G 云主机,内核参数调优后实测 |
| 消息峰值 QPS | 80W+ | 批量合并 + 零拷贝推送 |
| P99 推送延迟 | < 8ms | 心跳 + 写缓冲区水位控制 |
| GC 停顿 | < 10ms | 对象池 + 堆外内存 + ZGC |
01 架构全景:四层网关模型
生产级推送网关不是简单的 WebSocket Server。它需要处理接入层负载均衡、连接状态管理、消息路由、下游保护四个核心职责。
─────────────────────────────────────────────────────────────────
接入层:iOS/Android/Web/小程序 → HAProxy → Spring Cloud Gateway
─────────────────────────────────────────────────────────────────
连接层:ChannelGroup ↔ UserId→Channel 本地索引 ↔ IdleStateHandler
─────────────────────────────────────────────────────────────────
路由层:Kafka push.topic → Consumer Pool → 本地路由表 / 广播
─────────────────────────────────────────────────────────────────
保护层:写水位 + 令牌桶限流 + Resilience4j 熔断 + Prometheus
─────────────────────────────────────────────────────────────────
关键设计决策:
- 有状态服务:长连接必须落在固定 Netty 节点,L4 负载均衡采用源地址哈希,避免七层再路由。
- 本地索引优先:用户在线状态先查本地
ConcurrentHashMap,未命中再回查 Redis,降低 90% 以上远程调用。 - 广播转局部:全量推送通过 Kafka 分片消费,只推送本节点挂载的连接,避免跨节点 RPC 风暴。
02 连接治理:百万连接的内存与线程模型
Netty 的线程模型是性能基石。生产环境使用 EpollEventLoopGroup(Linux)并设置合理的 SO_BACKLOG、TCP_NODELAY、SO_KEEPALIVE。
public class PushGatewayServer implements Lifecycle {
private final EventLoopGroup bossGroup = new EpollEventLoopGroup(1);
private final EventLoopGroup workerGroup = new EpollEventLoopGroup(
0, new DefaultThreadFactory("netty-worker"));
public void start(int port) throws InterruptedException {
ServerBootstrap b = new ServerBootstrap();
b.group(bossGroup, workerGroup)
.channel(EpollServerSocketChannel.class)
.option(ChannelOption.SO_BACKLOG, 8192)
.option(ChannelOption.SO_REUSEADDR, true)
.childOption(ChannelOption.TCP_NODELAY, true)
.childOption(ChannelOption.SO_KEEPALIVE, true)
.childOption(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT)
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ch.config().setWriteBufferWaterMark(
new WriteBufferWaterMark(32 * 1024, 256 * 1024));
ch.pipeline()
.addLast("idle", new IdleStateHandler(90, 30, 0))
.addLast("codec", new PushProtocolCodec())
.addLast("auth", new AuthHandshakeHandler(jwtVerifier, sessionStore))
.addLast("biz", new PushBusinessHandler(connectionManager, pushRouter));
}
});
b.bind(port).sync();
}
}
连接管理器需要解决三个问题:线程安全、快速查找、优雅下线。采用用户 ID 与 Channel 的多级索引:
public class ConnectionManager {
// userId -> Channel 主索引
private final ConcurrentHashMap<String, Channel> userChannelMap = new ConcurrentHashMap<>();
// ChannelId -> userId 反向索引,用于断线时清理
private final ConcurrentHashMap<String, String> channelUserMap = new ConcurrentHashMap<>();
public void bind(String userId, Channel channel) {
channel.attr(Attributes.USER_ID).set(userId);
Channel prev = userChannelMap.put(userId, channel);
if (prev != null && prev.isActive()) {
// 同一用户新登录,踢掉旧连接
prev.writeAndFlush(new KickoutMessage("new_login"))
.addListener(ChannelFutureListener.CLOSE);
}
channelUserMap.put(channel.id().asShortText(), userId);
Metrics.CONNECTIONS.increment();
}
public void unbind(Channel channel) {
String userId = channel.attr(Attributes.USER_ID).getAndSet(null);
if (userId != null) {
userChannelMap.remove(userId, channel);
channelUserMap.remove(channel.id().asShortText());
Metrics.CONNECTIONS.decrement();
}
}
}
03 背压与限流:防止客户端拖垮整个集群
推送网关最常见的故障模式是:下游某个客户端接收极慢,TCP 发送缓冲区堆积,最终把服务内存撑爆。生产方案需要业务层背压 + 令牌桶限流 + 慢连接熔断三位一体。
public class BackPressurePushHandler extends ChannelOutboundHandlerAdapter {
private final Semaphore globalInflight = new Semaphore(500_000);
private final RateLimiter globalRateLimiter = RateLimiter.create(800_000.0);
@Override
public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) {
Channel ch = ctx.channel();
// 1. 全局 QPS 限流,保证 CPU 不跑满
if (!globalRateLimiter.tryAcquire(1, TimeUnit.MILLISECONDS)) {
Metrics.RATE_LIMITED.increment();
ReferenceCountUtil.release(msg);
promise.setFailure(new PushException("global rate limit"));
return;
}
// 2. 写缓冲区水位背压,单个 channel 排队超过阈值直接丢弃
if (!ch.isWritable()) {
Metrics.BACK_PRESSURE_DROP.increment();
ReferenceCountUtil.release(msg);
promise.setFailure(new PushException("channel not writable"));
return;
}
// 3. 全局在途消息数限制,防止内存无限增长
if (!globalInflight.tryAcquire()) {
Metrics.INFLIGHT_REJECT.increment();
ReferenceCountUtil.release(msg);
promise.setFailure(new PushException("inflight overflow"));
return;
}
ctx.write(msg, promise).addListener(f -> globalInflight.release());
}
}
三个层级的保护:
- 全局 QPS 限流:基于 Guava RateLimiter,保护 Netty Worker 线程不被打满。
- 写缓冲区背压:
isWritable()判断的是 Netty 写水位,低层且高效。 - 全局在途计数:通过 Semaphore 限制未确认的消息总量,避免瞬时洪峰导致 OOM。
04 熔断降级:下游抖动时的自愈机制
推送网关对接的业务系统(如企微回调、订单中心)偶尔会出现超时或错误率飙升。如果网关无脑重试,会把故障放大。使用 Resilience4j 对按用户维度聚合后的批量推送接口做熔断,并配合降级策略。
public class ProtectedPushService {
private final CircuitBreakerRegistry registry;
private final PushMetrics metrics;
public Mono<Void> pushBatch(String bizType, List<PushMessage> messages) {
CircuitBreaker cb = registry.circuitBreaker(bizType, "default");
return Mono.fromCallable(() -> doPushBatch(messages))
.transformDeferred(CircuitBreakerOperator.of(cb))
.doOnSuccess(v -> metrics.recordSuccess(bizType, messages.size()))
.doOnError(e -> metrics.recordFailure(bizType, e.getClass().getSimpleName()))
.onErrorResume(Throwable.class, e -> fallback(bizType, messages, e));
}
private Mono<Void> fallback(String bizType, List<PushMessage> messages, Throwable e) {
if (e instanceof CallNotPermittedException) {
// 熔断开启:写入死信队列,稍后重推
return Mono.fromRunnable(() -> deadLetterQueue.offer(bizType, messages));
}
// 其他异常:按用户维度降级为只推在线用户
List<PushMessage> onlineOnly = messages.stream()
.filter(m -> connectionManager.isOnline(m.getUserId()))
.toList();
return Mono.fromRunnable(() -> doPushBatch(onlineOnly));
}
}
熔断配置核心参数(按业务类型隔离):
| 参数 | 默认值 | 说明 |
|---|---|---|
| failureRateThreshold | 50% | 50% 失败率开启熔断 |
| slowCallRateThreshold | 80% | 慢调用比例阈值 |
| slowCallDurationThreshold | 500ms | 超过即视为慢调用 |
| waitDurationInOpenState | 20s | 熔断后等待半开时间 |
| permittedNumberOfCallsInHalfOpenState | 10 | 半开探针数量 |
05 上下文传播与可观测性:定位线上问题不抓瞎
长连接服务的问题定位非常困难:一条消息可能经过 Kafka、Netty、业务 Handler 多个线程。要求每个阶段都携带 TraceId,并通过 Micrometer 暴露连接数、推送 QPS、 延迟、错误率等核心指标。
public class TraceContextHandler extends ChannelDuplexHandler {
private static final AttributeKey<String> TRACE_ID =
AttributeKey.valueOf("traceId");
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
if (msg instanceof PushPacket packet) {
String traceId = packet.getTraceId() != null
? packet.getTraceId()
: TraceIdGenerator.next();
ctx.channel().attr(TRACE_ID).set(traceId);
try (MDC.MDCCloseable ignored = MDC.putCloseable("traceId", traceId)) {
ctx.fireChannelRead(packet);
}
} else {
ctx.fireChannelRead(msg);
}
}
@Override
public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) {
if (msg instanceof PushPacket packet) {
String traceId = ctx.channel().attr(TRACE_ID).get();
if (traceId != null) packet.setTraceId(traceId);
}
ctx.write(msg, promise);
}
}
指标埋点 RED 四类:
| 类型 | 指标名 | 说明 |
|---|---|---|
| Rate | push_messages_total | 按 bizType / status 标签聚合 |
| Errors | push_errors_total | 区分 timeout / backpressure / circuit_open |
| Duration | push_latency_seconds | P50 / P99 / P999 直方图 |
| Saturation | netty_connections | Gauge 实时连接数与水位 |
06 生产踩坑:真金白银买来的经验
坑 1:Epoll 不可用却未兜底,连接数上不去
部分容器镜像缺少 native epoll 库,Netty 会静默回退到 NIO,但性能直接腰斩。
修复:启动时检测 Epoll.isAvailable(),不可用时告警;同时用 -Dio.netty.noUnsafe=false 开启堆外内存。
坑 2:只读空闲不检测,"僵尸连接"耗尽文件句柄
客户端断网不会立即触发 TCP FIN,导致服务端维持大量死连接。
修复:IdleStateHandler 必须同时配置读/写空闲,读空闲超 90s 强制关闭,并配合应用层心跳确认。
坑 3:ByteBuf 引用计数泄漏,凌晨 OOM
自定义 Handler 中忘记 release() 或重复释放都会触发内存泄漏。
修复:启用 ResourceLeakDetector.Level.PARANOID 在测试环境抓泄漏;生产使用 SimpleChannelInboundHandler 自动释放。
坑 4:发布时直接 kill -9,消息丢失 + 连接雪崩
滚动发布时粗暴退出,未写出的消息和内存队列全部丢失。
修复:注册 JVM ShutdownHook,先标记节点为 offline、停止接收新连接、等待 30s 让在途消息 flush,再优雅关闭 EventLoop。
坑 5:全量广播没有做分片,瞬间打满内网带宽
百万用户同时推送时,如果不做本地过滤,所有节点会互相同步用户在线状态。
修复:Kafka 按 userId 取模路由到 Partition,消费者只推送本节点持有的连接,实现"本地广播"。
07 总结
百万长连接推送网关的核心 checklist:
- 连接治理:用户-Channel 双向索引 + 单点登录踢人 + 优雅下线。
- 背压限流:全局 QPS 限流 + 写水位 + 在途消息数三重保护。
- 熔断降级:按业务类型隔离熔断,失败消息入死信队列。
- 可观测性:TraceId 全链路传递 + RED 指标 + 慢连接/死连接监控。
- 内核调优:ulimit、tcp_keepalive、epoll、零拷贝、对象池缺一不可。
Netty 本身只是工具,真正决定上线稳定性的是对边界条件的敬畏:慢客户端、断网、发布、广播、内存泄漏,每一项都可能在凌晨把你叫起来。希望这篇实战能帮你少踩几个坑。
更多推荐




所有评论(0)