ストリーム処理のウィンドウイングと時間意味論(イベント時間・ウォーターマーク・厳密に一度)
1. 概要
A. 定義
ストリーム処理(Stream Processing)とは、開始と終了が定まらない無限(unbounded)のデータが到着した瞬間に連続的に変換・集計・分析し、低遅延で結果を算出するデータ処理パラダイムである。
ウィンドウイング(Windowing)とは、無限ストリームを有限の処理単位に切るために、時間・件数・セッションなどの基準でイベントをまとめる境界設定の技法であり、時間意味論(Time Semantics)とは、イベントが「実際に発生した時刻(event time)」と「エンジンが処理した時刻(processing time)」を区別して集計の正確性を保証するモデルである。
ストリーム処理の本質的な難しさは「終わらないデータに対して、いつ・何を・どれだけ正確に計算を確定するか」という問いにある。バッチ処理は入力がすべて揃った後に計算するため境界が自明であるが、ストリームはデータが絶え間なく流れ、ネットワーク遅延によって順序が入れ替わって(out-of-order)到着するため、集計を確定する時点を自ら定義しなければならない。ウィンドウイングは「何を一つのまとまりと見るか(where in event time)」を、ウォーターマークは「いつそのまとまりを閉じて結果を出すか(when in processing time)」を決める。技術士の観点では、この主題は単なるAPIの利用ではなく、遅延・正確性・完全性(completeness)のトレードオフを設計するリアルタイムデータパイプラインのガバナンス問題として理解すべきである。
本質的に、ストリーム処理は「データが完全に揃うまで待てば遅く、待たなければ不完全である」という矛盾を扱う。バッチが完全性のために遅延を甘受し、純粋なリアルタイムが遅延をなくすために完全性を手放すのに対し、現代のストリーム処理はウィンドウ・ウォーターマーク・トリガーという三つのつまみでその中間点をパラメータ化し、同じコードで「速いが近似の結果」と「遅いが正確な結果」の間を業務要求に合わせて調整できる点に、その特徴がある。
B. 登場の背景と必要性
第一に、ビジネス意思決定の時間的地平が「翌日」から「今この瞬間」へ移った。不正取引検知(FDS)、リアルタイム推薦、モノのインターネット(IoT)設備の異常検知、急変する在庫・相場の反映は、夜間バッチでは価値を失う。数秒〜数分の遅延の内に判断を下さねばならないため、データが積み上がるのを待たず流れる最中に処理する構造が必要になった。
第二に、分散収集環境ではデータは本質的に遅れ、入れ替わって到着する。モバイル端末がオフラインから復帰して数分前のイベントを遅れて送信したり、パーティションごとの処理速度差で順序が逆転したりする。処理時刻を基準に集計すると「10時に発生したが10時5分に到着した注文」が誤ったウィンドウに入り、結果が歪む。これを正すには、イベント自身が持つ発生時刻を基準とするイベント時間処理と、遅延を許容しつつ限界を引くウォーターマークの概念が必要であった。
第三に、リアルタイムだからといって正確性を手放すことはできない。障害でノードが落ちたり再起動したりする際にイベントが欠落・重複集計されると、課金・精算・規制報告のような領域では致命的である。したがって、障害が起きても各イベントが結果に厳密に一度だけ反映される厳密に一度(exactly-once)の処理保証が、ストリームエンジンの中核要件として定着した。
第四に、運用・コスト面の要求も高まった。24時間止まらず回るパイプラインは状態が無限に肥大しないよう管理せねばならず、ロジックを変えれば過去データを流し直して結果を再生成(reprocessing)できねばならない。また瞬間的な急増(トラフィックスパイク)時に上流へ負荷を伝播して自らを守る逆圧(backpressure)制御がなければ、メモリ枯渇で全体が崩壊する。このようにストリーム処理は、正確性のみならず持続的な運用性までを同時に求める技術へと進化した。
C. 処理保証水準のスペクトラム
ストリーム処理の信頼性は「障害時に各イベントが結果へ何回反映されるか」で規定される。最小保証のat-most-onceは再送をしないため欠落を許容する代わりに重複はない。at-least-onceは障害時の再処理で欠落を防ぐが、同じイベントが二度反映されうるため、合計・カウントが過大計上される。exactly-onceは状態(state)と出力が障害の前後で一貫して厳密に一度だけ反映されることを保証する。ただしこれはネットワーク転送そのものを物理的に一度だけ行うのではなく、チェックポイントとトランザクションで「効果(effect)が一度」になるようにするものであるため、effectively-onceとも呼ばれる。この区別は、後の深掘りの節で扱うチェックポイント・トランザクションの機構で実現される。
保証水準の選択はただではない。上位の保証ほど状態スナップショット、トランザクション調停、遅延コミットなどのオーバーヘッドが大きくなり遅延が増す。したがって、すべてのパイプラインを一律にexactly-onceで設計するより、課金・精算・規制報告のように正確性が金銭・法的責任に直結する経路のみexactly-onceとし、ダッシュボード指標や推薦シグナルのように少量の重複・欠落が許される経路はat-least-onceとしてコストを最適化する選択的適用が実務の定石である。
2. 全体構造と時間意味論
A. ストリーム処理パイプラインの構造
リアルタイムパイプラインは一般に、不変ログに基づくメッセージブローカーを中心に、状態を持つ処理オペレータとチェックポイントストア、そして冪等・トランザクションのシンクで構成される。
flowchart LR
SRC["イベントソース(アプリ・IoT・ログ)"] --> BROKER["メッセージブローカー(不変ログ・パーティション)"]
BROKER --> OP["ストリームオペレータ(ウィンドウ・集計)"]
OP --> ST["状態ストア(Keyed State)"]
OP --> CP["チェックポイント(スナップショット)"]
CP --> DFS["永続ストア(分散ファイルシステム)"]
OP --> SINK["シンク(冪等・トランザクション)"]
SINK --> SERVE["サービング(DB・ダッシュボード・通知)"]
メッセージブローカー(例: Apache Kafka)はパーティション単位で順序が保証される不変ログを提供し、障害時にオフセットを巻き戻して再生(replay)できる再処理の基盤となる。保存期間(retention)を例えば7日とすれば、その期間内の任意の時点へ消費位置を戻してロジックを再適用でき、パーティション数は並列処理単位を決めてスループット拡張の基準となる。ストリームオペレータはキーごとに状態を保ちつつウィンドウ集計を行い、定期的に状態スナップショットをチェックポイントとして取り分散ファイルシステムへ保存する。障害が起きると直近のチェックポイントから状態を復元し、その地点のオフセットから再び消費するため、at-least-onceの再処理が基本となる。ここに冪等またはトランザクションのシンクを組み合わせると重複反映まで除かれ、exactly-onceが完成する。この構造ではチェックポイント周期が復旧地点の密さを、状態バックエンドの性能が集計スループットを左右するため、両者を併せて調整することが運用設計の出発点となる。
B. イベント時間・処理時間・ウォーターマーク
時間意味論の出発点は三つの時刻を区別することである。イベント時間はイベントが実際に発生した時刻(例: 決済承認時刻)で、データ内にタイムスタンプとして埋め込まれている。取込時間(ingestion time)はブローカー・エンジンに入った時刻、処理時間はオペレータが実際に処理した時刻である。結果の再現性と正確性はイベント時間基準の処理から生まれる。同じ入力を流し直すと(再処理)処理時間は毎回変わるが、イベント時間基準の集計は常に同じ結果を与えるからである。
問題は、イベント時間でウィンドウを閉じるには「そのウィンドウに属するイベントがもう来ない」という判断が要る点である。これを推定する仕掛けがウォーターマークである。ウォーターマークW(t)は「イベント時間t以前のデータは(おおむね)すべて到着した」というエンジンの進行表示子であり、このウォーターマークがウィンドウ終了時刻を越えた瞬間に当該ウィンドウの集計を確定・発火(trigger)する。実務では、観測された最大イベント時間から許容遅延を引いた有界無秩序(bounded out-of-orderness)のウォーターマークがよく使われる。例えば遅延許容を5秒とすればW = max_event_time - 5sとなる。
sequenceDiagram
participant E as イベントストリーム
participant W as ウォーターマーク生成器
participant WIN as ウィンドウオペレータ
participant OUT as 結果シンク
E->>W: イベント到着(イベント時間タイムスタンプを含む)
W->>WIN: ウォーターマーク伝播("t以前の到着完了")
WIN->>WIN: ウィンドウ境界とウォーターマークを比較
alt ウォーターマーク > ウィンドウ終了
WIN->>OUT: ウィンドウ集計を確定・発火
else 遅延データ到着
WIN->>WIN: 許容遅延(allowed lateness)内なら更新
WIN->>OUT: サイド出力(side output)として分離
end
なお取込時間はイベント時間と処理時間の折衷であり、ソースに信頼できるタイムスタンプがない場合にブローカー進入時刻を基準として、ある程度の順序安定性を得る実用的な代替である。ただし再処理時に結果が変わりうるため、完全な再現性は保証しない。
ウォーターマークは本質的に完全性と遅延のトレードオフを調停するつまみである。許容遅延を大きく取れば遅れて来るデータまで含めて結果が正確になるが、ウィンドウ確定が遅れて遅延が増す。逆に小さく取れば速いが、遅刻データを取りこぼして結果が不完全になる。そこで多くのエンジンは、ウィンドウ確定後も許容遅延(allowed lateness)の期間は遅刻データで結果を更新し、それより遅いデータはサイド出力として分離し別の補正経路で処理する。
C. ウィンドウの種類
ウィンドウは無限ストリームをどの基準でまとめるかによって分かれる。種類の選択は問いの性質(周期的な集計か、活動区間の分析か)によって異なり、下表は補助的な要約にすぎず、各種類の選択理由は文章で説明する。
| ウィンドウ種類 | 定義 | 特徴 | 代表的用途 |
|---|---|---|---|
| タンブリング(Tumbling) | 固定サイズ・非重複 | 境界が重ならず各イベントが正確に1つのウィンドウ | 1分の売上集計、時間別レポート |
| スライディング(Sliding) | 固定サイズ・一定間隔で移動 | ウィンドウが重なりイベントが複数ウィンドウに属する | 直近5分の移動平均(1分ごと) |
| セッション(Session) | 活動間の空白(gap)基準 | サイズ可変、活動区間を自動分離 | ユーザーセッション、設備の稼働区間 |
| グローバル(Global) | 境界なし・カスタムトリガー | ユーザーがトリガーを定義 | 件数ベースの集計 |
タンブリングウィンドウは1分、1時間のように固定長で重ならないように切り、各イベントは正確に一つのウィンドウにのみ属する。「毎1分の取引件数」のような明確な周期レポートに適する。境界が単純なため状態管理コストが最も低く、ウィンドウが閉じればその状態を直ちに破棄でき、メモリが予測しやすい。一方で境界が固定であるため、閾値現象が二つのウィンドウにまたがって分かれると各々は閾値未満となり検知を取りこぼしうる点を設計時に考慮せねばならない。
スライディングウィンドウはサイズ(例: 5分)と移動間隔(例: 1分)を別に置いてウィンドウが重なるため、「1分ごとに更新される直近5分の移動平均」のような連続的な傾向監視に使う。一つのイベントが(サイズ÷間隔)個のウィンドウに重複して含まれるため、状態サイズと計算量がその分大きくなる。上の例では一つのイベントが5個のウィンドウに属するため、単純なタンブリングに比べ約5倍の状態・計算負荷が生じる。傾向の感度(間隔を短く)と資源コストの間の折衷が核心の設計ポイントである。
セッションウィンドウは固定境界ではなく活動の間の空白(gap)で区間を分ける。例えば空白閾値を30分とすれば、ユーザーの連続クリックが30分以上途切れるとセッションが終了し新たなセッションが始まる。ユーザー行動分析や設備稼働セッションの抽出に自然である。ウィンドウ長がデータに応じて動的に決まるため、遅刻イベントが二つのセッション間の空白を埋めると、既に発火した二つのセッションを一つに併合(merge)せねばならない複雑さが生じる。このためセッションウィンドウは状態・トリガー設計が最も厄介な種類である。
グローバルウィンドウは時間境界なしに「100件ごと」のようなユーザー定義トリガーで発火する。時間ベースの集計で表しにくい件数・閾値条件ベースの処理に柔軟に対応できるが、トリガーと状態整理を全面的にユーザーが担うため、濫用すると状態リークとメモリ増加の原因となる。
D. トリガーと累積モード
同じウィンドウでも「いつ、何回結果を出すか」はトリガー(Trigger)で、「再発火時に以前の結果とどう合わせるか」は累積モード(accumulation mode)で分離して設計する。遅刻データが重要な環境では、ウォーターマーク到達時に一次結果を速く出し(speculative)、その後遅刻分が来るとウィンドウを再発火して結果を訂正する方式が有用である。このとき累積(accumulating)モードは毎回全体の再計算値を送り出し、破棄後累積(accumulating & retracting)モードは以前の値を取り消す補正レコードと新たな値を併せて送り出し、下流が二重集計しないようにする。この「早期結果 + 事後訂正」のパターンはDataflowモデルの核心的な貢献であり、完全性を犠牲にせず低遅延を併せて得る鍵である。
3. 比較と適用事例
A. 代表的エンジンの比較と差異の理由
| 区分 | Apache Flink | Kafka Streams | Spark Structured Streaming |
|---|---|---|---|
| 処理モデル | 真のレコード単位ストリーミング | レコード単位(ライブラリ) | マイクロバッチ(既定)・連続(実験的) |
| 状態・チェックポイント | 非同期バリアスナップショット(ABS) | チェンジログトピック + RocksDB | チェックポイント + WAL |
| exactly-once | 状態+トランザクションシンクで保証 | exactly_once_v2(Kafka間) |
冪等・トランザクションシンクに基づく |
| 遅延特性 | ミリ秒級の低遅延 | ミリ秒級(アプリ内蔵) | バッチ周期の分だけ(数百ms〜秒) |
三つのエンジンの差異は「ストリーミングを何と見るか」という根本的な観点に由来する。Flinkはストリームを一級市民と見てレコード単位で連続処理するため、低遅延と精緻なイベント時間・ウォーターマーク制御に強い。Spark Structured Streamingは伝統的にストリームを小さなバッチの連続として扱うマイクロバッチモデルであり、スループットとバッチエコシステムの統合に有利な代わりにバッチ周期の分だけの遅延がある(連続処理モードは遅延を下げるが保証水準に制約がある)。Kafka Streamsは別クラスタなしにアプリケーションへ内蔵されるライブラリであるため運用が単純で、状態をチェンジログトピックへ複製して復元し、Kafka内部の入出力についてexactly-onceを提供する。すなわち、低遅延・精緻な時間制御が核心ならFlink、Kafka中心の軽量な内蔵処理ならKafka Streams、バッチ・MLエコシステムとの統合が重要ならSparkが合理的な選択となる。
B. 具体的な適用事例
リアルタイム不正取引検知(FDS)を見よう。カード承認イベントをカード番号キーでまとめ、1分のタンブリングウィンドウで承認件数・金額を集計し、閾値を超えると遮断通知を発火する。モバイル決済の特性上、一部のイベントが数秒遅れて到着するためウォーターマーク遅延を5秒として遅刻分まで含める一方、それより遅いイベントはサイド出力へ送り事後精算に反映する。処理時間ではなく承認時刻(イベント時間)基準であってこそ「同じ1分区間」の判定が一貫する。
産業IoTの事例として、数千台の設備のセンサーストリームからセッションウィンドウ(空白閾値10分)で稼働セッションを分離すれば、一つのセッション内の振動・温度傾向で異常を早期に検知できる。稼働と停止が不規則な設備では、固定ウィンドウよりセッションウィンドウが現実の運転区間を自然に反映する。グローバルな配車サービスが、リアルタイムの需給マッチングとETA予測にFlinkベースのストリーム処理を大規模に運用しているのも同じ文脈として知られている。
電子商取引の推薦ではスライディングウィンドウが使われる。ユーザーごとに直近10分のクリックを1分間隔のスライディングで集計し関心商品の傾向をリアルタイムに反映すれば、セッションが長くなるほど推薦が最新の行動に敏感に反応する。このときイベントが複数ウィンドウに重複して含まれ状態が肥大するため、状態TTLとチェックポイント周期(例: 10秒)を併せて設計しメモリ・復旧時間を管理する。例えば活性ユーザー100万人 × 平均状態1KBなら約1GBのキーごとの状態が常時維持されるため、状態バックエンドをメモリからRocksDBへ移し、遊休ユーザーの状態をTTLで失効させる決定がスループット・コストへ直接的に影響する。
4. 深掘り: 厳密に一度の処理の実装機構と最新動向
exactly-onceは「転送を一度だけ」ではなく「効果が一度だけ」を保証する問題である。Flinkは非同期バリアスナップショット(Asynchronous Barrier Snapshotting)でこれを実装する。ソースがデータの流れに特殊なバリア(barrier)を定期的に挿入すると、各オペレータはバリアが通過する際に自らの状態をスナップショットとして取り保存する。すべてのオペレータのスナップショットが揃うと全域的に一貫した(globally consistent)チェックポイントが完成し、障害時にはこの地点へ状態とソースオフセットを併せて戻す。処理を止めず非同期にスナップショットを取るため性能低下が小さいことが核心である。これは分散スナップショット理論(Chandy-Lamportアルゴリズム)をストリームデータの流れに合わせて変形したものである。
状態復元だけではシンクへ既に出た出力の重複を防げないため、出力段では冪等書き込みまたはトランザクション書き込み(2相コミット)を組み合わせる。例えばFlinkのTwoPhaseCommitSinkFunctionはチェックポイント完了の時点に合わせて外部トランザクションをコミットし、チェックポイントと出力可視性を原子的に束ねる。Kafkaはトランザクション(プロデューサトランザクション)とコンシューマオフセットコミットを一つのトランザクションに束ねるread-process-writeパターンでKafka内部の入出力に対するexactly-onceを提供し、Kafka Streamsはprocessing.guarantee=exactly_once_v2の設定でこれを単純化した(KIP-447系列の改善で多数パーティション処理時の効率が向上したと知られている)。
一方、Spark Structured Streamingはチェックポイントにオフセットと状態を記録するWAL(Write-Ahead Log)に基づいて障害復旧を保証し、冪等シンクやトランザクションシンクと組み合わせてend-to-endの正確性を確保する。Kafka Streamsは状態をチェンジログトピックへ複製しておき別のインスタンスで復元するため、状態を持つアプリケーションを無状態サービスのように水平拡張・再配置できる。エンジンごとに機構の実装は異なるが、「一貫した状態スナップショット + 出力の冪等・トランザクション保証」という共通原理は同一であるという点を理解することが核心である。
最新動向としては幾つかの流れが際立つ。第一に、GoogleのDataflowモデルで確立された「何を(What)・どこで(Where, ウィンドウ)・いつ(When, ウォーターマーク・トリガー)・どのように(How, 遅刻データの訂正)」の四つの軸が、Apache Beamを通じてエンジン中立の抽象として広がり、一つのパイプライン定義を複数の実行エンジンで再利用する方向が強まっている。第二に、ストリームをテーブルのようにSQLで扱うストリーミングSQL(Flink SQLなど)が成熟し、非開発職の接近性が高まった。第三に、状態ストアをクラウドオブジェクトストレージへ分離してコンピュートと状態を独立に拡張する状態・コンピュート分離アーキテクチャと、ストリーム・バッチを単一のテーブルフォーマット(例: Apache Icebergなどのレイクハウス)上で統合しようとする試みが活発である。ただしこれらの詳細な数値・バージョンは速く変わるため、設計時に最新の公式文書で確認するのが望ましい。
5. 考慮事項および示唆
第一に、遅延と完全性のトレードオフをビジネス要求に換算して設計する。 ウォーターマークの許容遅延とウィンドウサイズは、技術パラメータである以前に「どれだけ遅いデータまで結果に反映する価値があるか」という業務判断である。精算のように完全性が重要なら許容遅延を十分に取り事後訂正の経路を設け、警報のように即時性が重要なら短く取り遅刻分は別途補正する。技術士の観点では、SLA・正確性要求をまず定義しパラメータを逆算する接近が望ましい。
第二に、exactly-onceの範囲とコストを明確にする。 exactly-onceは通常「エンジン状態と特定のシンク」に限られ、外部システム全体へ自動的に拡張されない。冪等設計が不可能な外部連携(例: 外部へのメール・決済)はトランザクションで束ねにくいため、at-least-once + 冪等キー設計で補わねばならない。またチェックポイント周期を短くすれば復旧時間は減るがオーバーヘッドが増すため、復旧目標時間(RTO)とスループットの均衡点を見つけねばならない。
第三に、状態管理と再処理戦略を併せて設計する。 キーごとの状態が無限に肥大しないようTTL・状態バックエンド(例: RocksDB)・状態サイズ監視を置き、ロジック変更時に過去データを安全に再処理できるようブローカー保存期間とオフセット巻き戻し戦略、結果ビューの冪等再生成構造をあらかじめ用意する。カッパアーキテクチャ([[lambda-kappa-architecture]])の再処理哲学と連携して設計すれば運用複雑度を下げられる。
第四に、運用観測性と逆圧(backpressure)対応を内在化する。 ウォーターマーク遅延の急増、チェックポイント失敗率、消費遅延(consumer lag)、逆圧の発生区間はリアルタイムパイプラインの健全性指標である。リアクティブストリームの背圧制御([[reactive-streams-backpressure]])とメッセージキュー([[message-queue-kafka]])のパーティション・オフセット観測を組み合わせ、急増時にスループットを自動調整し状態が復旧可能な範囲を外れないようにせねばならない。
第五に、展望と連携技術。 ストリーミングSQLとBeamベースのエンジン中立抽象、レイクハウスベースのストリーム・バッチ統合は「リアルタイムとバッチの境界を曖昧にする」方向へ収束している。最終一貫性([[eventual-consistency]])・CQRSとイベントソーシング([[cqrs-event-sourcing]])のようなデータ一貫性モデルと併せて見れば、ストリーム処理は単一の機能ではなくリアルタイムデータプラットフォーム全般を貫く設計の軸として理解すべきである。
参考資料
- Apache Flink Documentation, "Timely Stream Processing / Windows": https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/time/
- Tyler Akidau et al., "The Dataflow Model", VLDB 2015: https://research.google/pubs/pub43864/
- Apache Kafka Documentation, "Exactly Once Semantics": https://kafka.apache.org/documentation/#semantics
- Apache Beam Programming Guide, "Windowing & Watermarks": https://beam.apache.org/documentation/programming-guide/
一言まとめ: ストリーム処理はウィンドウイングで無限データを有限単位にまとめ、イベント時間・ウォーターマークで集計の確定時点を定め、チェックポイントとトランザクションシンクで厳密に一度の処理を保証して、遅延・正確性・完全性のトレードオフを設計するリアルタイムデータ処理技術である。