Apache Flink — ストリームを本物のエンジンで動かす
クラスターを起動し、ジョブを一つ最後まで追う
目標
Flinkのローカルクラスターを起動してREST APIで構造を確認したあと、バッチジョブ1つを提出から終了まで追います。並列度・スロット・失敗の履歴を自分で変えてみて、その結果を報告書にまとめます。
なぜ重要なのか
運用で受ける質問(ジョブがなぜ止まったのか、並列度を上げたのになぜそのままなのか、なぜ再起動されなかったのか)は、SQLではなくエンジンの構造についてのものです。JobManagerは、ジョブごとにJobMasterを作り、スロットを分け与え、すべての記録をRESTで出します。このラボの採点ツールは、クラスターに問い合わせません。皆さんがファイルとして保存したRESTレスポンスとsql-clientの出力だけを読み、集計の期待値は元のCSVから直接計算して照合します。そのため、レスポンスは手で直さず、curlの出力をそのまま保存してください。
ステップ
flink-upでクラスターを起動し、/overviewのレスポンスを、/root/flink/cluster/overview.jsonに保存してください(TaskManager 1つ・スロット2つが見える必要があります)。/taskmanagersのレスポンスを、/root/flink/cluster/taskmanagers.jsonに保存し、そのmemoryConfigurationをMiBに直して四捨五入した値7つを、/root/flink/cluster/memory.jsonに書いてください。- /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に保存してください。 - ジョブが終わったあと、
/jobs/overviewのレスポンスを、/root/flink/cluster/jobs.jsonに保存してください。 - /root/flink/cluster/p2.sqlに、同じ集計を並列度2・ジョブ名
flk-first-p2で動かすSQLを作成し、出力を、/root/flink/cluster/p2.outに、そのジョブの/jobs/<jid>のレスポンスを、/root/flink/cluster/p2-job.jsonに保存してください。頂点の並列度に2が見える必要があります。 /opt/flink/conf/config.yamlのtaskmanager.numberOfTaskSlotsを4に変えてクラスターを再起動し、/overviewを、/root/flink/cluster/overview-4.jsonに保存してください。- /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に保存してください。 - /root/flink/cluster/report.jsonに、
slots_total・first_job_id・p2_job_id・failed_job_id・restart_strategy・root_causeを書いてください。
参考
- 元の列:
order_id BIGINT, user_id STRING, status STRING, amount INT, order_time TIMESTAMP(3)(ヘッダーなしのCSV)。 - SQLファイルの実行:
sql-client.sh -f 파일.sql > 파일.out 2>&1(プレースホルダーはファイル名です)。エラーが起きた文で止まり、出力ファイルに[ERROR]が残ります。 - 結果は、表形式(tableau)モードで出力されるように設定されています。バッチで動くと
op列がなく、ストリーミングで動くと先頭にop列(+I・-U・+U)が付きます。 - ジョブidの探し方:
curl -s localhost:8081/jobs/overview | jq -r '.jobs[] | select(.name=="잡이름") | .jid'(プレースホルダーはジョブ名です) - 設定ファイルでスロット数は、
taskmanager:の下にインデントされたnumberOfTaskSlots:の行です。ファイルだけを直しても反映されません(flink-downのあとにflink-up)。 - よくある間違い: バッチジョブは、アダプティブスケジューラーがデータサイズで並列度を決め直します。小さなファイルなら、1になります。
- このPodにはインターネットがありません。メモリの上限は2Giなので、クラスターとSQLクライアントが合わせて1.5GB前後を使います。
- 公式ドキュメント: Flink Architecture・REST API・TaskManagerのメモリ・Adaptive Batch・Task Failure Recovery
クラスターを起動して概要を取得する
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から写せばよいです。