Netty 百万长连接推送网关生产实战

📅 2026/7/30 10:15:38 👁️ 阅读次数 📝 编程学习
Netty 百万长连接推送网关生产实战

引言

企微推送、电商秒杀通知、IoT 指令下发……这些场景都有一个共同挑战:如何在单台机器上维持几十万甚至百万级长连接,并在下游抖动时保证系统不雪崩。

本篇文章基于 Netty 构建一套生产级推送网关,从连接治理到可观测性全链路展开。

指标数值说明
单机连接数1.2M+16C32G 云主机,内核参数调优后实测
消息峰值 QPS80W+批量合并 + 零拷贝推送
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_BACKLOGTCP_NODELAYSO_KEEPALIVE

publicclassPushGatewayServerimplementsLifecycle{privatefinalEventLoopGroupbossGroup=newEpollEventLoopGroup(1);privatefinalEventLoopGroupworkerGroup=newEpollEventLoopGroup(0,newDefaultThreadFactory("netty-worker"));publicvoidstart(intport)throwsInterruptedException{ServerBootstrapb=newServerBootstrap();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(newChannelInitializer<SocketChannel>(){@OverrideprotectedvoidinitChannel(SocketChannelch){ch.config().setWriteBufferWaterMark(newWriteBufferWaterMark(32*1024,256*1024));ch.pipeline().addLast("idle",newIdleStateHandler(90,30,0)).addLast("codec",newPushProtocolCodec()).addLast("auth",newAuthHandshakeHandler(jwtVerifier,sessionStore)).addLast("biz",newPushBusinessHandler(connectionManager,pushRouter));}});b.bind(port).sync();}}

连接管理器需要解决三个问题:线程安全、快速查找、优雅下线。采用用户 ID 与 Channel 的多级索引:

