实战篇:Netty 通信层完整实现
概述
本篇把前六篇的知识串成一条可运行的完整通信链路:自定义协议编解码 → Session 管理 → 心跳检测 → 消息路由分发。最终产物是一个能接受客户端连接、解析二进制消息、按 Opcode 分发到业务 Handler、并支持心跳保活与断线清理的 Netty 网关。所有代码基于 Java 17 + Netty 4.1,可在 logic 模块直接运行。
一、工程目标与结构
1.1 目标
本篇实现:
自定义协议编解码(Magic + Version + Serializer + Opcode + Length + Body)
Session 管理与连接生命周期
心跳检测(IdleStateHandler 读超时)
消息路由分发(Opcode → Handler)
一条登录消息全链路打通
复用骨架工程的六模块结构:
gateway 模块承载全部通信层代码1.2 代码结构
com.game.gateway
├── codec
│ ├── GameMessage.java // 消息对象
│ ├── GameMessageDecoder.java // 解码器
│ └── GameMessageEncoder.java // 编码器
├── session
│ ├── GameSession.java // 会话
│ └── SessionManager.java // 会话管理
├── dispatch
│ ├── MessageDispatcher.java // 分发器
│ ├── MessageHandler.java // 处理器接口
│ └── MessageMapping.java // 注解
├── handler
│ ├── HeartbeatHandler.java // 心跳
│ └── DispatchHandler.java // 分发入口
├── netty
│ ├── GameServerInitializer.java // Pipeline 装配
│ └── GameNettyServer.java // 服务器启动
└── GameServerApplication.java // 启动入口二、协议层实现
2.1 消息对象
java
package com.game.gateway.codec;
public class GameMessage {
public static final int HEADER_LENGTH = 10;
public static final short MAGIC_NUMBER = (short) 0x4A47;
private byte version;
private byte serializer;
private short opcode;
private byte[] body;
public GameMessage(short opcode, byte[] body) {
this.version = 1;
this.serializer = 0;
this.opcode = opcode;
this.body = body;
}
public short getOpcode() { return opcode; }
public byte[] getBody() { return body; }
}2.2 解码器
java
package com.game.gateway.codec;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.ByteToMessageDecoder;
import java.util.List;
public class GameMessageDecoder extends ByteToMessageDecoder {
private static final int MAX_BODY_LENGTH = 16 * 1024;
@Override
protected void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) {
if (in.readableBytes() < GameMessage.HEADER_LENGTH) {
return;
}
in.markReaderIndex();
short magic = in.readShort();
if (magic != GameMessage.MAGIC_NUMBER) {
ctx.close();
return;
}
byte version = in.readByte();
byte serializer = in.readByte();
short opcode = in.readShort();
int length = in.readInt();
if (length < 0 || length > MAX_BODY_LENGTH) {
ctx.close();
return;
}
if (in.readableBytes() < length) {
in.resetReaderIndex();
return;
}
byte[] body = new byte[length];
in.readBytes(body);
out.add(new GameMessage(version, serializer, opcode, body));
}
}2.3 编码器
java
package com.game.gateway.codec;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.MessageToByteEncoder;
public class GameMessageEncoder extends MessageToByteEncoder<GameMessage> {
@Override
protected void encode(ChannelHandlerContext ctx, GameMessage msg, ByteBuf out) {
out.writeShort(GameMessage.MAGIC_NUMBER);
out.writeByte(msg.getVersion());
out.writeByte(msg.getSerializer());
out.writeShort(msg.getOpcode());
out.writeInt(msg.getBody().length);
out.writeBytes(msg.getBody());
}
}三、Session 层实现
3.1 会话对象
java
package com.game.gateway.session;
import com.game.gateway.codec.GameMessage;
import io.netty.channel.Channel;
public class GameSession {
public static final int STATE_CONNECTED = 0; // 未登录
public static final int STATE_AUTHED = 1; // 已登录
private final Channel channel;
private final long sessionId;
private volatile long playerId;
private volatile int state = STATE_CONNECTED;
private volatile long lastActiveTime;
public GameSession(Channel channel, long sessionId) {
this.channel = channel;
this.sessionId = sessionId;
this.lastActiveTime = System.currentTimeMillis();
}
public void send(GameMessage msg) {
if (channel != null && channel.isActive()) {
channel.writeAndFlush(msg);
}
}
public void close() {
if (channel != null) {
channel.close();
}
}
public long getPlayerId() { return playerId; }
public void setPlayerId(long playerId) { this.playerId = playerId; }
public int getState() { return state; }
public void setState(int state) { this.state = state; }
public long getLastActiveTime() { return lastActiveTime; }
public void touch() { this.lastActiveTime = System.currentTimeMillis(); }
}3.2 会话管理器
java
package com.game.gateway.session;
import com.game.gateway.codec.GameMessage;
import io.netty.channel.Channel;
import io.netty.channel.ChannelId;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
public class SessionManager {
private final ConcurrentHashMap<ChannelId, GameSession> byChannel = new ConcurrentHashMap<>();
private final ConcurrentHashMap<Long, GameSession> byPlayer = new ConcurrentHashMap<>();
private final AtomicLong idGenerator = new AtomicLong();
public GameSession create(Channel channel) {
GameSession session = new GameSession(channel, idGenerator.incrementAndGet());
byChannel.put(channel.id(), session);
return session;
}
public GameSession getByChannel(Channel channel) {
return byChannel.get(channel.id());
}
public void bindPlayer(GameSession session, long playerId) {
// 顶号处理:旧连接下线
GameSession old = byPlayer.put(playerId, session);
if (old != null && old != session) {
old.send(new GameMessage(0, new byte[0])); // 通知下线
old.close();
}
session.setPlayerId(playerId);
}
public void remove(Channel channel) {
GameSession session = byChannel.remove(channel.id());
if (session != null && session.getPlayerId() != 0) {
byPlayer.remove(session.getPlayerId(), session);
}
}
public int onlineCount() {
return byPlayer.size();
}
}四、分发层实现
4.1 处理器接口与注解
java
package com.game.gateway.dispatch;
import com.game.gateway.session.GameSession;
import io.netty.channel.ChannelHandlerContext;
public interface MessageHandler {
void handle(ChannelHandlerContext ctx, GameSession session, byte[] body);
}java
package com.game.gateway.dispatch;
import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
public @interface MessageMapping {
short opcode();
}4.2 分发器
java
package com.game.gateway.dispatch;
import com.game.gateway.codec.GameMessage;
import com.game.gateway.session.GameSession;
import io.netty.channel.ChannelHandlerContext;
import java.util.concurrent.ConcurrentHashMap;
public class MessageDispatcher {
public static final short OP_HEARTBEAT = 1;
public static final short OP_LOGIN = 1001;
private final ConcurrentHashMap<Short, MessageHandler> handlers = new ConcurrentHashMap<>();
public void register(short opcode, MessageHandler handler) {
handlers.put(opcode, handler);
}
public void dispatch(ChannelHandlerContext ctx, GameSession session, GameMessage msg) {
MessageHandler handler = handlers.get(msg.getOpcode());
if (handler == null) {
// 未注册消息:记录日志(心跳等系统消息在网关层拦截,不走这里)
return;
}
handler.handle(ctx, session, msg.getBody());
}
}4.3 业务 Handler 示例
java
package com.game.gateway.handler;
import com.game.gateway.dispatch.MessageHandler;
import com.game.gateway.dispatch.MessageMapping;
import com.game.gateway.session.GameSession;
import io.netty.channel.ChannelHandlerContext;
@MessageMapping(opcode = 1001)
public class LoginHandler implements MessageHandler {
@Override
public void handle(ChannelHandlerContext ctx, GameSession session, byte[] body) {
// 1. 校验登录信息(token 等)
// 2. 加载玩家数据
long playerId = resolvePlayerId(body);
// 3. 绑定会话(顶号处理在 bindPlayer 内)
session.getManager().bindPlayer(session, playerId);
session.setState(GameSession.STATE_AUTHED);
session.touch();
// 4. 回登录成功
session.send(new GameMessage(1002, "{\"code\":0}".getBytes()));
}
private long resolvePlayerId(byte[] body) {
return 10001L; // 实际解析 Protobuf 后返回真实玩家 ID
}
}说明:业务 Handler 由 Spring 扫描 + 注解自动注册
每个 Handler 一个 @Component + @MessageMapping(opcode=xxx)
启动时 HandlerRegistry 统一 register 到 dispatcher五、心跳检测实现
5.1 心跳处理器
java
package com.game.gateway.handler;
import com.game.gateway.codec.GameMessage;
import com.game.gateway.dispatch.MessageDispatcher;
import com.game.gateway.session.GameSession;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.handler.timeout.IdleState;
import io.netty.handler.timeout.IdleStateEvent;
public class HeartbeatHandler extends ChannelInboundHandlerAdapter {
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
if (evt instanceof IdleStateEvent) {
IdleStateEvent event = (IdleStateEvent) evt;
if (event.state() == IdleState.READER_IDLE) {
// 读超时:60s 无任何数据,判定死连接
ctx.close();
}
} else {
ctx.fireUserEventTriggered(evt);
}
}
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) {
// 任何消息到达都刷新活跃时间
GameSession session = SessionHolder.get(ctx);
if (session != null) {
session.touch();
}
ctx.fireChannelRead(msg);
}
}心跳消息处理位置:
心跳 Opcode(1)在 DispatchHandler 前拦截
直接回一个空心跳响应,不进入业务分发
玩家挂机但连接健康 → 心跳持续,不被误杀六、Pipeline 装配与服务器启动
6.1 初始化器
java
package com.game.gateway.netty;
import com.game.gateway.codec.GameMessageDecoder;
import com.game.gateway.codec.GameMessageEncoder;
import com.game.gateway.dispatch.MessageDispatcher;
import com.game.gateway.handler.DispatchHandler;
import com.game.gateway.handler.HeartbeatHandler;
import com.game.gateway.session.SessionManager;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.ChannelPipeline;
import io.netty.channel.socket.SocketChannel;
import io.netty.handler.timeout.IdleStateHandler;
import java.util.concurrent.TimeUnit;
public class GameServerInitializer extends ChannelInitializer<SocketChannel> {
private final MessageDispatcher dispatcher;
private final SessionManager sessionManager;
public GameServerInitializer(MessageDispatcher dispatcher, SessionManager sessionManager) {
this.dispatcher = dispatcher;
this.sessionManager = sessionManager;
}
@Override
protected void initChannel(SocketChannel ch) {
ChannelPipeline p = ch.pipeline();
p.addLast("idle", new IdleStateHandler(60, 0, 0, TimeUnit.SECONDS));
p.addLast("decoder", new GameMessageDecoder());
p.addLast("encoder", new GameMessageEncoder());
p.addLast("heartbeat", new HeartbeatHandler());
p.addLast("dispatch", new DispatchHandler(dispatcher, sessionManager));
}
}6.2 服务器启动
java
package com.game.gateway.netty;
import com.game.gateway.session.SessionManager;
import com.game.gateway.dispatch.MessageDispatcher;
import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelOption;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class GameNettyServer {
private static final Logger log = LoggerFactory.getLogger(GameNettyServer.class);
private final int port;
private final int bossThreads;
private final int workerThreads;
public GameNettyServer(int port, int bossThreads, int workerThreads) {
this.port = port;
this.bossThreads = bossThreads;
this.workerThreads = workerThreads;
}
public void start(MessageDispatcher dispatcher, SessionManager sessionManager) {
NioEventLoopGroup boss = new NioEventLoopGroup(bossThreads);
NioEventLoopGroup worker = new NioEventLoopGroup(workerThreads);
try {
ServerBootstrap bootstrap = new ServerBootstrap();
bootstrap.group(boss, worker)
.channel(NioServerSocketChannel.class)
.childHandler(new GameServerInitializer(dispatcher, sessionManager))
.option(ChannelOption.SO_BACKLOG, 1024)
.childOption(ChannelOption.TCP_NODELAY, true)
.childOption(ChannelOption.SO_KEEPALIVE, true);
ChannelFuture future = bootstrap.bind(port).sync();
log.info("Game Netty Server started on port {}", port);
future.channel().closeFuture().sync();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
boss.shutdownGracefully();
worker.shutdownGracefully();
}
}
}6.3 启动入口
java
package com.game;
import com.game.gateway.dispatch.HandlerRegistry;
import com.game.gateway.dispatch.MessageDispatcher;
import com.game.gateway.netty.GameNettyServer;
import com.game.gateway.session.SessionManager;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.ConfigurableApplicationContext;
@SpringBootApplication(scanBasePackages = "com.game")
public class GameServerApplication {
public static void main(String[] args) throws Exception {
ConfigurableApplicationContext context = SpringApplication.run(GameServerApplication.class, args);
MessageDispatcher dispatcher = context.getBean(MessageDispatcher.class);
SessionManager sessionManager = context.getBean(SessionManager.class);
HandlerRegistry registry = context.getBean(HandlerRegistry.class);
registry.registerAll(dispatcher);
int port = context.getEnvironment().getProperty("game.netty.port", Integer.class, 8888);
GameNettyServer server = new GameNettyServer(port, 1, Runtime.getRuntime().availableProcessors() * 2);
server.start(dispatcher, sessionManager);
}
}七、验证与压测要点
7.1 功能验证清单
1. 连接建立:客户端 connect 成功
2. 登录链路:发 LoginReq(Opcode 1001)→ 收 LoginResp(1002)
3. 心跳保活:30s 心跳,连接持续在线
4. 断线清理:停发心跳 60s+ → 服务器主动断开
5. 顶号处理:同账号二次登录 → 旧连接被顶掉
6. 非法数据:发错误 Magic → 连接被关闭
7. 超长消息:Length 超限 → 连接被关闭调试工具:
自写 Java/Python 测试客户端
Wireshark 抓包核对协议字节
日志打印收发十六进制报文7.2 压测要点
关注指标:
建连数、在线数(Session 数)
单 Opcode 处理耗时(P99)
Worker 线程利用率
GC 频率与内存(ByteBuf 泄漏检查)
压测方法:
模拟客户端连接池,每个连接发登录 + 心跳
观察吞吐瓶颈在编解码还是业务
对比调整 Worker 线程数与业务线程池八、小结
本篇把协议、会话、心跳、分发四条链路合成一个可运行的 Netty 网关:GameMessageDecoder/Encoder 处理字节与对象互转,GameSession 封装连接与玩家状态,HeartbeatHandler 借助 IdleStateHandler 清理死连接,MessageDispatcher 按 Opcode 路由到各业务 Handler。整条链路职责单一、边界清晰,通过 Spring 注解自动注册 Handler 让后续新增玩法消息几乎零成本。运行后依次验证登录、心跳、断线、顶号、非法数据五种场景,即证明通信层已就绪,可以承载第 3 周起的玩家系统与房间对战业务。