【実務・中級編】Fiberを用いたリアクティブプログラミングフレームワークの構築:イベントストリームとオペレーター – PHPコア・内部エンジンと高速化・並行処理の極意解析バイブル

PHPを掌握する極限の知見:Fiberで構築するリアクティブ・イベントストリームの真髄

PHPにおける非同期処理の歴史は、長らく「苦難の歴史」であった。`pcntl_fork`によるプロセス並行処理はメモリフットプリントを肥大化させ、`ext-async`やAmp、ReactPHPなどのイベントループは、コールバック地獄(Pyramid of Doom)やジェネレータ(Generator)の煩雑な委譲構文という悪夢を伴っていた。

だが、PHP 8.1でFiber(ファイバー)が導入されたことで、言語のパラダイムは決定的に変わった。スタックフルコルーチンを手に入れたPHPは、コールバックを排除し、同期コードの見た目を維持したまま、完全に非同期なリアクティブ・プログラミングを実行できる。

今回は、Zend VMのコールスタックの挙動を脳内に描きながら、Fiberを基盤としたRx(Reactive Extensions)スタイルのイベントストリームおよび、実務に耐えうるリアクティブフレームワークの核を実装・解説する。

—

1. なぜ「コールバック」ではなく「Fiber」なのか:内部挙動の真実

従来のイベントループは、I/O待ちが発生するたびに処理を中断し、再開時のコールバックを登録していた。これにより、ビジネスロジックが細切れになり、スタックトレースの追跡が極めて困難になった。

一方、Fiberは独自のコールスタック(Cレベルのzend_execute_dataチェーン)をヒープ上に保持する。
`Fiber::suspend()`が呼ばれた瞬間、Zend VMは現在の実行コンテキストを退避させ、親スコープ(イベントループ)へ制御を返す。そして`Fiber::resume()`が呼ばれた時、中断された正確なバイトコードの位置から実行が再開される。

このメカニズムを利用すれば、無限に続くイベントストリーム(Stream)の各要素の処理を、あたかも上から下に流れる同期処理のように記述できるのだ。

—

2. 設計の中核:Observable、Observer、そしてFiber Scheduler

リアクティブプログラミングの基本構成要素は以下の3つである。
1. Observable(ストリーム): データの発生源であり、オペレーターを通じて変換される。
2. Observer(オブザーバー): データのコンシューマー(`onNext`, `onError`, `onCompleted`)。
3. Scheduler / EventLoop: Fiberの切り替えとI/O多重化(`stream_select`等)を司る心臓部。

今回は、余計なサードパーティ製ライブラリに依存せず、PHP 8.2+の厳格な型システムとメモリ効率を意識した実用的なミニマム・フレームワークを構築する。

実装:堅牢なリアクティブ・コアエンジン

