メインコンテンツまでスキップ

Flink を理解する

· 約6分

ビッグデータシステムの発展の過程には、ずっと一つの核心的な問題がありました。データはますます増えるのに、処理速度はますます遅くなる。 この記事ではバッチ処理から話を始め、Flink がなぜ登場したのか、いったい何を解決したのかを整理します。

バッチ処理からストリーム処理へ

従来のデータ処理は、そのほとんどがオフライン処理でした。毎晩ジョブを 1 回実行し、ユーザー行動を集計し、広告データを計算し、レポートを生成し、ビジネス指標を分析する。この方式は Batch Processing(バッチ処理) と呼ばれ、代表的なツールには次のようなものがあります。

  • Hadoop MapReduce
  • Hive
  • Spark

これらのシステムは大量の過去データの処理に向いています。しかしインターネットの発展にともない、多くのビジネスでリアルタイムなデータ処理が必要になってきました。たとえば:

  • リアルタイム不正検知
  • リアルタイムレコメンド
  • リアルタイム監視
  • リアルタイム広告入札
  • リアルタイムログ分析

これらのシーンには共通の特徴があります。データを「生成されると同時に処理」しなければならず、翌日まで待つことはできません。そこで新しい計算モデルが登場しました。Stream Processing(ストリーム処理) です。

初期のリアルタイム計算システムは主に Storm と Spark Streaming に依存していましたが、これらのシステムにはいくつかの問題がありました。

Storm はストリーム処理システムではあるものの、開発が複雑で、耐障害性はそこそこ、状態管理も困難でした。Spark Streaming は安定してはいるものの、その本質はマイクロバッチ処理(Micro Batch) です。つまり、本当の意味でのストリーム処理ではありません。たとえば:

每1秒处理一批数据

このモデルでは、多くのシーンで依然としてリアルタイム性が足りません。そこで Flink が登場しました。その目標は明確です。本物のストリーム処理システムを構築すること。

1) リアルタイムデータ処理

Flink は本当の意味でのミリ秒級のリアルタイム計算を実現できます。ユーザーが商品をクリックする例で見てみましょう。

点击 → Kafka → Flink → 实时推荐

この一連の流れは、数十ミリ秒しかかからないこともあります。

2) 状態管理

ストリーム処理において非常に重要な問題が状態(State) です。たとえばユーザーのクリック回数の集計:

用户A 点击 1次
用户A 点击 2次
用户A 点击 3次

システムはそれまでのデータを覚えておく必要があります。Flink には状態管理機構(Keyed State、Operator State)が組み込まれており、ローカル状態と RocksDB による状態ストレージをサポートしているため、超大規模な状態データを扱えます。

3) 耐障害性

分散システムでは、ノードのダウンは日常茶飯事です。Flink は Checkpoint 機構によってデータの安全を保証します。システムは定期的に計算状態を保存し、ノードがダウンしても Checkpoint から計算を復旧できます。この機構が Exactly Once(厳密に一度だけの処理) ——各データがちょうど一度だけ処理されること——を保証します。

4) イベントタイム

ストリーム処理には、もう一つ非常に難しい問題があります。イベントタイム(Event Time) です。ログデータは遅れて到着することがあり、システム時刻だけで計算すると集計に誤りが生じます。

Flink は Watermark(ウォーターマーク)機構を導入しました。Watermark はどのデータが「おおむね出そろった」かをシステムが判断する助けとなり、遅延データ、ウィンドウ計算、時間ベースの集約を正しく処理できるようになります。

1) Stream First

Flink の核心的な理念はすべてはストリームであるというものです。バッチ処理は単なる有限のデータストリームにすぎません。したがって、Flink の Batch と Stream は同じ一つのエンジンを使います。

2) State Driven

Flink は状態駆動の計算システムであり、すべての計算は状態を中心に展開されます。ユーザー行動の集計、リアルタイム集約、リアルタイム監視——これらはすべて状態管理に依存しています。

3) Exactly Once

Flink はデータ処理の正しさを非常に重視します。Checkpoint、分散スナップショット、2 フェーズコミットによって、Flink は Exactly Once セマンティクスを実現できます。これは多くのリアルタイムシステムにとって極めて重要な能力です。

Flink は多くのインターネット企業で幅広く活用されており、典型的なシーンには、リアルタイムレコメンドシステム、ユーザー行動分析、広告クリック集計、リアルタイム不正検知システム、ログのリアルタイム分析などがあります。よくあるデータパイプラインは次のような形です。

用户行为 → Kafka → Flink → ClickHouse

Flink はその中間で、データクレンジング、リアルタイム計算、リアルタイム集約を担当します。

まとめ

Flink の登場は、本質的には「本物のストリーム処理」を実現するためのものでした。同じ一つのエンジンでバッチとストリームを統一し、組み込みの状態管理で大規模なステートフル計算を支え、Checkpoint 機構と引き換えに Exactly Once の正しさを手に入れ、さらに Watermark で順序の乱れや遅延データに対処する。これらのポイントを理解すれば、Flink と初期のリアルタイム計算システムとの根本的な違いも理解できます。この後で具体的な API やクラスタのデプロイを見ていくときも、ずっとスムーズになるはずです。

COMMENTS