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

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

クラスターを起動し、ジョブを一つ最後まで追う

TT Labで続きを見る

目標

Flinkのローカルクラスターを起動してREST APIで構造を確認したあと、バッチジョブ1つを提出から終了まで追います。並列度・スロット・失敗の履歴を自分で変えてみて、その結果を報告書にまとめます。

なぜ重要なのか

運用で受ける質問(ジョブがなぜ止まったのか、並列度を上げたのになぜそのままなのか、なぜ再起動されなかったのか)は、SQLではなくエンジンの構造についてのものです。JobManagerは、ジョブごとにJobMasterを作り、スロットを分け与え、すべての記録をRESTで出します。このラボの採点ツールは、クラスターに問い合わせません。皆さんがファイルとして保存したRESTレスポンスとsql-clientの出力だけを読み、集計の期待値は元のCSVから直接計算して照合します。そのため、レスポンスは手で直さず、curlの出力をそのまま保存してください。

ステップ

  1. flink-upでクラスターを起動し、/overviewのレスポンスを、/root/flink/cluster/overview.jsonに保存してください(TaskManager 1つ・スロット2つが見える必要があります)。
  2. /taskmanagersのレスポンスを、/root/flink/cluster/taskmanagers.jsonに保存し、そのmemoryConfigurationをMiBに直して四捨五入した値7つを、/root/flink/cluster/memory.jsonに書いてください。
  3. /root/flink/cluster/first.sqlに、バッチモード・ジョブ名flk-first-ordersで/opt/lab/fixtures/data/cluster_orders.csvを読み、status別にorders(件数)・revenue(amountの合計)を出すSQLを書き、sql-client.sh -fの出力を、/root/flink/cluster/first.outに保存してください。
  4. ジョブが終わったあと、/jobs/overviewのレスポンスを、/root/flink/cluster/jobs.jsonに保存してください。
  5. /root/flink/cluster/p2.sqlに、同じ集計を並列度2・ジョブ名flk-first-p2で動かすSQLを作成し、出力を、/root/flink/cluster/p2.outに、そのジョブの/jobs/<jid>のレスポンスを、/root/flink/cluster/p2-job.jsonに保存してください。頂点の並列度に2が見える必要があります。
  6. /opt/flink/conf/config.yamlのtaskmanager.numberOfTaskSlotsを4に変えてクラスターを再起動し、/overviewを、/root/flink/cluster/overview-4.jsonに保存してください。
  7. /root/flink/cluster/fail.sqlで、ジョブ名flk-bad-castとしてstatusをINTにCASTしようとして死ぬジョブを実行し、そのジョブの/jobs/<jid>を、/root/flink/cluster/fail-job.jsonに、/jobs/<jid>/exceptionsを、/root/flink/cluster/fail-exceptions.jsonに保存してください。
  8. /root/flink/cluster/report.jsonに、slots_total・first_job_id・p2_job_id・failed_job_id・restart_strategy・root_causeを書いてください。

参考

クラスターを起動して概要を取得する

flink-upでクラスターを起動したあと、curl -s localhost:8081/overviewのレスポンスを、/root/flink/cluster/overview.jsonにそのまま保存してください。

flink-upはstart-cluster.shを呼び、TaskManagerがJobManagerにつながるまで待ちます。つながる前に取得したレスポンスには、taskmanagersが0と出ます。保存先のディレクトリを先に作ってください。

TaskManagerのメモリ予算を数字で読む

/taskmanagersのレスポンスを、/root/flink/cluster/taskmanagers.jsonに保存し、memoryConfigurationの値をMiBに直して四捨五入したtotal_process_mb・total_flink_mb・jvm_metaspace_mb・jvm_overhead_mb・task_heap_mb・managed_mb・network_mbを、/root/flink/cluster/memory.jsonに整数で書いてください。

値はバイト単位です。1048576で割って四捨五入します。プロセス全体は、Flinkメモリ + メタスペース + オーバーヘッドと等しくなければなりません。この恒等式が合わなければ、どこかを間違えて写しています。jqのroundを使うと、1回で作れます。

