【すとりーむしょり】
ストリーム処理 とは?
最終更新:
💡 流れ続けるデータを、全件待たずに処理する
終わりが定まらないデータの流れに対し、全体がそろうのを待たずに継続して処理する方式。センサー値やアクセス履歴などの分析に使う。遅延や結果を出すタイミングは構成による。
📌 このページのポイント
- データの流れを継続的に処理する。小さなバッチに区切る実装もある
- 時間などの区間をウィンドウとして扱い、集計できる
- 出来事が起きた時刻と処理した時刻を区別し、遅れて届くデータへの対応を決める
- 状態の保存・障害からの復旧・重複出力の扱いを、入出力も含めて設計する
バッチ処理と何が違うの?
バッチは、範囲が決まったデータをまとめて扱う。ストリームは、終わりが定まらない流れを継続して扱うんだ。全件が来るのを待てないので、注文が届くのに合わせて売上の集計を更新する、といった使い方をするよ。
届いた1件を必ず即座に処理するの?
処理や出力のタイミングは方式と設定によるよ。Spark Structured Streamingの既定は、流れを小さなバッチに区切って処理する方式。ストリーム処理という名前だけで、1件ずつの実行や一定の応答時間が保証されるわけではないんだ。
ウィンドウって何?
流れを集計する区間だよ。たとえば5分ごとで重ならない区間がタンブリング。過去5分を1分ずつずらすならスライディング。セッションウィンドウは、一定時間データが来ない区切りで活動をまとめる。区間と出力タイミングの設定を決めるんだ。
遅れて届くデータはどうするの?
出来事が起きたイベント時刻と、実際の処理時刻を区別するよ。Flinkではウォーターマークでイベント時刻の進み具合を扱い、許容する遅れなどを設定できる。すべてのデータが順番どおり届くと決めつけず、遅れたデータをどう扱うか設計するんだ。
障害があっても正確に1回だけ処理される?
チェックポイントと入力の再読み込みで、状態を復元する仕組みがあるよ。ただし保証は構成次第で、処理を再実行することもある。内部状態の整合性だけでなく、外部への書き込みで結果が重複しないかも、入出力先の対応や設定を含めて確認する必要があるんだ。
まとめ:ざっくりこれだけ覚えればOK!
「ストリーム処理」って出てきたら「流れるデータを、全件待たずに処理する方式」と思えばだいたいOK!
📖 おまけ:英語の意味
「Stream Processing」 = データの流れを処理すること
💬 Streamは流れ。終わりが定まらないデータを継続して扱うよ。必ず1件ずつ即座に処理が完了する、という保証ではないんだ。