Netty, un framework asynchrone piloté par les événements, repose sur plusieurs composants fondamentaux qui orchestrent son fonctionnement pour une communication réseau efficace et performante. Ces éléments clés incluent :
- Bootstrap (ou amorceur)
- EventLoopGroup (groupe de boucles d'événements)
- Channel (canal)
- ChannelHandler (gestionnaire de canal)
- ChannelPipeline (pipeline de canal)
- ChannelHandlerContext (contexte de gestionnaire de canal)
- ChannelOption (options de canal)
- ByteBuf (tampon d'octets)
- ChannelFuture (futur de canal)
Bootstrap
Le composant Bootstrap est l'outil principal pour configurer et lancer les applications Netty, qu'il s'agisse de serveurs ou de clients. Il simplifie l'assemblage des diverses pièces de Netty, évitant ainsi un travail fastidieux de configuration manuelle. On distingue deux types :
ServerBootstrap: Utilisé pour configurer et démarrer les serveurs Netty.Bootstrap: Destiné à la configuration des clients Netty.
Voici un exemple illustrant la configuration d'un client Netty via Bootstrap :
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.ChannelOption;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;
import io.netty.buffer.PooledByteBufAllocator;
import java.net.InetSocketAddress;
public class NettyClientLauncher {
public static void main(String[] args) throws InterruptedException {
startClient("localhost", 8080);
}
public static void startClient(String host, int port) throws InterruptedException {
Bootstrap clientBootstrap = new Bootstrap(); // Instanciation de l'amorceur client
ChannelFuture connectionResult = clientBootstrap.group(new NioEventLoopGroup()) // Affectation d'un EventLoopGroup
.channel(NioSocketChannel.class) // Spécification du type de canal (NIO Socket)
.remoteAddress(new InetSocketAddress(host, port)) // Définition de l'adresse du serveur distant
.option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT) // Configuration de l'allocateur de mémoire
.handler(new ChannelInitializer<SocketChannel>() { // Ajout d'un initialiseur de canal
@Override
protected void initChannel(SocketChannel ch) {
ch.pipeline().addLast(); // Personnalisation du pipeline du canal
}
})
.connect(); // Initialisation de la connexion
connectionResult.sync(); // Bloque jusqu'à ce que la connexion soit établie
// Le canal est maintenant disponible pour les opérations d'écriture/lecture
// connectionResult.channel().writeAndFlush("Hello Netty!");
}
}
Le Bootstrap client requiert la spécification d'un EventLoopGroup, du type de Channel (par exemple, NioSocketChannel), et de l'adresse du serveur cible. Des options TCP peuvent également être configurées via option(). Contrairement aux serveurs, les clients n'ont pas de concept de canaux "parents" et "enfants", d'où l'utilisation unique des méthodes option() et handler(). La méthode connect() initialise la connexion de manière asynchrone, retournant un ChannelFuture pour gérer le résultat. Ce ChannelFuture permet d'accéder au Channel une fois la connexion établie, autorisant alors les opérations d'E/S.
EventLoopGroup
L'EventLoopGroup étend le modèle de thread Reactor, gérant les opérations d'E/S avec un pool de threads pour exploiter efficacement la concurrence des CPU modernes. Dans une architecture serveur, il se divise en deux rôles :
- Boss EventLoopGroup : Sa tâche principale est d'accepter les nouvelles connexions clientes.
- Worker EventLoopGroup : Une fois qu'une connexion est acceptée par le Boss, le Worker prend le relais pour gérer toutes les opérations d'E/S et le transfert de données associées à cette connexion.
Netty propose différentes implémentations d'EventLoopGroup en fonction de la plateforme et des besoins :
NioEventLoopGroup: L'implémentation la plus courante, basée sur NIO pour les E/S non bloquantes.EpollEventLoopGroup: Spécifique à Linux, utilisant le mécanismeepollpour des performances d'E/S optimisées.KQueueEventLoopGroup: Pour les systèmes de type Unix comme macOS ou BSD, utilisantkqueue.OioEventLoopGroup: Une implémentation pour les E/S bloquantes (BIO), généralement considérée comme obsolète.EmbeddedEventLoop: Utilisé principalement pour les tests.
Channel
Le Channel de Netty est une abstraction d'un canal de communication réseau, similaire au SocketChannel de NIO, mais offrant des fonctionnalités enrichies. Il représente une connexion ouverte vers une entité réseau (fichier, socket, etc.) et est le point central pour toutes les opérations d'E/S. Ses fonctionnalités principales incluent :
- Accès à des objets liés comme l'
EventLoopou leChannelFuture. - API pour vérifier l'état du canal (
isOpen(),isRegistered(),isActive()). - Méthodes pour déclencher des événements d'E/S (
read(),write(),flush()).
ChannelHandler, ChannelHandlerContext et ChannelPipeline
Ces trois composants sont au cœur du traitemant de la logique métier dans Netty.
ChannelHandler
Un ChannelHandler est un composant qui intercepte et traite les événements d'E/S ou transforme les données circulant dans un Channel. Ils sont séparés en deux catégories :
ChannelInboundHandler: Traite les événements "entrants" (lectures, connexions, exceptions).ChannelOutboundHandler: Traite les événements "sortants" (écritures, déconnexions).
ChannelPipeline
Chaque Channel possède son propre ChannelPipeline, qui est une liste d'instances de ChannelHandler. Le Pipeline est comparable à un filtre en chaîne, où chaque gestionnaire peut intercepter et traiter les données ou les événements avant de les passer au gestionnaire suivant. Lors de la configuration d'un serveur, vous définissez les gestionnaires enfants via childHandler() :
serverBootstrap.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel socketChannel) {
socketChannel.pipeline().addLast(new MyBusinessLogicHandler());
}
});
Le ChannelPipeline est implémenté comme une liste doublement chaînée, avec des pointeurs head et tail. Tous les ChannelHandler, qu'ils soient entrants ou sortants, sont ajoutés à ce pipeline.
Les événements entrants (par exemple, la lecture de données) se propagent du head vers le tail. Chaque ChannelInboundHandler rencontre l'événement et peut le traiter, le modifier ou le propager au prochain gestionnaire. La propagation peut être interrompue si un gestionnaire ne transmet pas l'événement explicitement (par exemple, en n'appelant pas super.channelRead(ctx, msg) ou ctx.fireChannelRead(msg)).
Les événements sortants (par exemple, l'écriture de données) se propagent du tail vers le head. Chaque ChannelOutboundHandler rencontre l'événement et peut le traiter, le modifier ou le propager. Il n'est généralement pas possible d'interrompre le flux d'événements sortants dans le pipeline de la même manière que les événements entrants.
### ChannelHandlerContext
Chaque ChannelHandler est encapsulé dans un ChannelHandlerContext lorsqu'il est ajouté au ChannelPipeline. Le ChannelHandlerContext est le lien entre un ChannelHandler et son ChannelPipeline. Il fournit un moyen d'interagir avec d'autres gestionnaires dans le pipeline ou de déclencher des événements spécifiques. Ses méthodes fireXxx() permettent de propager un événement à travers le pipeline, en commençant par le prochain gestionnaire dans la direction appropriée. Par exemple :
ChannelHandlerContext fireChannelRegistered()ChannelHandlerContext fireChannelActive()ChannelHandlerContext fireChannelRead(Object var1)
Ces méthodes sont essentielles pour la gestion des événements et le contrôle du flux de données dans le pipeline.
ChannelOption
Les ChannelOption sont des paramètres qui permettent de configurer le comportement du socket sous-jacent du Channel, influençant ainsi des aspects du protocole TCP/IP. Par exemple, ils peuvent définir des propriétés telles que SO_KEEPALIVE, TCP_NODELAY, ou la taille du tampon de réception/envoi. Ces options sont cruciales pour optimiser les performances réseau et l'interaction avec le système d'exploitation.
ByteBuf
Le ByteBuf de Netty est une amélioration par rapport au ByteBuffer de NIO, offrant une API plus flexible et performante pour la manipulation des octets. Il est conçu pour la gestion des données binaires et offre une meilleure ergonomie.
La structure d'un ByteBuf se compose de quatre sections logiques :
- Déchets : Les données avant
readerIndex, considérées comme lues et disponibles pour être écrasées. - Données lisibles : Les octets entre
readerIndexetwriterIndex. - Espace inscriptible : L'espace disponible entre
writerIndexetcapacity. - Espace maximum inscriptible : L'espace entre
capacityetmaxCapacity, pouvant être étendu si nécessaire.
La capacity est la taille actuelle du tampon, qui peut augmenter dynamiquement jusqu'à maxCapacity si plus d'espace est nécessaire pour les opérations d'écriture.
Voici un exemple illustrant la gestion des indices et de la capacité du ByteBuf :
import io.netty.buffer.ByteBuf;
import io.netty.buffer.PooledByteBufAllocator;
import org.junit.jupiter.api.Test;
public class ByteBufManipulationTest {
private void displayBufferInfo(ByteBuf buffer) {
System.out.println("--- Infos ByteBuf ---");
System.out.println("readerIndex=" + buffer.readerIndex());
System.out.println("readableBytes=" + buffer.readableBytes());
System.out.println("writerIndex=" + buffer.writerIndex());
System.out.println("writableBytes=" + buffer.writableBytes());
System.out.println("capacity=" + buffer.capacity());
System.out.println("maxCapacity=" + buffer.maxCapacity());
System.out.println("maxWritableBytes=" + buffer.maxWritableBytes());
System.out.println("---------------------\n");
}
@Test
public void testByteBufGrowth() {
ByteBuf buffer = PooledByteBufAllocator.DEFAULT.buffer(8, 64); // Capacité initiale 8, max 64
displayBufferInfo(buffer); // Initial: reader=0, writer=0, capacity=8
buffer.writeInt(12345); // Écrit 4 octets
displayBufferInfo(buffer); // Après 1ère écriture: writer=4
buffer.writeInt(67890); // Écrit 4 octets
displayBufferInfo(buffer); // Après 2nde écriture: writer=8. Capacity peut être ajustée.
buffer.writeInt(101112); // Écrit 4 octets de plus, dépasse la capacité initiale de 8
displayBufferInfo(buffer); // Capacity devrait avoir augmenté automatiquement
}
}
Netty propose deux types d'allocateurs de mémoire pour ByteBuf :
PooledByteBufAllocator: Un allocateur avec mise en pool qui réutilise la mémoire déjà allouée, réduisant ainsi les frais d'allocation et de désallocation. Il alloue par défaut sur le tas Java.UnpooledByteBufAllocator: Un allocateur non mis en pool, qui alloue de la nouvelle mémoire à chaque fois. Il alloue par défaut sur le tas Java.
Les deux allocateurs peuvent créer des tampons "directs" (hors tas Java) via la méthode directBuffer(). Les tampons directs offrent une meilleure performance d'E/S mais nécessitent une gestion manuelle de la libération de mémoire pour éviter les fuites, contrairement aux tampons sur le tas Java qui sont gérés par le GC.
| Vitesse d'allocation | Efficacité d'utilisation | |
|---|---|---|
| Mémoire sur le tas (Java Heap) | Rapide | Faible (GC overhead) |
| Mémoire hors tas (Direct/Native) | Lente | Élevée (Direct E/S) |
Les méthodes de lecture de ByteBuf se distinguent : les méthodes readXxx() avancent le readerIndex, tandis que les méthodes getXxx() ne le modifient pas.
import io.netty.buffer.ByteBuf;
import io.netty.buffer.UnpooledByteBufAllocator;
import org.junit.jupiter.api.Test;
public class ByteBufReadGetTest {
private void displayBufferInfo(ByteBuf buffer) {
System.out.println("--- Infos ByteBuf ---");
System.out.println("readerIndex=" + buffer.readerIndex());
System.out.println("readableBytes=" + buffer.readableBytes());
System.out.println("writerIndex=" + buffer.writerIndex());
System.out.println("writableBytes=" + buffer.writableBytes());
System.out.println("capacity=" + buffer.capacity());
System.out.println("maxCapacity=" + buffer.maxCapacity());
System.out.println("maxWritableBytes=" + buffer.maxWritableBytes());
System.out.println("---------------------\n");
}
@Test
public void testReadAndGetOperations() {
ByteBuf buffer = UnpooledByteBufAllocator.DEFAULT.buffer(16, 128);
displayBufferInfo(buffer); // Initial
buffer.writeInt(100); // Écrit 4 octets
displayBufferInfo(buffer);
buffer.writeInt(200); // Écrit 4 octets supplémentaires
displayBufferInfo(buffer);
int firstValue = buffer.readInt(); // Lit 4 octets, avance readerIndex
System.out.println("Valeur lue (readInt): " + firstValue);
displayBufferInfo(buffer);
int secondValuePeek = buffer.getInt(4); // Lit 4 octets à l'index 4, n'avance PAS readerIndex
System.out.println("Valeur obtenue (getInt@4): " + secondValuePeek);
displayBufferInfo(buffer);
}
}
Zéro-Copie dans Netty
Le concept de "zéro-copie" dans Netty vise à minimiser ou éliminer les copies de données entre différentes zones de mémoire. Bien que le terme soit également utilisé au niveau du système d'exploitation pour éviter les copies entre espace noyau et espace utilisateur (par exemple via mmap ou sendfile), Netty l'applique au niveau de l'application de plusieurs manières :
-
CompositeByteBuf: Permet de combiner plusieursByteBufen une seule vue logique, évitant ainsi de copier les données des tampons originaux dans un nouveau tampon unifié.Considérons un scénario où vous devez envoyer un en-tête et un corps de message séparément, puis les assembler pour l'envoi :
import io.netty.buffer.ByteBuf; import io.netty.buffer.PooledByteBufAllocator; import io.netty.buffer.CompositeByteBuf; import io.netty.buffer.Unpooled; import org.junit.jupiter.api.Test; import java.nio.charset.StandardCharsets; public class ZeroCopyExamples { @Test public void testManualConcatenation() { ByteBuf header = PooledByteBufAllocator.DEFAULT.buffer(); ByteBuf body = PooledByteBufAllocator.DEFAULT.buffer(); header.writeBytes("HTTP/1.1 200 OK\r\n".getBytes(StandardCharsets.UTF_8)); body.writeBytes("Ceci est le corps du message.".getBytes(StandardCharsets.UTF_8)); // Nécessite une copie des données de 'header' et 'body' ByteBuf combinedBuffer = PooledByteBufAllocator.DEFAULT.buffer(header.readableBytes() + body.readableBytes()); combinedBuffer.writeBytes(header); combinedBuffer.writeBytes(body); // ... Envoyer combinedBuffer ... System.out.println("Buffer combiné (avec copie): " + combinedBuffer.toString(StandardCharsets.UTF_8)); header.release(); body.release(); combinedBuffer.release(); } @Test public void testCompositeByteBuf() { ByteBuf header = PooledByteBufAllocator.DEFAULT.buffer(); ByteBuf body = PooledByteBufAllocator.DEFAULT.buffer(); header.writeBytes("HTTP/1.1 200 OK\r\n".getBytes(StandardCharsets.UTF_8)); body.writeBytes("Ceci est le corps du message.".getBytes(StandardCharsets.UTF_8)); // Utilisation de CompositeByteBuf pour éviter la copie CompositeByteBuf composite = PooledByteBufAllocator.DEFAULT.compositeBuffer(); composite.addComponents(true, header, body); // 'true' ajuste writerIndex automatiquement // Le composite représente les deux ByteBufs sans les copier System.out.println("CompositeByteBuf (zéro-copie): " + composite.toString(StandardCharsets.UTF_8)); composite.release(); // Libère aussi header et body si addComponents a pris possession } @Test public void testUnpooledWrappedBuffer() { ByteBuf header = PooledByteBufAllocator.DEFAULT.buffer(); ByteBuf body = PooledByteBufAllocator.DEFAULT.buffer(); header.writeBytes("L'en-tête.\r\n".getBytes(StandardCharsets.UTF_8)); body.writeBytes("Le contenu du corps.".getBytes(StandardCharsets.UTF_8)); // Unpooled.wrappedBuffer est une autre forme de zéro-copie pour combiner des buffers ByteBuf wrapped = Unpooled.wrappedBuffer(header, body); System.out.println("WrappedBuffer (zéro-copie): " + wrapped.toString(StandardCharsets.UTF_8)); wrapped.release(); // Libère aussi header et body si Unpooled.wrappedBuffer a pris possession } } -
Opérations
wrap: La classe utilitaireUnpooledfournit des méthodeswrappedBuffer()qui permettent d'encapsuler des tableaux d'octets (byte[]), desByteBufferou d'autresByteBufexistants dans un nouveauByteBufsans copier les données. Le nouveauByteBufpartage la même mémoire que les données sources.Par exemple, au lieu de copier un
byte[]dans unByteBuf:// Copie les octets ByteBuf buffer = PooledByteBufAllocator.DEFAULT.buffer(); byte[] data = "Hello Netty".getBytes(StandardCharsets.UTF_8); buffer.writeBytes(data);Vous pouvez utiliser
Unpooled.wrappedBuffer()pour éviter la copie :// Encapsule les octets sans copie byte[] data = "Hello Netty".getBytes(StandardCharsets.UTF_8); ByteBuf buffer = Unpooled.wrappedBuffer(data);Unpooledoffre de nombreuses surcharges pourwrappedBuffer, adaptées à divers scénarios. -
Opérations
slice: À l'opposé dewrap, l'opérationslice()permet de créer des vues logiques (des sous-sections) d'unByteBufexistant. Ces "tranches" partagent la même mémoire sous-jacente que leByteBuforiginal. Chaque tranche a ses propresreaderIndexetwriterIndex, mais les modifications des données dans l'une se reflètent dans les autres.ByteBuf originalBuffer = ...;
ByteBuf headerSlice = originalBuffer.slice(0, 10);
ByteBuf bodySlice = originalBuffer.slice(10, originalBuffer.readableBytes() - 10);
-
FileRegionpour le transfert de fichiers : Netty utiliseFileRegion(qui encapsule unFileChannel) pour transférer des fichiers directement de la mémoire du noyau vers un socket, utilisant la fonctionnalitétransferTo()du système d'exploitation. Cela permet un transfert de fichiers à zéro-copie au niveau du noyau (si SSL n'est pas activé).import io.netty.channel.ChannelFutureListener; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.DefaultFileRegion; import io.netty.channel.SimpleChannelInboundHandler; import io.netty.handler.ssl.SslHandler; import io.netty.handler.stream.ChunkedFile; import java.io.RandomAccessFile; public class FileTransferServerHandler extends SimpleChannelInboundHandler<String> { @Override public void channelActive(ChannelHandlerContext ctx) { ctx.writeAndFlush("BIENVENUE: Entrez le chemin du fichier à récupérer.\n"); } @Override protected void channelRead0(ChannelHandlerContext ctx, String requestPath) throws Exception { RandomAccessFile targetFile = null; long fileSize = -1; try { targetFile = new RandomAccessFile(requestPath, "r"); fileSize = targetFile.length(); } catch (Exception e) { ctx.writeAndFlush("ERREUR: " + e.getClass().getSimpleName() + ": " + e.getMessage() + '\n'); return; } finally { if (fileSize < 0 && targetFile != null) { targetFile.close(); } } ctx.write("OK: Taille du fichier " + targetFile.length() + '\n'); if (ctx.pipeline().get(SslHandler.class) == null) { // Pas de SSL - Zéro-copie possible via FileRegion. ctx.write(new DefaultFileRegion(targetFile.getChannel(), 0, fileSize)); } else { // SSL activé - Zéro-copie impossible, utilise ChunkedFile. ctx.write(new ChunkedFile(targetFile)); } ctx.writeAndFlush("\n"); } @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); if (ctx.channel().isActive()) { ctx.writeAndFlush("ERREUR Interne: " + cause.getClass().getSimpleName() + ": " + cause.getMessage() + '\n').addListener(ChannelFutureListener.CLOSE); } } }
ChannelFuture
Toutes les opérations d'E/S dans Netty sont non bloquantes et asynchrones. Lorsqu'une opération d'E/S est lancée (par exemple, connect(), write(), flush()), elle retourne immédiatement un objet ChannelFuture. Ce "futur" représente le résultat d'une opération asynchrone qui n'est pas encore terminée.
L'interface ChannelFuture offre des méthodes pour :
- Obtenir le
Channelassocié à l'opération :Channel channel(). - Enregistrer des écouteurs (
GenericFutureListener) qui seront notifiés lorsque l'opération est terminée (succès ou échec). C'est le moyen préféré de gérer les résultats asynchrones. - Bloquer l'exécution du thread courant jusqu'à ce que l'opération soit terminée (
sync(),await()). Ces méthodes transforment une opération asynchrone en une opération synchrone bloquante, mais leur utilisation doit être minimisée pour préserver la nature non bloquante de Netty.
Un GenericFutureListener est un rappel qui est invoqué une fois l'opération d'E/S complétée, permetant de réagir au succès ou à l'échec de l'opération.