publicclassConnectionManager{// userId -> Channel 主索引privatefinalConcurrentHashMap<String,Channel>userChannelMap=newConcurrentHashMap<>();// ChannelId -> userId 反向索引,用于断线时清理privatefinalConcurrentHashMap<String,String>channelUserMap=newConcurrentHashMap<>();publicvoidbind(StringuserId,Channelchannel){channel.attr(Attributes.USER_ID).set(userId);Channelprev=userChannelMap.put(userId,channel);if(prev!=null&&prev.isActive()){// 同一用户新登录,踢掉旧连接prev.writeAndFlush(newKickoutMessage("new_login")).addListener(ChannelFutureListener.CLOSE);}channelUserMap.put(channel.id().asShortText(),userId);Metrics.CONNECTIONS.increment();}publicvoidunbind(Channelchannel){StringuserId=channel.attr(Attributes.USER_ID).getAndSet(null);if(userId!=null){userChannelMap.remove(userId,channel);channelUserMap.remove(channel.id().asShortText());Metrics.CONNECTIONS.decrement();}}}

03 背压与限流:防止客户端拖垮整个集群

推送网关最常见的故障模式是:下游某个客户端接收极慢,TCP 发送缓冲区堆积,最终把服务内存撑爆。生产方案需要业务层背压 + 令牌桶限流 + 慢连接熔断三位一体。

publicclassBackPressurePushHandlerextendsChannelOutboundHandlerAdapter{privatefinalSemaphoreglobalInflight=newSemaphore(500_000);privatefinalRateLimiterglobalRateLimiter=RateLimiter.create(800_000.0);@Overridepublicvoidwrite(ChannelHandlerContextctx,Objectmsg,ChannelPromisepromise){Channelch=ctx.channel();// 1. 全局 QPS 限流,保证 CPU 不跑满if(!globalRateLimiter.tryAcquire(1,TimeUnit.MILLISECONDS)){Metrics.RATE_LIMITED.increment();ReferenceCountUtil.release(msg);promise.setFailure(newPushException("global rate limit"));return;}// 2. 写缓冲区水位背压,单个 channel 排队超过阈值直接丢弃if(!ch.isWritable()){Metrics.BACK_PRESSURE_DROP.increment();ReferenceCountUtil.release(msg);promise.setFailure(newPushException("channel not writable"));return;}// 3. 全局在途消息数限制,防止内存无限增长if(!globalInflight.tryAcquire()){Metrics.INFLIGHT_REJECT.increment();ReferenceCountUtil.release(msg);promise.setFailure(newPushException("inflight overflow"));return;}ctx.write(msg,promise).addListener(f->globalInflight.release());}}

三个层级的保护:

  1. 全局 QPS 限流:基于 Guava RateLimiter,保护 Netty Worker 线程不被打满。
  2. 写缓冲区背压isWritable()判断的是 Netty 写水位,低层且高效。
  3. 全局在途计数:通过 Semaphore 限制未确认的消息总量,避免瞬时洪峰导致 OOM。

04 熔断降级:下游抖动时的自愈机制

推送网关对接的业务系统(如企微回调、订单中心)偶尔会出现超时或错误率飙升。如果网关无脑重试,会把故障放大。使用 Resilience4j 对按用户维度聚合后的批量推送接口做熔断,并配合降级策略。

publicclassProtectedPushService{privatefinalCircuitBreakerRegistryregistry;privatefinalPushMetricsmetrics;publicMono<Void>pushBatch(StringbizType,List<PushMessage>messages){CircuitBreakercb=registry.circuitBreaker(bizType,"default");returnMono.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));}privateMono<Void>fallback(StringbizType,List<PushMessage>messages,Throwablee){if(einstanceofCallNotPermittedException){// 熔断开启:写入死信队列,稍后重推returnMono.fromRunnable(()->deadLetterQueue.offer(bizType,messages));}// 其他异常:按用户维度降级为只推在线用户List<PushMessage>onlineOnly=messages.stream().filter(m->connectionManager.isOnline(m.getUserId())).toList();returnMono.fromRunnable(()->doPushBatch(onlineOnly));}}

熔断配置核心参数(按业务类型隔离):

参数默认值说明
failureRateThreshold50%50% 失败率开启熔断
slowCallRateThreshold80%慢调用比例阈值
slowCallDurationThreshold500ms超过即视为慢调用
waitDurationInOpenState20s熔断后等待半开时间
permittedNumberOfCallsInHalfOpenState10半开探针数量

05 上下文传播与可观测性:定位线上问题不抓瞎

长连接服务的问题定位非常困难:一条消息可能经过 Kafka、Netty、业务 Handler 多个线程。要求每个阶段都携带TraceId,并通过 Micrometer 暴露连接数、推送 QPS、 延迟、错误率等核心指标。

publicclassTraceContextHandlerextendsChannelDuplexHandler{privatestaticfinalAttributeKey<String>TRACE_ID=AttributeKey.valueOf("traceId");@OverridepublicvoidchannelRead(ChannelHandlerContextctx,Objectmsg){if(msginstanceofPushPacketpacket){StringtraceId=packet.getTraceId()!=null?packet.getTraceId():TraceIdGenerator.next();ctx.channel().attr(TRACE_ID).set(traceId);try(MDC.MDCCloseableignored=MDC.putCloseable("traceId",traceId)){ctx.fireChannelRead(packet);}}else{ctx.fireChannelRead(msg);}}@Overridepublicvoidwrite(ChannelHandlerContextctx,Objectmsg,ChannelPromisepromise){if(msginstanceofPushPacketpacket){StringtraceId=ctx.channel().attr(TRACE_ID).get();if(traceId!=null)packet.setTraceId(traceId);}ctx.write(msg,promise);}}

指标埋点 RED 四类:

类型指标名说明
Ratepush_messages_total按 bizType / status 标签聚合
Errorspush_errors_total区分 timeout / backpressure / circuit_open
Durationpush_latency_secondsP50 / P99 / P999 直方图
Saturationnetty_connectionsGauge 实时连接数与水位

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:

  1. 连接治理:用户-Channel 双向索引 + 单点登录踢人 + 优雅下线。
  2. 背压限流:全局 QPS 限流 + 写水位 + 在途消息数三重保护。
  3. 熔断降级:按业务类型隔离熔断,失败消息入死信队列。
  4. 可观测性:TraceId 全链路传递 + RED 指标 + 慢连接/死连接监控。
  5. 内核调优:ulimit、tcp_keepalive、epoll、零拷贝、对象池缺一不可。

Netty 本身只是工具,真正决定上线稳定性的是对边界条件的敬畏:慢客户端、断网、发布、广播、内存泄漏,每一项都可能在凌晨把你叫起来。希望这篇实战能帮你少踩几个坑。