【実務・中級編】Fiberを用いたリアルタイムデータストリーミング処理:WebSocketサーバーにおけるメッセージキューイングとFiberスケジューリング – PHPコア・内部エンジンと高速化・並行処理の極意解析バイブル

Fiberを用いたリアルタイムデータストリーミング処理:WebSocketサーバーにおけるメッセージキューイングとFiberスケジューリング

こんにちは。テクニカルリードの私だ。
コードレビューの際、「とりあえず動くから」という理由で、I/O待ちが発生する処理を同期的にブロックさせたり、無駄に重いプロセスフォークを乱発しているコードを見かけると、正直頭痛が覚える。

PHPは「リクエストライフサイクルが短く、共有ードメモリを持たないため非同期処理が苦手」というのは、もはや過去の遺物だ。PHP 8.1で導入された Fiber(ファイバー) を正しく掌握すれば、OSスレッドを消費せずに数千・数万の協調的マルチタスク(Cooperative Multitasking)をユーザーランドでー綺麗に制御できる。

今回は、WebSocketサーバーを題材に、イベントループとFiberを組み合わせた「ノンブロッキングなメッセージストリーミングとキューイング処理」の極意を、Zend VMの挙動やメモリ管理の視点を交えながら徹底的に解説しよう。

—

1. なぜ従来のPHPではリアルタイム処理が破綻するのか?

多くのWebエンジニアは、PHPを「HTTPリクエストを受けて、DBからデータを引いて、HTML/JSONを返して死ぬ言語」として認識している。しかし、RatchetやWorkerman、Ampといったリアクティブなエコシステムにおいて、PHPは長寿命プロセスとして常駐し、数万のソケットをハンドリングしている。

ここで問題になるのが I/Oブロッキング と コンテキストスイッチのコスト だ。

[クライアントA] ──┐
[クライアントB] ──┼─> [Socket Server (Event Loop)]
[クライアントC] ──┘ │
(同期I/Oでブロック)
▼
[DB / 外部API / キュー]

従来、重い外部API呼び出しやデータベースへのブロッキングクエリが走ると、そのスレッド(またはプロセス)全体の実行が完全に停止する。これを防ぐためにプロセスを数千個フォークすれば、OSのメモリ空間(RSS)とコンテキストスイッチのオーバーヘッドでCPUが溶ける。

Zend VMとFiberの内部挙動

Fiberは、「スタックを持つ独立した実行コンテキスト(zend_execute_data)」をユーザーランドで自由に中断(Suspend)および再開(Resume)する仕組みだ。

OSスレッドを切り替えるのではなく、PHPの仮想マシン(Zend VM)のコールスタックのポインタを退避・復元するため、コンテキストスイッチのコストが極めて低い。このFiberをイベントループ(Event Loop)のC10K問題(1万クライアント問題)の調停役として組み込むことで、「書いているコードは同期処理のように直感的でありながら、実行実態は完全に非同期ノンブロッキング」という究極のアーキテクチャが完成する。

—

2. 設計指針:Fiberベースのメッセージキューとイベント駆動

今回の設計の核心は以下の3点だ。

1. イベントループ駆動: `ext-ev` や `ext-uv`、あるいは純粋な `stream_select` をベースにしたイベントループでソケットの読み書きを監視する。
2. Fiberによるタスク分離: クライアントからのメッセージ受信、および重いメッセージ配信(外部APIへの転送やDBログ保存など)を、それぞれ独立したFiberとしてスケジューリングする。
3. バックプレッシャー(流量制御)付きキュー: 送信バッファが溢れた際にメモリが爆発(OOM)しないよう、Fiberのサスペンド機構を利用したメッセージキューイングを行う。

—

3. 実装:実務に耐えうるWebSocket・Fiberスケジューラー

