Nettyフレームワークの実践:アーキテクチャ理解と初期実装

標準NIOの課題とNettyの選定理由

Java標準のNIOライブラリは、低レイヤーのネットワーク処理を提供する一方で、以下の点で開発コストが高騰しやすい。

  • APIの複雑さ:Selector、ServerSocketChannel、SocketChannel、ByteBufferなどのクラス群を適切に組み合わせる必要がある。
  • スレッド管理の負担:Reactorパターンを正しく実装するには、並行処理とネットワークI/Oの深い知識が不可欠。
  • epoll空ループバグ:特定条件下でSelectorが空のスピンを引き起こし、CPU使用率100%状態に陥る問題が長期にわたり存在した。

Nettyはこれらの課題を抽象化し、以下の利点で生産性と信頼性を向上させる。

  • 直感的なAPI設計と低い習得コスト
  • HTTP/SSL/Protobuf/ZIPなど多数のプロトコルとコーデックを標準サポート
  • ゼロコピーやバッファ最適化による業界最高級クラスのパフォーマンス
  • 活発なコミュニティによる迅速なバグ修正と機能追加
  • DubboやElasticsearchなど主要ミドルウェアで実績が検証されている

Nettyのレイヤードアーキテクチャ

Nettyは公式サイトで公開されている図示通り、明確な役割分担を持つ3層構造で設計されている。

  • Core層:ゼロコピー技術、再利用可能なバッファ、拡張可能なイベント駆動モデルを基盤として提供する。
  • Protocol Support層:HTTP、WebSocket、SSL/TLS、Protobuf、gzip、大ファイル転送などのプロトコル処理を標準ライブラリとして備える。
  • Transport Services層:TCP/UDPソケット、HTTP Tunnel、Datagramなどの物理伝送路を抽象化し、プラットフォーム依存を隠蔽する。

Reactorモデルとイベントループの動作設計

Nettyは接続制御とデータ送受信用の処理を分離し、スケーラビリティを確保する。

  1. Accept用スレッドプール(BossGroup):新規クライアント接続を監視・受け付け、確立後はWorkerに委任する。
  2. I/O用スレッドプール(WorkerGroup):接続済みのソケットに対する読み書きとパイプライン処理を担当する。
  3. NioEventLoop:各グループ内部の実際の処理スレッド。セレクタを保有し、登録されたチャンネルのI/Oイベントをポーリングする。
  4. イベント処理フロー:Bossはacceptイベントを検知しNioSocketChannelを生成・Workerにバインドする。Workerはread/writeイベントを処理し、パイプライン内のハンドラーを実行する。
  5. 非同期タスクの委譲:ハンドラー内で重いビジネスロジックを実行するとI/Oブロックが生じるため、EventLoopのタスクキューやスケジューラに非同期的に処理を投げる設計が推奨される。

依存関係とビルド設定

Mavenプロジェクトでは、以下の依存を追加してNettyの全機能を有効化する。

<dependency>
    <groupId>io.netty</groupId>
    <artifactId>netty-all</artifactId>
    <version>4.1.90.Final</version>
</dependency>

サーバー側の実装

接続を受け付け、リクエストを処理するサーバー構成を定義する。ブートストラップは設定を連鎖メソッドで構築し、イベントループの終了を同期的に待機する。

package org.example.tcp;

import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.*;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioServerSocketChannel;

public class TcpServerStarter {

    public static void main(String[] args) {
        new TcpServerStarter().launch();
    }

    public void launch() {
        EventLoopGroup acceptorPool = new NioEventLoopGroup();
        EventLoopGroup ioWorkerPool = new NioEventLoopGroup();

        try {
            ServerBootstrap serverBoot = new ServerBootstrap();
            serverBoot.group(acceptorPool, ioWorkerPool)
                    .channel(NioServerSocketChannel.class)
                    .option(ChannelOption.SO_BACKLOG, 256)
                    .childOption(ChannelOption.SO_KEEPALIVE, true)
                    .childHandler(new ChannelInitializer<SocketChannel>() {
                        @Override
                        protected void initChannel(SocketChannel ch) {
                            ch.pipeline().addLast(new ServerRequestProcessor());
                        }
                    });

            System.out.println("[SERVER] リッスン準備完了");
            ChannelFuture bindFuture = serverBoot.bind(7000).sync();
            bindFuture.channel().closeFuture().sync();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        } finally {
            acceptorPool.shutdownGracefully();
            ioWorkerPool.shutdownGracefully();
        }
    }
}