以下のコードは、イベントループを内包し、Fiber上で動的にストリームを処理するリアクティブエンジンの実例である。

  • オブザーバーのインターフェース
  • /
    interface ObserverInterface
    {
    public function onNext(mixed $value): void;
    public function onError(Throwable $e): void;
    public function onCompleted(): void;
    }

    /

    • ストリームの購読を管理するサブスクリプション

    /
    class Subscription
    {
    private bool $isUnsubscribed = false;

    public function __construct(private readonly \Closure $unsubLogic) {}

    public function unsubscribe(): void
    {
    if (!$this->isUnsubscribed) {
    $this->isUnsubscribed = true;
    ($this->unsubLogic)();
    }
    }

    public function isUnsubscribed(): bool
    {
    return $this->isUnsubscribed;
    }
    }

    /

    • イベントストリームの基底クラス Observable

    /
    class Observable
    {
    private readonly \Closure $subscribeLogic;

    public function __construct(callable $subscribeLogic)
    {
    // $subscribeLogic は (ObserverInterface $observer): Subscription を受け取る
    $this->subscribeLogic = \Closure::fromCallable($subscribeLogic);
    }

    public function subscribe(ObserverInterface $observer): Subscription
    {
    return ($this->subscribeLogic)($observer);
    }

    /

    • map オペレーター:流れてくるデータを変換する

    /
    public function map(callable $mapper): self
    {
    return new self(function (ObserverInterface $observer) use ($mapper) {
    return $this->subscribe(new class($observer, $mapper) implements ObserverInterface {
    public function __construct(
    private readonly ObserverInterface $destination,
    private readonly $mapper
    ) {}

    public function onNext(mixed $value): void
    {
    try {
    // マッパー関数を適用して下流へ流す
    $mapped = ($this->mapper)($value);
    $this->destination->onNext($mapped);
    } catch (Throwable $e) {
    $this->onError($e);
    }
    }

    public function onError(Throwable $e): void
    {
    $this->destination->onError($e);
    }

    public function onCompleted(): void
    {
    $this->destination->onCompleted();
    }
    });
    });
    }

    /

    • filter オペレーター:条件に合致するデータのみを通す

    /
    public function filter(callable $predicate): self
    {
    return new self(function (ObserverInterface $observer) use ($predicate) {
    return $this->subscribe(new class($observer, $predicate) implements ObserverInterface {
    public function __construct(
    private readonly ObserverInterface $destination,
    private readonly $predicate
    ) {}

    public function onNext(mixed $value): void
    {
    try {
    if (($this->predicate)($value)) {
    $this->destination->onNext($value);
    }
    } catch (Throwable $e) {
    $this->onError($e);
    }
    }

    public function onError(Throwable $e): void
    {
    $this->destination->onError($e);
    }

    public function onCompleted(): void
    {
    $this->destination->onCompleted();
    }
    });
    });
    }

    /

    • Fiberを活用した非同期遅延(delay)オペレーター

    /
    public function delay(int $milliseconds): self
    {
    return new self(function (ObserverInterface $observer) use ($milliseconds) {
    return $this->subscribe(new class($observer, $milliseconds) implements ObserverInterface {
    public function __construct(
    private readonly ObserverInterface $destination,
    private readonly int $milliseconds
    ) {}

    public function onNext(mixed $value): void
    {
    // 現在実行中のFiberを一時停止し、指定時間後に再開するスケジューリング
    $fiber = Fiber::getCurrent();
    if ($fiber === null) {
    // Fiber外から呼ばれた場合は同期的に処理(フォールバック)
    usleep($this->milliseconds 1000);
    $this->destination->onNext($value);
    return;
    }

    // イベントループ側のタイマーキューにタスクを登録する想定の疑似遅延
    // 実務ではイベントループの timer API と連携させる
    usleep($this->milliseconds 1000);
    $this->destination->onNext($value);
    }

    public function onError(Throwable $e): void
    {
    $this->destination->onError($e);
    }

    public function onCompleted(): void
    {
    $this->destination->onCompleted();
    }
    });
    });
    }
    }

    —

    3. コードレビュー:なぜこの設計が実務で安全なのか

    上記のコードにおいて、テクニカルリードとして以下のアーキテクチャ上のこだわりと安全性を指摘しておかねばならない。

    1. 例外伝播の堅牢性 (`try-catch` の境界)

    オペレーター(`map`や`filter`)の内部でユーザー定義のクロージャが例外を投げた場合、それがそのままイベントループの根幹をクラッシュさせてはならない。上記のコードでは、各オペレーター層で`try-catch`を完全にカプセル化し、例外が発生した場合は即座に`onError`チャネルへルーティングしている。これにより、ストリーム全体が安全に破棄(Teardown)される。

    2. メモリリーク(循環参照)の回避

    リアクティブプログラミングにおいて最も危険なのは、`Observable` と `Observer` 間でクロージャが互いを参照し合うことによる循環参照(Circular Reference)である。PHPの参照カウント方式(RC)において、クロージャが `$this` や外部変数を暗黙的にキャプチャすると、ガベージコレクタ(GC)が走るまでメモリが解放されず、長期間稼働するAPIサーバーや常駐プロセスでは致命的なメモリリーク(OOM)を引き起こす。
    これを防ぐため、明示的に必要な変数のみを `use` でキャプチャし、不要な参照を作らない設計を徹底している。

    —

    4. 実戦投入:Fiberベースの非同期イベントストリームの駆動

    では、このフレームワークを実際のWebアプリケーションやAPIバックエンドの文脈でどのように動かすのか。以下に、複数の非同期タスクをFiberで並行実行しながら、イベントストリームを消費する統合サンプルを示す。

    onNext($item);
    }
    $observer->onCompleted();
    });

    // 初回起動
    $fiber->start();

    return new Subscription(function () use ($fiber) {
    // キャンセル時のクリーンアップ処理
    echo “[System] ストリームがキャンセルされました。\n”;
    });
    });

    // 2. パイプラインの構築 (Rxスタイル)
    // 偶数のみをフィルタリング -> 10倍に変換 -> 100ms遅延
    $pipeline = $source
    ->filter(fn(int $val) => $val % 2 === 0)
    ->map(fn(int $val) => $val 10)
    ->delay(100);

    // 3. 購読(Subscribe)と実行
    echo “[System] イベントストリームの処理を開始します。\n”;

    $subscription = $pipeline->subscribe(new class implements ObserverInterface {
    public function onNext(mixed $value): void
    {
    echo “[Observer] 受信データ: {$value} (時刻: ” . microtime(true) . “)\n”;
    }

    public function onError(Throwable $e): void
    {
    echo “[Observer] エラー発生: ” . $e->getMessage() . “\n”;
    }

    public function onCompleted(): void
    {
    echo “[Observer] すべてのイベント処理が完了しました。\n”;
    }
    });

    実行結果のイメージ

    [System] イベントストリームの処理を開始します。
    [Observer] 受信データ: 20 (時刻: 1680000000.1234)
    [Observer] 受信データ: 40 (時刻: 1680000000.2345)
    [Observer] 受信データ: 60 (時刻: 1680000000.3456)
    [Observer] 受信データ: 80 (時刻: 1680000000.4567)
    [Observer] 受信データ: 100 (時刻: 1680000000.5678)
    [Observer] すべてのイベント処理が完了しました。

    —

    5. アーキテクトからの最終提言

    PHPにおけるFiberの導入は、単に「コードが綺麗にかける」というレベルのメリットにとどまらない。従来の「1リクエスト=1プロセス(またはスレッド)」というPHPの古典的な制約を脱し、単一のPHPプロセス内で数千の同時リクエストや非同期タスクを効率的にハンドリングするための強力な武器となる。

    しかし、強大な力には相応の責任が伴う。Fiber内部でブロッキングなI/O(通常の同期PDO接続や外部API同期コールなど)を実行してしまえば、イベントループ全体がフリーズし、すべてのFiberがブロックされる。非同期・リアクティブな世界に踏み込むのであれば、ドライバ層に至るまでノンブロッキング(`ext-uv`や非同期Streams)で統一されているか、コードレビュー時に厳しく監査されなければならない。

    Zend VMのメモリ構造と、コンテキストスイッチのコストを正確に理解した上でFiberを使いこなせた時、あなたのPHPアプリケーションは、Node.jsやGoにも引けを取らない高スループットな非同期エンジンへと生まれ変わる。その設計の舵取りを、あなたに託す。

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