バッチジョブを1つ提出する

/root/flink/cluster/first.sqlに、SET 'execution.runtime-mode' = 'batch';・SET 'pipeline.name' = 'flk-first-orders';・元のCSVを読むCREATE TABLE・status別にorders(件数)とrevenue(amountの合計)を出すSELECTを書き、sql-client.sh -f first.sql > first.out 2>&1で、/root/flink/cluster/first.outを作成してください。

filesystemコネクターで、pathはfile:///opt/lab/fixtures/data/cluster_orders.csv、formatはcsvです。列名は指示文の元の列のまま使い、結果の列にはAS orders・AS revenueで別名を付けます。出力にop列が見えたら、バッチではなくストリーミングで動いています。

終わったジョブを一覧から探す

curl -s localhost:8081/jobs/overviewのレスポンスを、/root/flink/cluster/jobs.jsonに保存してください。名前がflk-first-ordersのジョブがFINISHEDと見える必要があります。

JobManagerは、終わったジョブをしばらく記憶しています(クラスターを停止すると消えます)。名前が見えなければ、first.sqlのpipeline.nameがSELECTより前にあるかを確認してください。SETは、そのあとに提出されるジョブから適用されます。

並列度2でもう一度動かす

/root/flink/cluster/p2.sqlに、同じ集計をSET 'parallelism.default' = '2';・ジョブ名flk-first-p2で動かすSQLを作成し、出力を、/root/flink/cluster/p2.outに保存して、そのジョブの/jobs/<jid>のレスポンスを、/root/flink/cluster/p2-job.jsonに保存してください。頂点(vertices)の並列度の最大値が2で、集計結果が並列度1のときと同じである必要があります。

並列度を2にしたのに頂点がすべて1なら、バッチジョブのアダプティブスケジューラーがデータサイズ(小さなファイル)を見て、並列度を決め直したのです。その自動決定を切る設定は、execution.batch.adaptive.auto-parallelismの下にあります。結果を集めるシンクが1のまま残るのは、正常です。

スロットを4つに変えて再起動する

/opt/flink/conf/config.yamlのtaskmanager.numberOfTaskSlotsを4に変えて、クラスターを停止してからもう一度起動し、/overviewのレスポンスを、/root/flink/cluster/overview-4.jsonに保存してください。

スロット数は、TaskManagerプロセスが起動するときに1回だけ読む値です。ファイルだけを直して概要を取得し直すと、今も2が出ます。flink-downで停止し、flink-upで起動してください。同じ名前の行が別のブロックにないかも確認してください。

わざと失敗させて履歴を取得する

ジョブ名をflk-bad-castにして、statusをINTにCASTするSELECTを、/root/flink/cluster/fail.sqlで実行してください。そのジョブの/jobs/<jid>を、/root/flink/cluster/fail-job.jsonに、/jobs/<jid>/exceptionsを、/root/flink/cluster/fail-exceptions.jsonに保存してください。

文字列'paid'を整数に変換した瞬間にタスクが例外を投げ、チェックポイントが無効なジョブは、再起動せず、すぐにFAILEDになります。sql-clientの出力にもエラーが出力されますが、採点ツールは、JobManagerが残した例外履歴(exceptionHistory)を見ます。

報告書: 原因と再起動戦略を区別する

/root/flink/cluster/report.jsonに、slots_total(現在のクラスターのスロット数)、first_job_id・p2_job_id・failed_job_id(各ジョブのjid)、restart_strategy(一番上の例外の文で再起動を止めた戦略の名前)、root_cause(スタックの最後のCaused by:の例外クラス、パッケージを含む)を書いてください。

一番上の例外は、「Recovery is suppressed by …」で始まるJobExceptionです。それは原因ではなく、「再起動しなかった」という結果です。本当の原因は、スタックをたどって下った最後のCaused byにあります。jidは、前に保存したJSONから写せばよいです。