TT Lab
はじめる
学ぶ 学習パス コース

Apache Flink — ストリームを本物のエンジンで動かす

クラスター一式 — JobManager が決め、TaskManager が動かす

TT Labで続きを見る

一言でいうと

Flinkクラスターは、決めるプロセス(JobManager)1つと、動くプロセス(TaskManager)複数です。SQL 1行はジョブグラフになり、JobManagerがそれをTaskManagerのスロットに分けて載せます。何がどこで動いたか、なぜ死んだかは、すべてREST APIに残ります。

なぜ必要なのか

ストリーム処理を初めて任されると、たいてい「SQLを書けば結果が出る」から出発します。ところが、運用で受ける質問はSQLではありません。「ジョブがなぜ止まったのか」「並列度を上げたのになぜそのままなのか」「TaskManagerのメモリを4GBにしたのに、ヒープはなぜ1GBしかないのか」「再起動はなぜされなかったのか」。これらの質問はすべて、エンジンの構造についてのものです。

バッチスクリプトは、死んだらもう一度動かせば済みます。ストリーミングジョブは何週間も稼働し続け、その間、状態を抱えています。だから「どのプロセスが何を担当するのか」を知っていてはじめて、障害のときにどこを見るかを決められます。このモジュールは、その地図を先に広げます。

どう動くのか

JobManagerの中にDispatcher・ResourceManager・JobMasterがあり、TaskManager 1つにスロットが2つあります。SQLクライアントがDispatcherにジョブを提出するとJobMasterができ、ResourceManagerがスロットを割り当てて、ソースから集計を経てシンクまでつながったタスクがスロットに載ります。REST 8081はDispatcherが開きます

JobManagerは、3つの部品でできています(公式ドキュメントのFlink Architecture)。

部品 役割
Dispatcher RESTインターフェース(デフォルト8081)とWeb UIを開き、ジョブが入ってきたらJobMasterを1つ作る
ResourceManager スロットを管理する。スタンドアロン(standalone)デプロイでは、あるTaskManagerのスロットを分けるだけで、新しいTaskManagerを起動できない
JobMaster ジョブ1つの実行を担当する。ジョブごとに1つずつできる

TaskManagerは、実際に演算を実行し、データをやり取りします。リソース割り当ての最小単位がタスクスロットで、スロット数がそのまま同時に動けるタスク数です。注意点が2つあります。1つ目は、スロットはCPUを分けないことです。ドキュメントの表現どおり、現在のスロットはマネージドメモリ(managed memory)だけを分けます。2つ目は、複数のオペレーターがチェーンで結ばれて、1つのスレッドで動けることです。ソース → フィルター → 変換が1つのタスクになる理由です。スレッド間の受け渡しとバッファリングがなくなり、スループットが上がります。

並列度は、オペレーター1つを何個に複製して動かすかです。並列度2の集計は、サブタスク2つに分かれ、それぞれスロットを1つずつ占めます。ここでバッチとストリーミングが分かれます。バッチジョブは、デフォルトのスケジューラーがアダプティブバッチスケジューラーなので、別に指定していないオペレーターの並列度を、入ってくるデータのサイズで決め直します。基準は、タスク1つあたり平均16MB(execution.batch.adaptive.auto-parallelism.avg-data-volume-per-task)です。小さなファイルにparallelism.default = 2を指定しても並列度が1になるのは、バグではなく、この設計です。自動決定を切るスイッチは、execution.batch.adaptive.auto-parallelism.enabledです。

メモリは、このモジュールで最もよく誤解される部分です。taskmanager.memory.process.sizeは、JVMプロセス全体のサイズです。その中から、JVMメタスペースとJVMオーバーヘッドを先に引いて、残ったFlinkメモリを、さらにフレームワークヒープ・タスクヒープ・マネージドメモリ・ネットワークメモリに分けます。

프로세스 전체 = Flink 메모리 + JVM 메타스페이스 + JVM 오버헤드
Flink 메모리 = 프레임워크 힙 + 태스크 힙 + 프레임워크 오프힙 + 태스크 오프힙 + 관리 메모리 + 네트워크 메모리

このPodは上限が2Giなので、TaskManagerを1024MiBに設定しました。すると、タスクヒープは200MiBちょっとしか残りません。「メモリを与えたのにOOMが出る」という報告は、たいていこの表を逆に読んだところから生まれます。実際の値は、/taskmanagersのレスポンスのmemoryConfigurationに、バイト単位で出ます。

失敗と再起動の場面です。ジョブが死ぬと、JobMasterが再起動戦略を見ます。ドキュメント(Task Failure Recovery)によると、チェックポイントを有効にしなければ、デフォルトは再起動しないで、有効にすると、デフォルトはexponential-delayです。そのため、チェックポイントのないジョブの例外履歴の一番上には、原因ではなくRecovery is suppressed by NoRestartBackoffTimeStrategyが出力されます。本当の原因は、その下のCaused by:の連鎖の末尾にあります。

curl -s localhost:8081/overview                  # TaskManager 수 · 슬롯 수 · 버전
curl -s localhost:8081/jobs/overview             # 잡 목록과 상태(FINISHED · FAILED · RUNNING)
curl -s localhost:8081/jobs/<jid>                # 정점(vertex)별 병렬도
curl -s localhost:8081/jobs/<jid>/exceptions     # 실패 기록(exceptionHistory)

このコードブロックの4つの韓国語コメントは、順に、TaskManagerの数・スロット数・バージョン、ジョブ一覧と状態(FINISHED・FAILED・RUNNING)、頂点(vertex)ごとの並列度、失敗の履歴(exceptionHistory)を表しています。

現場での姿

最もよくある場面は、「並列度を上げたのに速くならない」です。原因は、たいてい次の3つのうちの1つです。スロットが足りずにジョブがリソースを待っているか、バッチジョブなのでアダプティブスケジューラーがデータサイズで並列度を削ったか、ソースの分割(ファイル数・パーティション数)が並列度より少なくて、サブタスクの一部が遊んでいるかです。3つとも、/jobs/<jid>の頂点ごとの並列度とサブタスク数を見れば、見分けられます。ダッシュボードのスループットのグラフだけを見ていると、3つが同じに見えます。

2つ目は、設定を変えたのに反映されない場合です。numberOfTaskSlotsやメモリサイズのように、プロセスが起動するときに読む値は、ファイルだけを直しても何も起きません。TaskManagerを再起動する必要があります。逆に、SET 'parallelism.default'のようにジョブの提出時に読む値は、次のジョブからすぐに効きます。どちらなのかを知らないと、「再起動したらうまくいった」という経験談ばかりがたまります。

3つ目は、障害報告書です。アラートは、「ジョブがFAILEDになった」という1行で届きます。例外履歴の一番上の文をそのまま貼り付けると、全員が再起動戦略を原因と誤解します。連鎖の末尾のCaused byを探して書く習慣が、報告書の質を分けます。

次のラボですること

クラスターを起動してRESTでTaskManagerとスロットを確認し、メモリの予算を数字に書き換えます。注文ファイルをバッチで集計するジョブを動かしてジョブ一覧から見つけ、同じジョブを並列度2でもう一度動かして、頂点の並列度を確認します。スロットを4つに変えて再起動したあと、文字列を整数に変換しようとして死ぬジョブをわざと作って、本当の原因と再起動戦略を報告書にまとめます。