引言

企微推送、电商秒杀通知、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_BACKLOGTCP_NODELAYSO_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());
  }
}

三个层级的保护:

  1. 全局 QPS 限流:基于 Guava RateLimiter,保护 Netty Worker 线程不被打满。
  2. 写缓冲区背压isWritable() 判断的是 Netty 写水位,低层且高效。
  3. 全局在途计数:通过 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:

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

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


Logo

电商企业物流数字化转型必备!快递鸟 API 接口,72 小时快速完成物流系统集成。全流程实战1V1指导,营造开放的API技术生态圈。

更多推荐