リアクティブストリームとバックプレッシャー(Backpressure)制御
1. 概要
A. 定義
リアクティブストリーム(Reactive Streams)とは、非同期・ノンブロッキングなデータストリームを生産者と消費者の間で処理しつつ、消費者が処理できる速度でのみデータを流すようバックプレッシャー(Backpressure)信号を標準化した仕様であり、プログラミングモデルである。
リアクティブストリームの本質は「プッシュ(push)」と「プル(pull)」を一つのプロトコルに結合した点にある。
従来のイベントベース処理は、生産者がデータを消費者へ一方的に押し出すプッシュモデルであるため、消費者が遅いとバッファが溢れるかメッセージが失われる。
逆に、消費者が必要なときにデータを要求するプルモデルは安全であるが、往復遅延が大きくなり処理量が落ちる。
リアクティブストリームは、消費者が「いま n 個まで受け取れる」という需要(demand)を生産者へ伝える request(n) 信号を導入し、平時はプッシュの効率を保ちながら消費者が飽和すると自動的に流れが制御される動的なプッシュ・プルモデルを実現する。
ここでバックプレッシャーとは、消費者の処理能力が生産者の生産速度より低いとき、その不足分を上流へ戻して生産速度そのものを遅くする制御信号を意味する。 水道管で下流が詰まると上流の水圧が上がるように、ソフトウェアのパイプラインでも下流の遅延が上流へ伝わってこそ、システム全体が崩壊しない。 バックプレッシャーのないシステムは、負荷が集中するとキューが無限に大きくなりメモリ不足でプロセスが終了するか、キューを有限にすればメッセージが静かに捨てられるかという二つの失敗のいずれかに至る。
B. 登場背景と必要性
リアクティブストリームが標準として確立された2013〜2015年は、ストリーミングデータとマイクロサービスが急増した時期である。
RxJava、Akka Streams、Project Reactor など複数のライブラリがそれぞれの方式で非同期ストリームを扱う中で、異なるライブラリを組み合わせる際にバックプレッシャー信号が互換しない問題が繰り返された。
これを解決するために Netflix、Lightbend、Pivotal などが参加し、Publisher、Subscriber、Subscription、Processor の四つのインターフェースからなる最小仕様に合意し、これはのちに Java 9 の java.util.concurrent.Flow として標準ライブラリに取り込まれた。
必要性は、資源効率と安定性という二つの軸で説明される。 第一に、スレッド・パー・リクエスト(thread-per-request)モデルはリクエストごとにスレッドを占有するため、同時接続が数万件に達するとスレッドスタックメモリとコンテキストスイッチのコストが爆発する。 たとえばスレッド一つが1MBのスタックを使うと1万個のスレッドは10GBを消費し、I/O待機中のスレッドが大半であるため CPU は遊びメモリだけが枯渇する。 ノンブロッキングなリアクティブモデルは、少数のイベントループスレッド(通常は CPU コア数)で数万個の同時接続を処理し、メモリを大きく節約する。
第二に、分散パイプラインでは、ある区間の遅延が上流へ伝播しなければ障害が増幅される。 Kafka コンシューマがデータベース書き込みの遅延で遅くなったにもかかわらず、ブローカーから継続してデータを引き取ると、コンシューマのヒープが溢れる。 バックプレッシャーはこの状況で、コンシューマが処理可能な分だけポーリングするようにして、遅延をパイプライン全体で分担して吸収させる。 したがってリアクティブストリームは、単なる性能最適化ではなく、負荷急増時の回復力(resilience)設計の中核的手段である。
C. 特徴
リアクティブストリームの特徴は、非同期性、ノンブロッキング、バックプレッシャーに基づく流量制御、合成可能性(composability)に要約される。 演算子(operator)を関数型で連結してパイプラインを宣言的に記述し、各演算子は上流の需要と下流の需要を調整する役割も兼ねる。 次の表は従来モデルとの違いを整理したものだが、各項目の「なぜ」は続く節で散文として説明する。
| 区分 | ブロッキング同期モデル | 純粋プッシュイベントモデル | リアクティブストリーム |
|---|---|---|---|
| 流量制御 | 呼び出しスタックが自然なバックプレッシャー | なし(消失・OOM の危険) | request(n) による明示的バックプレッシャー |
| スレッド使用 | リクエストごとにスレッド占有 | イベントループ | イベントループ + ノンブロッキング |
| 遅延伝播 | 同期呼び出しで伝播 | 伝播しない | 需要信号で伝播 |
| 合成性 | 低い | 中程度 | 高い(演算子チェーン) |
2. リアクティブストリーム仕様と構成要素
A. 四つの中核インターフェース
リアクティブストリーム仕様は、Publisher、Subscriber、Subscription、Processor という四つのインターフェースで相互作用の規約を定義する。
Publisher はデータを生産する源であり、subscribe() 呼び出しで消費者を接続する。
Subscriber はデータを消費し、onSubscribe、onNext、onError、onComplete の四つのコールバックを持つ。
Subscription は生産者と消費者の間の一回の接続を表し、消費者が需要を要求する request(n) と接続を断つ cancel() を提供する。
Processor は Publisher でありながら同時に Subscriber でもある中間段階で、演算子チェーンの各ノードに相当する。
全体構造を概念図で表すと次のとおりである。
flowchart LR
P["Publisher(生産者)"] -->|onSubscribe| SUB["Subscription(購読)"]
SUB -->|request n| P
P -->|onNext データ| PR["Processor(演算子チェーン)"]
PR -->|onNext 変換| S["Subscriber(消費者)"]
S -->|request n 需要| PR
PR -->|request n 再調整| P
S -.->|onError / onComplete| END["終了信号"]
この構造の核心は、データが左から右へ流れる間、需要信号(request(n))は右から左へ逆向きに流れる点にある。
消費者が自らが処理できる個数だけ要求すると、その信号が演算子を経て生産者まで伝わり、生産者は要求された個数以下でのみ onNext を放出する。
この規約のおかげで、生産者は消費者の処理能力を知らなくても、決して要求量を超えてデータを押し出さない。
B. 信号規約と不変条件
仕様は、信号の順序と個数に関する厳格な不変条件を要求する。
onSubscribe は正確に一度だけ先に呼び出されなければならず、onNext は要求された数を超えてはならず、onError や onComplete の後にはいかなる信号も発生してはならない。
これらの規約は異なるライブラリを安全に連結するための最小契約であり、TCK(Technology Compatibility Kit)という標準テストスイートで遵守可否を検証する。
規約を守らない実装は合成時に競合状態や資源リークを起こすため、自分で Publisher を実装するより、検証済みライブラリのファクトリメソッドを用いるのが実務上の推奨事項である。
C. コールドストリームとホットストリーム
ストリームは購読時点に応じてコールド(cold)とホット(hot)に分かれる。 コールドストリームは購読者が付くたびにデータ生産が最初から新たに始まり、HTTP リクエスト・ファイル読み込み・データベース照会のようにリクエスト単位で完結する作業に適する。 ホットストリームは購読者の有無と無関係にデータが流れ続け、センサーイベント・株価・ユーザークリックのようなリアルタイム放送的な源に該当する。 ホットストリームはバックプレッシャーをそのまま適用しにくいが、これは生産源が消費者の需要を待ってくれないためである。 この場合、後述する緩衝・標本抽出・最新値保持のような消失許容戦略を併せて設計しなければならない。
3. バックプレッシャー制御メカニズムと動作手順
A. 需要ベースの流量制御手順
バックプレッシャーの標準メカニズムは、消費者が処理余力の分だけ request(n) を呼び出し、生産者はその範囲内でのみ放出する往復の対話である。
消費者は一度に大きな需要を要求してから消尽するたびに補充するか(例:256個要求後に半分消尽時に128個追加)、正確に一つずつ要求して処理速度に合わせる方式を選べる。
次のシーケンス図は、消費者が遅くなるとき需要がどのように調整されるかを示す。
sequenceDiagram
participant S as Subscriber(消費者)
participant O as Operator(バッファ)
participant P as Publisher(生産者)
S->>O: request(256)
O->>P: request(256)
P-->>O: onNext x256
O-->>S: onNext x256
Note over S: 処理遅延が発生
S->>O: request(64) のみ補充
O->>P: バッファ余裕分のみ request
Note over P: 放出速度が自動的に減少
S->>O: cancel() または完了
この手順で注目すべき点は、消費者が遅延すると補充要求が遅れ、その結果として生産者の放出も自然に遅くなることである。
いかなる明示的な「速度制限命令」もなく、ただ需要信号の遅延だけで、パイプライン全体の速度が消費者に合わせられる。
実務では Project Reactor の Flux が既定で256個の大きさのプリフェッチバッファを置き、75%が消尽されると次のバッチを要求する形でこの対話を自動化する。
B. バックプレッシャーがないときの崩壊シナリオ
バックプレッシャーが適切に伝播されないと、システムは二つの方向のいずれかへ崩れる。
無限バッファを使うと、負荷急増時にキューが増え続けてヒープが枯渇し、OutOfMemoryError でプロセスが終了する。
有限バッファを使いつつ溢れたときに例外を投げると、MissingBackpressureException のようなエラーで要求が失敗する。
2020年代初頭の複数のストリーミングサービスの事後分析で、消費遅延が上流へ伝播できずブローカー・コンシューマ間のキューが急増した事例が繰り返し報告された。
したがってバックプレッシャー設計は、「通常負荷でうまく回るか」ではなく「最大負荷が消費能力を超えるときにいかに優雅に劣化するか」を基準に検証しなければならない。
C. 消失許容戦略と緩衝
消費者を無制限に待てないホットストリームでは、データの一部を意図的に捨てるか要約するオーバーフロー戦略が必要である。 代表的に、バッファ(buffer)は一定サイズまで貯めてからバッチで渡し、ドロップ(drop)は消費者が忙しいとき新しいデータを捨て、最新値(latest)は最も最近の値のみ保持して以前の値を上書きする。 標本抽出(sample)やウィンドウ(window)は時間・個数基準でデータをまとめ、下流負荷を下げる。 次の表は戦略ごとの特性を比較する。
| 戦略 | 動作 | データ消失 | 適する状況 |
|---|---|---|---|
| buffer | 有限キューに積載後バッチ | 溢れるとエラー/遮断 | 短い瞬間の負荷、消失不可 |
| drop | 超過分を廃棄 | あり | 最新性より順序が重要なとき |
| latest | 最新1件のみ保持 | あり | ダッシュボード・相場表示 |
| error | 即時失敗信号 | - | 速い失敗が良いバッチ |
戦略選択は業務要件から出発する。 決済イベントのように一件も消失してはならないデータは buffer + 永続キュー + 再処理で無損失を保証しなければならず、リアルタイム温度ゲージのように最新値のみ意味のあるデータは latest で古い値を捨てるほうが資源効率とユーザー体験の双方に有利である。
4. 比較と実務適用事例
A. 命令型モデルとのトレードオフ
リアクティブモデルは高い処理量と資源効率を与えるが、デバッグと学習曲線という代価を伴う。 非同期パイプラインでは、スタックトレースが実際の実行経路を捉えられずエラー原因の追跡が難しく、ブロッキング呼び出し一つがイベントループを塞ぐと処理量全体が急落する。 一方、スレッド・パー・リクエストモデルはコードが直感的でデバッグが容易であり、近年は Java の仮想スレッド(Virtual Thread, Project Loom)が登場して、ブロッキングコードでも高い並行性を得られるようになった。 したがって「無条件にリアクティブ」ではなく、I/O バウンドでありストリーミング・バックプレッシャー制御が本質的に必要な区間へ選択的に適用する判断が重要である。
B. Web サービス事例 — Spring WebFlux
ある API ゲートウェイが毎秒数万件のリクエストを少数のイベントループで処理しなければならない状況を想定する。
Spring MVC(スレッド・パー・リクエスト)では Tomcat スレッドプール200個が外部 API 応答を待って全て占有され、実際の CPU 使用率が20%であってもリクエストがキューで待機してタイムアウトする現象が発生する。
Spring WebFlux(Reactor Netty ベース)へ切り替えると、CPU コア数分のイベントループがノンブロッキングで数万接続を扱い、下流 API が遅くなると request(n) が上流クライアントへ伝播して自然に流入を遅くする。
実測事例では、同一ハードウェアで同時接続の収容量が数倍増加し、P99 遅延が安定する結果が報告される。
C. データパイプライン事例 — Kafka コンシューマ
Kafka ベースのイベントパイプラインで、コンシューマがダウンストリームのデータベースへ毎秒5,000件を書けるのに、トピックには毎秒20,000件が流入するとしよう。
バックプレッシャーなしに無制限ポーリングすると、コンシューマのヒープに未処理レコードが積もり、GC 暴走と OOM につながる。
Reactor Kafka や Akka Streams の Kafka コネクタは、ダウンストリームのシンクの需要に合わせて poll 呼び出しとオフセットコミットの速度を調整し、パーティションごとのプリフェッチを制限してコンシューマが処理可能な分だけ引き取る。
その結果、瞬間流入が消費能力を超えても、ブローカーにデータが安全に保管されたままコンシューマは一定の速度で追いつき、パイプライン全体が流入の急増を緩衝する。
5. 深化 — 最新動向と類似技術の連携
リアクティブストリーム標準は Java 9 の Flow API へ取り込まれて言語標準の一部となり、R2DBC(リアクティブ関係データベースアクセス)やリアクティブ gRPC のように、バックプレッシャーがスタック全層へ拡張されている。
データベースドライバがバックプレッシャーを支援しなければ、パイプラインの一点がブロッキングされて全体の利点が崩れるため、「エンドツーエンド(end-to-end)ノンブロッキング」が実務の中核的関心事となった。
一方、Java 21 で正式化された仮想スレッドは、リアクティブモデルの代替として浮上した。 仮想スレッドはブロッキングのように見えるコードをノンブロッキングで実行し、リアクティブの資源効率を命令型コードで得させるが、仮想スレッド自体はバックプレッシャーを提供しない点が重要である。 すなわち並行性の問題(スレッド不足)は仮想スレッドが解くが、流量制御の問題(生産・消費の速度不一致)は依然としてリアクティブストリームやセマフォ・制限キューのような別途のメカニズムが必要である。 したがって二つの技術は競争ではなく、「仮想スレッドで単純な I/O を、リアクティブストリームでバックプレッシャーが必要なストリーミングを」扱う相補関係へ定着する趨勢である。
バックプレッシャーの概念はリアクティブストリームに限らない。 TCP の受信ウィンドウ、gRPC の HTTP/2 フロー制御、Kafka のコンシューマラグに基づく調整、サーキットブレーカの負荷遮断は、すべて広い意味のバックプレッシャー・流量制御系列に属する。 技術士の観点からこれらを一つの「流量制御」系譜としてまとめて説明すれば、個別技術の暗記ではなく原理中心の答案を構成できる。
6. 考慮事項および示唆点
第一に、適用範囲を冷静に判断しなければならない。 リアクティブストリームは I/O バウンド・高並行・ストリーミング区間で大きな利点を与えるが、CPU バウンドのバッチや単純な CRUD には複雑さだけを増やす。 仮想スレッドが成熟した今は、「ブロッキングコードが問題か、流量制御が問題か」をまず区分し、後者のときのみリアクティブを選択する意思決定基準が必要である。
第二に、エンドツーエンドのノンブロッキングを設計原則とすべきである。 パイプラインのどこか一点でもブロッキング呼び出し(JDBC、同期ファイル I/O など)がイベントループで実行されると、少数のイベントループが塞がり処理量全体が崩壊する。 避けられないブロッキング作業は必ず別途の境界スレッドプール(bounded elastic scheduler)へ隔離し、データアクセス層は R2DBC のようなノンブロッキングドライバへ統一する方向が望ましい。
第三に、オーバーフロー戦略を業務のリスク度に合わせて明示的に選択すべきである。 データ消失が許されないドメインは無損失バッファと永続キュー・再処理を結合し、最新性のみ重要なドメインは消失許容戦略で資源を節約するといった形で、オーバーフロー方針自体をアーキテクチャ決定記録(ADR)として残し、運用者が理解できるようにすべきである。
第四に、観測性とテストを併せて設計すべきである。 バックプレッシャーは通常負荷では見えず最大負荷で表れるため、プリフェッチサイズ・バッファ占有率・要求に対する放出量・コンシューマラグのような指標を常時収集し、負荷テスト・障害注入で「消費能力を超えたときの劣化曲線」を事前に確認すべきである。 非同期スタックトレースを補完するためのチェックポイント・コンテキスト伝播も運用準備の一部である。
第五に、組織の力量と標準化を並行すべきである。
リアクティブコードは学習曲線が急で、誤って使うとかえって性能が悪くなるため、検証済みライブラリの慣用句をチーム標準として定め、自分で Publisher を実装するよりファクトリ・演算子を組み合わせる規則をコードレビューに含めるのが安全である。
参考資料
- Reactive Streams Specification: https://www.reactive-streams.org/
- Project Reactor Reference — Handling Backpressure: https://projectreactor.io/docs/core/release/reference/
- Spring Framework, Web on Reactive Stack (WebFlux): https://docs.spring.io/spring-framework/reference/web/webflux.html
- OpenJDK, JEP 444: Virtual Threads: https://openjdk.org/jeps/444
一言まとめ: リアクティブストリームは、
request(n)に基づくバックプレッシャー信号を上流へ戻して消費者の処理能力に合わせ生産速度を制御することで、ノンブロッキングの資源効率と負荷急増時の回復力を同時に達成する非同期ストリーム処理の標準である。