受信データをパースし、重い処理はイベントループのタスクキューに委譲するハンドラー。

package org.example.tcp;

import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.util.CharsetUtil;

import java.util.concurrent.TimeUnit;

public class ServerRequestProcessor extends ChannelInboundHandlerAdapter {

    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) {
        ByteBuf payload = (ByteBuf) msg;
        String text = payload.toString(CharsetUtil.UTF_8);
        System.out.println("[SERVER] 受信: " + ctx.channel().remoteAddress() + " -> " + text);

        // ブロックする可能性のある処理は非同期で実行
        ctx.channel().eventLoop().submit(() -> {
            try { TimeUnit.SECONDS.sleep(2); } catch (InterruptedException ignored) {}
            System.out.println("[SERVER] バックグラウンド処理完了: " + text);
        });

        // 遅延実行キューへの登録
        ctx.channel().eventLoop().schedule(() -> {
            System.out.println("[SERVER] 遅延処理実行: " + text);
        }, 5, TimeUnit.SECONDS);
    }

    @Override
    public void channelReadComplete(ChannelHandlerContext ctx) {
        ctx.writeAndFlush(Unpooled.copiedBuffer("ACK: メッセージ受信確認", CharsetUtil.UTF_8));
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
        cause.printStackTrace();
        ctx.close();
    }
}

クライアント側の実装

サーバーへ接続し、初期データの送信とレスポンスの受信を担うクライアント構成。

package org.example.tcp;

import io.netty.bootstrap.Bootstrap;
import io.netty.channel.*;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;

public class TcpClientConnector {

    public static void main(String[] args) {
        new TcpClientConnector().launch();
    }

    public void launch() {
        NioEventLoopGroup clientLoopGroup = new NioEventLoopGroup();

        try {
            Bootstrap clientBoot = new Bootstrap();
            clientBoot.group(clientLoopGroup)
                    .channel(NioSocketChannel.class)
                    .handler(new ChannelInitializer<SocketChannel>() {
                        @Override
                        protected void initChannel(SocketChannel ch) {
                            ch.pipeline().addLast(new ClientResponseHandler());
                        }
                    });

            System.out.println("[CLIENT] 接続準備完了");
            ChannelFuture connFuture = clientBoot.connect("127.0.0.1", 7000).sync();
            connFuture.addListener(f -> {
                if (f.isSuccess()) System.out.println("[CLIENT] 接続成功");
                else System.out.println("[CLIENT] 接続失敗");
            });
            connFuture.channel().closeFuture().sync();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        } finally {
            clientLoopGroup.shutdownGracefully();
        }
    }
}

チャンネルが有効になった時点でデータを投入し、サーバーからの応答を記録するハンドラー。

package org.example.tcp;

import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.util.CharsetUtil;

public class ClientResponseHandler extends ChannelInboundHandlerAdapter {

    @Override
    public void channelActive(ChannelHandlerContext ctx) {
        String payload = "初回接続テストデータ";
        System.out.println("[CLIENT] 送信: " + payload);
        ctx.writeAndFlush(Unpooled.copiedBuffer(payload, CharsetUtil.UTF_8));
    }

    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) {
        ByteBuf resp = (ByteBuf) msg;
        System.out.println("[CLIENT] 受信: " + ctx.channel().remoteAddress() + " -> " + resp.toString(CharsetUtil.UTF_8));
    }
}

動作確認とログ出力

サーバーを先に起動し、その後クライアントを実行すると以下のような相互作用が観察できる。

サーバー側ログ

[SERVER] リッスン準備完了
[SERVER] 受信: /127.0.0.1:54321 -> 初回接続テストデータ
[SERVER] バックグラウンド処理完了: 初回接続テストデータ
[SERVER] 遅延処理実行: 初回接続テストデータ

クライアント側ログ

[CLIENT] 接続準備完了
[CLIENT] 送信: 初回接続テストデータ
[CLIENT] 接続成功
[CLIENT] 受信: /127.0.0.1:7000 -> ACK: メッセージ受信確認

タグ: Netty NIO Java イベントループ ネットワークプログラミング

7月31日 04:37 投稿