Dans un projet Spring Boot, ajoutez la dépendance Netty via Maven :
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-all</artifactId>
<version>4.1.25.Final</version>
</dependency>
Serveur Netty avec gestion de heartbeat :
public class ServeurNetty {
public static void main(String[] args) {
EventLoopGroup groupePrincipal = new NioEventLoopGroup();
EventLoopGroup groupeTravail = new NioEventLoopGroup();
ServerBootstrap config = new ServerBootstrap();
config.group(groupePrincipal, groupeTravail)
.channel(NioServerSocketChannel.class)
.childHandler(new CanalInitialiseur() {
@Override
protected void initChannel(SocketChannel canal) {
ChannelPipeline flux = canal.pipeline();
flux.addLast(new LineBasedFrameDecoder(18));
flux.addLast(new StringDecoder());
flux.addLast(new StringEncoder());
flux.addLast(new IdleStateHandler(3, 0, 0, TimeUnit.SECONDS));
flux.addLast(new GestionnaireHeartbeatServeur());
}
});
try {
ChannelFuture futur = config.bind(8888).sync();
futur.channel().closeFuture().sync();
} finally {
groupePrincipal.shutdownGracefully();
groupeTravail.shutdownGracefully();
}
}
}
Getsionnaire des événements de heartbeat côté serveur :
public class GestionnaireHeartbeatServeur extends SimpleChannelInboundHandler<String> {
private int compteurInactivite;
private long dernierTimestamp;
@Override
protected void channelRead0(ChannelHandlerContext ctx, String msg) {
if ("heartbeat".equals(msg)) {
if (System.currentTimeMillis() - dernierTimestamp >= 3000) {
compteurInactivite = 0;
}
}
}
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
if (evt instanceof IdleStateEvent) {
dernierTimestamp = System.currentTimeMillis();
compteurInactivite++;
if (compteurInactivite >= 3) {
ctx.channel().close();
}
}
}
}
Client Netty pour l'envoi de signaux heartbeat :
public class ClientNetty {
public static void main(String[] args) {
EventLoopGroup groupe = new NioEventLoopGroup();
Bootstrap configClient = new Bootstrap();
configClient.group(groupe)
.channel(NioSocketChannel.class)
.handler(new CanalInitialiseur() {
@Override
protected void initChannel(SocketChannel canal) {
canal.pipeline()
.addLast(new StringEncoder())
.addLast(new StringDecoder());
}
});
try {
ChannelFuture futur = configClient.connect("localhost", 8888).sync();
Canal canal = futur.channel();
while (canal.isActive()) {
Thread.sleep(new Random().nextInt(6) * 1000);
canal.writeAndFlush("heartbeat\n");
}
futur.channel().closeFuture().sync();
} finally {
groupe.shutdownGracefully();
}
}
}