それでは、純粋なPHP(外部非同期フレームワークに依存せず、コアの sockets 拡張と Fiber を使用)で構築した、リアルタイムメッセージストリーミング・サーバーのリファレンスコードを提示する。

  • 簡易イベントループとFiberスケジューラーを統合したWebSocketサーバーエンジン
  • /
    class FiberWebSocketServer
    {
    private \Socket $serverSocket;
    / @array /
    private array $clients = [];
    / @array /
    private array $fibers = [];
    / @array> メッセージ送信キュー /
    private array $sendQueues = [];
    private bool $isRunning = true;

    public function __construct(string $host, int $port)
    {
    // サーバーソケットの初期化
    $socket = socket_create(AF_INET, SOCK_STREAM, SOL_TCP);
    socket_set_option($socket, SOL_SOCKET, SO_REUSEADDR, 1);
    socket_bind($socket, $host, $port);
    socket_listen($socket, 128);
    socket_set_nonblock($socket);

    $this->serverSocket = $socket;
    echo “[INFO] WebSocket Server started on {$host}:{$port}\n”;
    }

    public function run(): void
    {
    while ($this->isRunning) {
    $read = [$this->serverSocket];
    $write = [];

    foreach ($this->clients as $id => $client) {
    $read[] = $client;
    // 送信キューにデータがある場合のみ、書き込み監視対象に追加する
    if (!empty($this->sendQueues[$id])) {
    $write[] = $client;
    }
    }

    $except = null;
    // タイムアウトを短く設定し、イベントループを高速回転させる (例: 5ms)
    if (socket_select($read, $write, $except, 0, 5000) === false) {
    continue;
    }

    // 1. 新規接続のハンドリング
    if (in_array($this->serverSocket, $read, true)) {
    $clientSocket = socket_accept($this->serverSocket);
    if ($clientSocket !== false) {
    socket_set_nonblock($clientSocket);
    $clientId = (int)$clientSocket;
    $this->clients[$clientId] = $clientSocket;
    $this->sendQueues[$clientId] = [];

    // 接続ハンドシェイク&メッセージ処理をFiberとして独立起動
    $this->spawnClientFiber($clientId, $clientSocket);
    }
    }

    // 2. 読み込み可能ソケットの処理(Fiberのレジューム)
    foreach ($read as $socket) {
    if ($socket === $this->serverSocket) continue;
    $clientId = (int)$socket;

    if (isset($this->fibers[$clientId]) && $this->fibers[$clientId]->isSuspended()) {
    // イベント発生をFiberに伝えて再開
    $this->fibers[$clientId]->resume($socket);
    }
    }

    // 3. 書き込み可能ソケットの処理(バックプレッシャー制御付きフラッシュ)
    foreach ($write as $socket) {
    $clientId = (int)$socket;
    $this->flushSendQueue($clientId);
    }
    }
    }

    private function spawnClientFiber(int $clientId, \Socket $client): void
    {
    $fiber = new Fiber(function () use ($clientId, $client) {
    // 1. WebSocketハンドシェイクの処理
    $headersRead = false;
    $buffer = ”;

    while (!$headersRead) {
    // データが来るまでFiberを中断して待機
    $socket = Fiber::suspend();
    $chunk = @socket_read($socket, 2048);

    if ($chunk === false || $chunk === ”) {
    $this->disconnectClient($clientId);
    return;
    }

    $buffer .= $chunk;
    if (str_contains($buffer, “\r\n\r\n”)) {
    $this->performHandshake($client, $buffer);
    $headersRead = true;
    echo “[CLIENT] Handshake completed for ID: {$clientId}\n”;
    }
    }

    // 2. メインのメッセージ送受信ループ
    while (true) {
    $socket = Fiber::suspend(); // I/Oイベントを待つ
    $data = @socket_read($socket, 2048);

    if ($data === false || $data === ”) {
    echo “[CLIENT] Connection lost for ID: {$clientId}\n”;
    $this->disconnectClient($clientId);
    break;
    }

    // WebSocketフレームのデコード(簡易版:実際にはマスク解除が必要)
    $message = $this->decodeWebSocketFrame($data);
    if ($message !== null) {
    echo “[MESSAGE] Received from #{$clientId}: {$message}\n”;

    // 非同期メッセージキューイング&ブロードキャストのシミュレーション
    $this->enqueueBroadcast($clientId, “Echo: ” . $message);
    }
    }
    });

    $this->fibers[$clientId] = $fiber;
    $fiber->start(); // 初回実行
    }

    private function enqueueBroadcast(int $senderId, string $message): void
    {
    $payload = $this->encodeWebSocketFrame($message);

    foreach ($this->clients as $id => $client) {
    // 全クライアントの送信キューに積む(バックプレッシャー管理)
    $this->sendQueues[$id][] = $payload;
    }
    }

    private function flushSendQueue(int $clientId): void
    {
    if (!isset($this->clients[$clientId], $this->sendQueues[$clientId])) {
    return;
    }

    $client = $this->clients[$clientId];
    while (!empty($this->sendQueues[$clientId])) {
    $payload = array_shift($this->sendQueues[$clientId]);
    $sent = @socket_write($client, $payload, strlen($payload));

    if ($sent === false) {
    $err = socket_last_error($client);
    // EAGAIN / EWOULDBLOCK の場合はバッファが一杯なのでキューに戻して次回に託す
    if ($err === SOCKET_EAGAIN || $err === SOCKET_EWOULDBLOCK) {
    array_unshift($this->sendQueues[$clientId], $payload);
    break;
    } else {
    $this->disconnectClient($clientId);
    break;
    }
    }
    }
    }

    private function disconnectClient(int $clientId): void
    {
    if (isset($this->clients[$clientId])) {
    @socket_close($this->clients[$clientId]);
    unset($this->clients[$clientId]);
    }
    unset($this->fibers[$clientId], $this->sendQueues[$clientId]);
    }

    private function performHandshake(\Socket $client, string $headers): void
    {
    // 簡易ハンドシェイク実装
    if (preg_match(“/Sec-WebSocket-Key: (.)\r\n/”, $headers, $matches)) {
    $key = trim($matches[1]);
    $acceptKey = base64_encode(pack(‘H’, sha1($key . ‘258EAFA5-E914-47DA-95CA-C5AB0DC85B11’)));
    $upgrade = “HTTP/1.1 101 Switching Protocols\r\n” .
    “Upgrade: websocket\r\n” .
    “Connection: Upgrade\r\n” .
    “Sec-WebSocket-Accept: $acceptKey\r\n\r\n”;
    socket_write($client, $upgrade, strlen($upgrade));
    }
    }

    private function decodeWebSocketFrame(string $data): ?string
    {
    // 実務ではRFC 6455に準拠したマスク解除やオプコード判定を厳密に行うこと
    // ここでは簡易的にペイロード部分のみを抽出するスタブ
    if (strlen($data) < 2) return null; $length = ord($data[1]) & 127; if ($length === 126) { return substr($data, 4); } elseif ($length === 127) { return substr($data, 10); } return substr($data, 6); // マスクキー4バイト含む場合の簡易オフセット } private function encodeWebSocketFrame(string $text): string { $b1 = 0x80 | (0x1 & 0x0f); // FIN + text frame $length = strlen($text); if ($length <= 125) { $header = pack('CC', $b1, $length); } elseif ($length > 125 && $length < 65536) { $header = pack('CCn', $b1, 126, $length); } else { $header = pack('CCJ', $b1, 127, $length); } return $header . $text; } } ---

    4. コードレビュー:なぜこの設計がプロダクションで生き残るのか?

    このコードをただの「おもちゃの非同期スクリプト」だと思っては困る。大規模なトラフィックを捌くシステムにおいて、以下の設計思想が組み込まれている。

    ① `socket_select` と Fiber の美しい協調

    イベントループ(`socket_select`)は、OSカーネルからの「このソケット読み書きできるよ」という通知を受け取る。従来のコールバック地獄(Node.jsの初期に見られたような複雑なネスト)や、複雑なPromiseチェーンを、Fiberの `Fiber::suspend()` と `resume()` によって「手続き型(上から下へ流れる)の美しいコード」に昇華させている。
    これにより、複雑なステートマシンを手動で管理する必要が消え、可読性と保守性が劇的に向上する。

    ② メモリ爆発を防ぐ「バックプレッシャー制御」

    リアルタイム配信で最も恐ろしいのは、「配信先のクライアントのネットワークが遅い(あるいは死んでいる)せいで、送信バッファにデータが無限に溜まり続け、PHPプロセスのメモリが枯渇してクラッシュ(OOM)する現象」だ。
    この実装では、`socket_write` が `SOCKET_EAGAIN` を返した瞬間、送信を諦めてキューに保持し、イベントループが「書き込み可能」と判断した瞬間(`$write` 配列に入ってきた時)にのみ送信を再開する。キューのサイズやメモリ使用量を適切に監視・制限すれば、巨大なトラフィックの波が来てもサーバーが沈没しない。

    —

    5. 運用時の注意点とさらなる高みへ

    このアーキテクチャを本番環境(Production)に投入するにあたり、以下の鉄則を忘れないでほしい。

    1. エラーハンドリングと例外のバブリング:
    Fiber内でキャッチされない例外(`Throwable`)が発生した場合、そのFiberは破棄され、該当するクライアント接続がゾンビ化する恐れがある。必ずFiberのクロージャ内全体を `try-catch` で囲み、エラーログの出力と接続のクリーンアップを行わせること。
    2. CPUバウンドな処理の排除:
    Fiberはあくまで「協調的」マルチタスクである。ひとつのFiber内で画像処理や重い暗号化ループ(数ミリ秒以上かかる処理)を回してしまうと、イベントループ全体がブロックされ、他のすべてのクライアントの応答が停止する。重い処理は `ext-parallel` や外部ワーカープロセス(RabbitMQやRedis Streams等)へオフロードせよ。

    PHPは、もはや「ダサいWebスクリプト言語」ではない。Zend VMの低レイヤの挙動とFiberの仕組みを完全に掌握したエンジニアが書いたコードは、Node.jsやGo言語のサーバーに匹敵する極めて高いスループットと美しさを両立する。
    さあ、あなたの次のプロジェクトで、この非同期並行処理の圧倒的なポテンシャルを証明して見せろ。

    タイトルとURLをコピーしました