Apache Flink(アパッチ・フリンク)は、終わりのないデータストリームと終わりのあるデータセットの両方を、状態を持ったまま分散処理するオープンソースのエンジンです。不正検知のように「イベントが届いた瞬間に判断する」処理と、日次集計のような「溜まったデータをまとめて計算する」処理を、同じAPIで書けます。
この記事では、公式サイトが挙げるユースケース3分類、JobManagerとTaskManagerの役割、2.0で削除された機能、2026年6月25日公開の最新版2.3.0を手元で動かした出力を順に説明します。最後にKafka StreamsやSpark Structured Streamingとの使い分けと、Flinkを選ぶべきでない条件をまとめました。
まとめ:Apache Flinkの要点と採用判断
- Flinkは有界・無界のストリームを状態付きで処理する分散エンジンです。最新安定版は2.3.0(2026年6月25日)で、ライセンスはApache License 2.0です。
- 公式のユースケースは、イベント駆動アプリケーション、データ分析、データパイプラインの3分類です。代表例は不正検知・異常検知・継続ETLです。
- クラスタはJobManager(ResourceManager・Dispatcher・JobMaster)とTaskManagerで構成され、リソースの最小単位はタスクスロットです。
- 2.0でDataSet API、Scala API、per-jobデプロイモード、flink-conf.yamlが廃止されました。1.x時代の手順書はそのままでは動きません。
- 入出力がKafkaだけでKafka Streamsの機能で足りるならKafka Streams、数分〜1時間遅れの集計で足りるならバッチ処理のほうが運用は軽くなります。Flinkは大きな状態、イベント時刻での集計、Kafka以外との入出力が必要なときに選びます。
Apache Flinkの定義と名前の由来
有界・無界のストリームを状態付きで処理するエンジン
Flinkが扱うデータは2種類です。センサーの計測値やクリックログのように終わりが決まっていない「無界ストリーム」と、昨日分のログファイルのように始まりと終わりがある「有界ストリーム」です。Flinkはバッチ処理を「有界ストリームの処理」とみなすため、同じDataStream APIやFlink SQLのまま実行モードだけを切り替えられます。この違いは後の実行例で出力行数として確認します。
もう一つの特徴が「状態(state)」です。直近5分間の取引額の合計、ユーザーごとの最終ログイン時刻といった途中経過を、Flink自身が保持して障害時にも復元します。外部のデータベースに毎回問い合わせる設計と比べ、計算とデータが同じプロセスにあるため、公式サイトはスループットと遅延の両面で有利だと説明しています。ストリーム処理そのものの考え方はストリーム処理の仕組みとバッチ処理との使い分けで整理しています。
「Flink」の意味とStratosphereからの経緯
flinkはドイツ語の形容詞で、Duden(ドイツ語辞典)は「素早く器用に動く、または働く(sich rasch und geschickt bewegend oder arbeitend)」と定義しています。語源は低地ドイツ語で、もとは「光る、つやのある」を意味したと同辞典は記載しています。
起源はベルリン工科大学(TU Berlin)で2009年に始まった研究プロジェクトStratosphereです。2014年4月14日にApache Incubatorへ入り、ASFは2015年1月12日にトップレベルプロジェクトへの昇格を発表しました。
Apache Flinkのユースケース3分類と代表例
公式サイトのUse Casesページは、Flinkで作られるアプリケーションを次の3つに分けています。どの分類かで、使うAPIと重視する機能が変わります。
| 分類 | 出力先 | 公式が挙げる代表例 | 主に使う機能 |
|---|---|---|---|
| イベント駆動アプリケーション | 外部アクション・状態更新 | 不正検知、異常検知、ルールベースのアラート | ProcessFunction、CEP、セーブポイント |
| データ分析 | DB・ダッシュボード | 通信網の品質監視、モバイルアプリの実験評価 | Flink SQL、ウィンドウ集計 |
| データパイプライン | 別のストレージ | ECサイトの検索インデックス構築、継続ETL | コネクタ、Table API |
3分類は排他ではありません。Kafkaから注文イベントを読み、不正の疑いを即時に止めつつ、集計結果をダッシュボード用DBに書き出す、という1本のジョブに2分類が同居することもよくあります。
イベント駆動アプリケーション:不正検知・異常検知・アラート
公式の定義は「1つ以上のイベントストリームを取り込み、計算・状態更新・外部アクションを起こして反応するステートフルなアプリケーション」です。代表例として不正検知、異常検知、ルールベースのアラート、業務プロセス監視、SNSのようなWebアプリケーションが挙がっています。
公式がイベント駆動アプリケーション向けの際立った機能として挙げるのは、セーブポイントです。状態の整合したスナップショットから、アプリケーションを更新したり並列度を変えたり、複数バージョンを並行起動してA/Bテストしたりできます。「同じカードで10分以内に3カ国から決済」のような時間をまたぐ条件は、パターン検出ライブラリのCEP(Complex Event Processing)で書けます。
データ分析:ダッシュボードの継続更新とアドホック分析
バッチ分析では、新しいデータを取り込むたびにクエリを再実行します。ストリーミング分析では、クエリがイベントを受け取り続けて結果を更新し続けます。公式の代表例は、通信網の品質監視、モバイルアプリのアップデートや実験の評価、ライブデータのアドホック分析、大規模グラフ分析です。
Flink SQLはバッチとストリーミングで同じ意味論を持つため、記録済みデータに対しても流れてくるデータに対しても同じ結果になると公式は説明しています。SQLでの書き方と状態の肥大を避ける設計はFlink SQLの動的テーブルと継続クエリの仕組みで詳しく扱っています。
データパイプライン:継続ETLとCDC取り込み
従来のETLは定期起動でデータをコピーします。Flinkのデータパイプラインは同じ変換・拡充を常時動かすため、転送先に反映されるまでの遅延が短くなります。公式の例は、ECサイトのリアルタイム検索インデックス構築と継続ETLです。
業務DBの変更をそのまま流したい場合は、サブプロジェクトのFlink CDCを使います。2026年9月時点の最新安定版は3.6.0(2026年3月30日公開)で、配布バイナリはFlink 1.20.x向けと2.2.x向けです。CDC 3.6.0の公式対応系列はFlink 1.20.xと2.2.xです。2.xで使う場合は2.2系に合わせ、対応する配布物やコネクタ依存関係を選びます。
JobManagerとTaskManagerの役割とタスクスロット
JobManagerを構成する3つのコンポーネント
Flinkクラスタは、ジョブ全体を調整するJobManagerと、実際に計算するTaskManager(ワーカー)で構成されます。JobManagerはタスクのスケジューリング、チェックポイントの調整、障害からの復旧を担い、内部は次の3つに分かれています。
| コンポーネント | 役割 |
|---|---|
| ResourceManager | タスクスロットの割り当てと解放 |
| Dispatcher | REST APIでジョブを受け付け、Web UIを提供 |
| JobMaster | 1つのジョブ(JobGraph)の実行を管理 |
JobMasterはジョブごとに1つ起動します。StandaloneクラスタのResourceManagerは既存TaskManagerのスロットを配るだけで、TaskManagerを自分で増やせません。YARNやKubernetesでは足りない分を新しく起動できます。この差が、後述する実行基盤の選び方に直結します。
TaskManagerとスロット:ローカル既定値の実測
TaskManagerは、タスクの実行とタスク間のデータのバッファリング・受け渡しを担うプロセスです。1つのTaskManagerはタスクスロットに分割され、スロットがリソース割り当ての最小単位になります。並列度4のオペレーターを動かすには、合計4スロット以上が必要です。
配布版flink-2.3.0のconf/config.yamlと、起動後のREST API(/taskmanagers)で確認した既定値は次のとおりです。
| 項目 | 既定値 |
|---|---|
| taskmanager.numberOfTaskSlots | 1 |
| parallelism.default | 1 |
| jobmanager.memory.process.size | 1600m |
| taskmanager.memory.process.size | 1728m |
| タスクヒープ(REST表示) | 383MB |
| マネージドメモリ(REST表示) | 512MB |
| ネットワークメモリ(REST表示) | 128MB |
スロットは1つだけなので、並列度を2にしたジョブはこの構成ではスロットが足りず実行できません。検証で並列度を上げるときは、taskmanager.numberOfTaskSlotsも同時に増やしてください。マネージドメモリはRocksDB系の状態バックエンドやバッチのソート処理に使われる領域で、ヒープとは別枠です。
Flink 2.xの変更点と最新版2.3.0
2.0で削除されたAPI・設定ファイル・デプロイモード
2025年3月に公開されたFlink 2.0は、互換性を切った大型リリースでした。1.x向けの記事や社内手順書を流用すると、次の箇所で止まります。
- DataSet APIの削除。バッチ処理はDataStream APIのBATCHモードかTable API/SQLへ移行します。
- Scala版のDataStream API・DataSet APIの削除。Java APIを使います。
- per-jobデプロイモードの削除。デプロイモードはApplication ModeとSession Modeの2つになりました。
- flink-conf.yamlの廃止。標準YAML形式の
config.yamlだけが読まれ、変換用にbin/migrate-config-file.shが同梱されています。 - Java 8のサポート終了。既定かつ推奨はJava 17で、最小はJava 11です。
配布物にも変化が表れています。2.3.0のexamplesディレクトリにあるのはstreaming・table・pythonだけで、DataSet時代のバッチ用サンプルはありません。1.x系の解説によく出てくるflink run-applicationを2.3.0で実行すると"run-application" is not a valid action.と返り、有効なアクションはrun・list・info・savepoint・stop・cancelの6つだけでした。
2.3.0時点の版と周辺プロジェクト
本体は複数の系列に並行して修正版が出ています。周辺プロジェクトとマネージドサービスの対応状況とあわせて整理しました(2026年9月14日時点)。
| 対象 | 最新版 | 公開日 |
|---|---|---|
| Flink 2.3系(最新) | 2.3.0 | 2026-06-25 |
| Flink 2.2系 | 2.2.1 | 2026-05-15 |
| Flink 2.1系 | 2.1.3 | 2026-06-14 |
| Flink 1.20系 | 1.20.5 | 2026-06-08 |
| Flink Kubernetes Operator | 1.15.0 | 2026-05-26 |
| Flink CDC | 3.6.0 | 2026-03-30 |
2.3.0の主な追加は、変更ログを変換するSQL演算子FROM_CHANGELOG・TO_CHANGELOG、マテリアライズドテーブルの更新戦略の細かな制御、AWS SDK v2で作り直した実験的なネイティブS3ファイルシステム(flink-s3-fs-native)です。1.20系にも2026年6月に1.20.5が出ているため、2.xのAPI移行が済むまで1.20に留まる選択は現実的です。
ローカルで2.3.0を動かす手順とストリーミング・バッチの出力差
動作を理解する一番の近道は、配布版に同梱のWordCountを2つの実行モードで動かして出力を比べることです。公式の推奨Javaは17で、Java 21は実験的サポートです。以下はmacOS(x86_64)、Eclipse Temurin JDK 21.0.12.1、flink-2.3.0-bin-scala_2.12.tgz(約605MB)で2026年9月14日に実行した手順です。
tar xzf flink-2.3.0-bin-scala_2.12.tgz
cd flink-2.3.0
./bin/start-cluster.sh
# クラスタの状態を確認
curl -s localhost:8081/overview
# 同じjarをストリーミング(既定)とBATCHで実行
./bin/flink run examples/streaming/WordCount.jar --output /tmp/wc_stream
./bin/flink run examples/streaming/WordCount.jar --execution-mode BATCH --output /tmp/wc_batch
./bin/stop-cluster.sh
起動直後の/overviewは次の内容を返しました。TaskManagerが1つ、スロットが1つの構成です。Web UIはhttp://localhost:8081で開けます。
{"taskmanagers":1,"slots-total":1,"slots-available":1,"jobs-running":0,"jobs-finished":0,"jobs-cancelled":0,"jobs-failed":0,"flink-version":"2.3.0","flink-commit":"c0f8d1a"}
2つのジョブはどちらもFINISHEDで終わりましたが、出力は大きく違いました。
| 実行モード | jobType | 出力行数 | 単語 the の出力 |
|---|---|---|---|
| 既定(ストリーミング) | STREAMING | 287行 | (the,1)〜(the,22)の22行 |
--execution-mode BATCH |
BATCH | 170行 | (the,22)の1行 |
入力はどちらも同じ文章で、異なる単語は170語です。ストリーミング実行は単語が現れるたびに途中経過を1行ずつ出し、バッチ実行は入力の終わりを知っているので最終結果だけを出します。コードは1行も変えていません。
ここから分かる実務上の注意があります。ストリーミングの結果をそのままファイルやDBに追記すると、同じキーの古い値が残ります。書き込み先ではキーで上書き(upsert)するか、ウィンドウで区切って確定値だけを出す設計が必要です。区切り方はストリームウィンドウ処理の4方式と遅延データの締め切り設計を参照してください。
デプロイモードと実行基盤の選び方
Application ModeとSession Modeの違い
2.xのデプロイモードは2つです。違いは、クラスタの寿命とリソースの分離度、そしてアプリケーションのmain()をクライアントとクラスタのどちらで実行するかにあります。
| 観点 | Application Mode | Session Mode |
|---|---|---|
| クラスタ | 1アプリ専用 | 複数アプリで共有 |
| main()の実行場所 | クラスタ側 | クライアント側 |
| 障害の影響 | そのアプリだけ | 同居ジョブに波及しうる |
| 起動の速さ | 毎回クラスタを起動 | 起動済みに投入 |
本番の常時稼働ジョブはApplication Mode、短いジョブを頻繁に投げる開発・検証環境はSession Mode、が基本の分け方です。Session ModeはJobManagerとTaskManagerを複数のジョブで共有するため、1つのジョブがメモリを使い切るなど資源を圧迫すると、同じクラスタの他のジョブにも影響が及びます。
Kubernetes Operatorとマネージドサービス
Kubernetes上で自前運用する場合は、Flink Kubernetes Operator(最新1.15.0)が選択肢になり、FlinkDeploymentリソースでジョブとセーブポイントを管理できます。1.15.0のリリース告知が新たな対応として挙げているのはFlink 2.2で、2.3.0には触れていません。2.3.0をOperatorで動かす場合は、事前に検証環境で確認してください。
クラスタを持ちたくない場合の選択肢は2つです。AWSのAmazon Managed Service for Apache Flinkは、2023年8月30日にAmazon Kinesis Data Analyticsから改称したサービスで、2026年9月時点で2.3.0、2.2.1、1.20.5などをサポートしています。Kafka基盤ごとマネージドにするなら、2024年3月19日にAWS・Google Cloud・Azureで一般提供が始まったConfluent Cloud for Apache Flinkがあります。KafkaをAWSで持つ構成はAmazon MSKの仕組みとKafka自前運用との違い、Confluentの製品境界はConfluent KafkaのPlatformとCloudの構成差で比較しています。
チェックポイント・セーブポイントと状態バックエンドの選び方
チェックポイントは、有効化するとFlinkが設定した間隔で作成する障害復旧用のスナップショットです。既定では無効で、復旧には入力の再読込能力とチェックポイントの保存先も必要です。保持・削除は設定に従います。セーブポイントは、アプリの更新や並列度の変更に備えてユーザーが作成し、ユーザーが削除するスナップショットです。公式ドキュメントはこの関係を、データベースの「リカバリーログ」と「バックアップ」の違いにたとえています。コードをデプロイし直すときはセーブポイントを取ってから停止する、が運用の基本です。
状態をどこに置くかは状態バックエンドで決まります。2.3で選べるのは次の3つで、何も指定しなければHashMapStateBackendになります。
| 設定値 | 実装 | 状態の置き場所 | 向く規模 |
|---|---|---|---|
| hashmap(既定) | HashMapStateBackend | JVMヒープ | ヒープに収まる状態 |
| rocksdb | EmbeddedRocksDBStateBackend | ローカルディスク | ヒープを超える大きな状態 |
| forst | ForStStateBackend | リモートストレージ | クラウドで計算と状態を分離 |
ForSt(For Streamingの略)は2.0の主要機能として発表された分離型の状態バックエンドで(1.20にも実験的な実装がありました)、状態をS3などに置いて計算ノードのディスク容量から切り離します。まずはhashmapで始め、状態がヒープに収まらなくなったらrocksdb、ローカルディスクを超える状態や計算とストレージの分離が必要ならforstを検討します。ただし、2.3のForStは実験段階で、canonicalセーブポイントなどに未対応のため、本番採用前に復旧・更新手順まで検証してください。重複も欠落も出さない処理保証の条件はExactly-onceの仕組みと実装条件、イベント時刻の進め方はウォーターマークの生成戦略で解説しています。
Kafka・Kafka Streams・Spark Structured Streamingとの違いと採用判断
「FlinkとKafkaの違い」は、比較の前提から整理が必要です。Apache Kafkaはイベントを保存して配信する基盤で、Flinkはそのイベントを読んで計算するエンジンです。競合ではなく、Kafkaから読みFlinkで処理してKafkaやDBへ書き戻す組み合わせが一般的です。Kafka本体の構成はApache Kafkaの仕組みとKafka 4系のKRaft構成にまとめています。
実際に選択肢として並ぶのは、Kafkaに付属する処理ライブラリKafka Streamsと、Spark Structured Streamingです。
| 観点 | Apache Flink | Kafka Streams | Spark Structured Streaming |
|---|---|---|---|
| 実行形態 | 専用クラスタ | アプリに組み込むライブラリ | Sparkクラスタ |
| 既定の処理方式 | レコード単位 | レコード単位 | マイクロバッチ |
| 入出力 | Kafka以外も多数 | Kafkaトピック | Kafka以外も多数 |
| SQLでの記述 | Flink SQL | ksqlDB(別製品) | Spark SQL |
Kafka Streamsは公式に「クライアントライブラリ」と位置付けられ、別の処理クラスタを必要としません。Spark Structured Streamingは既定でマイクロバッチ方式を使い、公式ドキュメントは最短100ミリ秒程度のエンドツーエンド遅延と説明しています(Continuous Processingモードでは約1ミリ秒・at-least-once)。
判断は次の条件で分けられます。
- 入力も出力もKafkaで、Kafka Streamsの演算と状態管理で要件を満たせるなら、専用クラスタを増やさないKafka Streamsを第一候補にします。Kafka Streamsも状態を複数インスタンスに分散できるため、1プロセスに収まるかだけでは判断しません。Kafka Streamsなら既存のJavaアプリにライブラリを足すだけで、JobManagerやチェックポイント保存先の運用が増えません。SQLで書きたい場合の比較はksqlDBの仕組みとKafka Streamsとの違いを参照してください。
- 既にSparkでバッチ基盤を持ち、秒単位の遅延が許容できるなら、Spark Structured Streamingで同じチームが運用するほうが合理的です。
- 1プロセスに収まらない大きな状態(公式は数TBまでの状態を扱えると説明)、遅れて届くイベントをイベント時刻で正しく集計する処理、Kafka以外のDBやファイルとの入出力が重なるなら、Flinkを選びます。
- 集計が1時間ごとで足りる要件にFlinkを入れるのは過剰です。常時稼働のクラスタとセーブポイント運用のコストに見合いません。
クラスローダー競合(child-first)の原因と対処
ジョブを投入した途端にClassCastExceptionやNoSuchMethodErrorで落ちる場合、原因の多くは依存ライブラリの競合で、その背景にはFlinkのクラス読み込み順があります。Javaの標準は親クラスローダーを先に探す「parent-first」ですが、Flinkはユーザーのjarを先に探す「child-first」を既定にしています(classloader.resolve-orderの既定値)。ユーザーがFlink本体と異なる版のライブラリを同梱できるようにするためです。
ただし、Flink本体と共有する必要があるパッケージは常にparent-firstで読まれます。2.3のclassloader.parent-first-patterns.defaultにはjava.、scala.、org.apache.flink.、org.apache.hadoop.、org.slf4j、org.apache.log4j、ch.qos.logbackなどが含まれます。これらのパッケージは親側のクラスが優先されますが、親側に該当クラスがなければユーザー側から読み込まれる場合があります。Flink本体やログ関連の依存関係は、クラスタ提供分との重複を避けて梱包してください。公式ドキュメントも、Java APIとランタイムモジュールはFlinkが提供するのでジョブのuber JARに含めないよう求めています。
X cannot be cast to Xが出たときの切り分け
com.foo.X cannot be cast to com.foo.Xという一見矛盾したエラーは、同じクラスの複数の版が別々のクラスローダーから読み込まれ、それらを互いに代入しようとしたことを示します。公式ドキュメントは、よくある原因としてライブラリがFlinkの反転したクラス読み込み順(child-first)に対応していないケースを挙げています。対処候補は次の3つです。設定例の1と2は代替案なので、選んだ行だけを設定します。両方を有効にすると2の全体parent-firstが適用され、1のパッケージ限定という意図は失われます。
# config.yaml
# 1. 問題のパッケージだけ親から読ませる(推奨)
classloader.parent-first-patterns.additional: com.foo.
# 2. 全体を Java 標準の parent-first に戻す(影響が大きい)
classloader.resolve-order: parent-first
3つ目は、アプリ側でmaven-shade-pluginを使い、衝突するライブラリのパッケージ名をビルド時に書き換える方法です。公式ドキュメントはAWS SDKのcom.amazonawsを例に挙げています。設定変更はクラスタ全体に効くため、複数チームで共有するSession Modeのクラスタでは、まずshadeで自分のjar内に閉じ込める方法を試してください。
Apache Flinkに関するよくある質問
Apache Flinkは無料で使えますか?
本体はApache License 2.0のオープンソースで、ソフトウェア自体の利用料はかかりません。商用利用や改変も可能です。費用が発生するのは、実行するサーバーやKubernetesクラスタ、チェックポイントを保存するS3などのストレージです。クラスタを運用せずに使いたい場合は、Amazon Managed Service for Apache FlinkやConfluent Cloud for Apache Flinkのような従量課金のマネージドサービスを選ぶことになり、その場合はサービスの料金体系に従います。
AWSでFlinkを使うにはどのサービスを選べばよいですか?
マネージドで動かすならAmazon Managed Service for Apache Flinkです。2023年8月30日にAmazon Kinesis Data Analyticsから改称したサービスで、2026年9月時点では2.3.0、2.2.1、1.20.5などの版を選べます。EC2やEKSに自分で構築する方法もありますが、その場合はクラスタやチェックポイントの保存先を自分で管理します。EKSなどのKubernetes環境では、Flink Kubernetes Operatorを使ってジョブを管理する方法もあります。入力側のKafkaもAWSで持つなら、Amazon MSKとの組み合わせが一般的です。
Flink 1.20から2.xへすぐ移行する必要はありますか?
急ぐ必要はありません。1.20系には2026年6月8日に1.20.5の修正版が出ており、AWSのマネージドサービスも1.20.5をサポート対象に含めています。ただし2.0でDataSet API、Scala API、per-jobモード、flink-conf.yamlが廃止されたため、移行はコード修正を伴います。DataSet APIを使うジョブが残っているなら、1.20のうちにDataStream APIのBATCHモードかTable APIへ書き換えておくと移行準備になります。ただし、2.xへの移行ではコネクタ・設定・削除APIへの対応も必要で、1.xからの状態互換性も保証されないため、状態の復元を含めた移行検証が必要です。
Flink SQLだけでアプリケーションを作れますか?
集計、結合、ウィンドウ処理、別ストレージへの転送といったデータ分析とデータパイプラインの多くは、Flink SQLだけで書けます。公式もデータパイプラインの一般的な変換はSQLとユーザー定義関数で対応できると説明しています。一方、イベントごとに外部APIを呼ぶ、タイマーで独自の状態遷移を管理するといった細かな制御には、DataStream APIのProcessFunctionが選択肢になります。2.3では、SQLやTable APIから利用するProcess Table Functionでも状態やタイマーを扱えますが、関数本体の実装が必要です。SQLで始めて、足りない部分だけDataStream APIで補う構成が現実的です。
TaskManagerのスロット数はいくつに設定すればよいですか?
配布版の既定は1スロットです。既定のスロット共有では、1ジョブに必要なスロット数は、そのジョブ内の最大並列度が目安になります。複数ジョブを同時実行する場合は各ジョブの必要数を合算し、スロット共有グループを分ける場合は別途見積もります。公式ドキュメントによると、スロットが分けるのはマネージドメモリだけで、CPUは分離しません。1台に多くのスロットを詰めるとJVMなどの固定費は薄まりますが、CPUを多く使うタスク同士が同じコアを奪い合います。状態が大きいジョブや計算の重いジョブは、スロットを増やすよりTaskManagerの台数を増やす方向で調整してください。