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

Apache Spark — 遅いジョブの答えは実行計画とイベントログにある

ストリーミングは小さなバッチの連続で、記憶はチェックポイントが担う

TT Labで続きを見る

一言でいうと

Structured Streamingは、際限なく増えるテーブルに同じクエリを小さなバッチで繰り返し実行するエンジンで、どこまで処理したかはチェックポイントが、何を書き出したかはシンクの記録が覚えています。遅れて届いたデータをどれだけ待つかは、ウォーターマークが決めます。

なぜバッチを毎日やり直してはいけないのか

ランディングフォルダーに、パートナーのファイルが1日に何十回も届くとします。バッチジョブで解くなら、道は2つしかありません。毎回フォルダー全体を読み直すか(データが溜まるほど遅くなります)、すでに読んだファイルの一覧を自分で管理するか(ジョブが途中で死ぬと、一覧と結果がずれます)です。2つ目の道を正しく作ろうとすると、結局「処理することにしたものを先に書き、書き終えたあとで終わったと書く」仕組みを、手で組むことになります。

Structured Streamingは、その仕組みをエンジンに組み込んだものです。入門ドキュメントは、入ってくるデータを、行が追加され続ける入力テーブルと見なし、クエリをその上の結果テーブルと見なします。トリガーごとに、新しい行だけで結果を更新します。クエリはバッチとまったく同じように書き、増分で動かす仕事はエンジンが引き受けます。概要ドキュメントによると、デフォルトの実行方式は、この仕事を小さなバッチジョブの連続として処理するマイクロバッチです。

どう動くのか(ファイルソース)

APIドキュメントのファイルソースは、ディレクトリに新しく現れたファイルを読みます。知っておくべきルールが4つあります。

3つ目のルールが、実務の重複を生みます。パートナーが昨日送ったファイルを、名前だけ変えて再送すると、エンジンはそれを新しいファイルとして受け入れます。ファイル単位の「1回ずつ」は保証されますが、内容単位の1回ずつではありません。それは、あとで重複排除で解決します。

チェックポイントとシンクの記録

マイクロバッチ1つが動く順序を描いた図です。ランディングフォルダーから新しいファイルを選んだあと、チェックポイントのoffsetsに今回のバッチが処理する範囲を先に書き、処理した結果を出力フォルダーにファイルとして書いて、_spark_metadataにそのファイル一覧を書き、最後にcommitsに終わったと書きます。途中で死ぬと、commitsがないバッチを同じ範囲でもう一度実行します

入門ドキュメントの耐障害性の節は、設計を1文で要約しています。ソースごとに読んだ位置を表すオフセットがあり、エンジンはチェックポイントと先行書き込みログ(write-ahead log)で、トリガーごとに処理するオフセットの範囲を記録します。シンクは、同じバッチを再び受け取っても結果が同じになるように(冪等に)設計されています。再読み込みできるソースと冪等なシンクが組み合わさって、エンドツーエンドで正確に1回になります。

チェックポイントディレクトリを開くと、この設計がファイルとして見えます。offsets/0は、バッチ0が処理することにした範囲で、処理の前に書かれます。commits/0は、そのバッチが終わったという印で、処理の後に書かれます。再起動したクエリは、offsetsにはあるのにcommitsにはないバッチを探して、同じ範囲でもう一度実行します。ファイルシンク側にも、対になる記録があります。出力フォルダーの_spark_metadata/0に、バッチ0が書き出したファイルの一覧が記録され、Sparkでそのフォルダーを読むと、この一覧にあるファイルだけが見えます。死んだバッチが残した半端なファイルが結果に混ざらない理由です。APIドキュメントの表で、ファイルシンクが「正確に1回」と表示されているのも、この記録のおかげです。

そのため、チェックポイントを消すのは、記憶を消すことです。同じクエリを新しいチェックポイントで起動すると、フォルダーのすべてのファイルを最初から処理し直します。同じドキュメントは、再起動の間にファイルシンクの出力パスや重複排除のカラムを変えることも、許可されない変更として挙げています。

トリガー(いつバッチを動かすのか)

トリガーを指定しなければ、前のバッチが終わり次第、次のバッチを実行します。間隔を指定すれば、その間隔ごとに動きます。このラボのように、あるものだけを処理して止まる使い方には、availableNowが合います。APIドキュメントによると、実行時点にあるデータをすべて処理したあと、自分で止まりますが、ソースのオプション(ファイルソースならmaxFilesPerTrigger)に従って複数のバッチに分けて処理し、前回の実行でコミットできなかったバッチを先に処理します。以前の1回だけのトリガー(once)は、廃止予定です。

q = (spark.readStream.schema(schema).json("/root/landing")
       .withWatermark("event_time", "10 minutes")
       .dropDuplicates(["event_id", "event_time"])
       .writeStream.format("parquet")
       .option("path", "/root/out")
       .option("checkpointLocation", "/root/chk")
       .trigger(availableNow=True)
       .start())
q.awaitTermination()

重複排除とウォーターマーク

ストリーミングのdropDuplicatesは、バッチと意味は同じですが、代償が違います。dropDuplicatesのドキュメントによると、ストリーミングでは、すでに見たキーをトリガーをまたいですべて状態として保持する必要があります。ウォーターマークがなければ、その状態は際限なく大きくなります。

ウォーターマークは、「これより遅れて来るデータは、もう待たない」という線です。withWatermarkのドキュメントは、その線を、これまでに見た最大のイベント時刻からしきい値を引いた値と定義します。この線は2つの仕事をします。ウィンドウ集計のどのウィンドウが確定したかを知らせ、確定したウィンドウの状態を片付けます。そのため、Appendモードのウィンドウ集計は、ウィンドウが終わってもすぐには出力されません。ウォーターマークがウィンドウの終端を過ぎてから1回だけ出力されます。

保証は一方向だけだという点を、必ず覚えておきます。APIドキュメントは、10分のウォーターマークが10分以内の遅れのデータは決して捨てないと保証していますが、それより遅れたデータが必ず捨てられるとは言っていません。たいていは捨てられますが、集計されることもあります。ウォーターマークは、遅れたデータを除外するフィルターではなく、状態を片付ける基準です。そして、集計で状態を片付けるには、ウォーターマークを集計の前に、集計に使う同じ時刻カラムにかける必要があり、出力モードはAppendかUpdateでなければなりません。

現場での姿

1つ目は、チェックポイントを一時フォルダーに置くことです。再起動1回で記憶が消え、クエリがすべてのファイルを処理し直します。チェックポイントは、出力と同じくらい大切なデータです。

2つ目は、ランディングフォルダーでその場書き込みをすることです。アップロードツールがファイルを書いている最中にトリガーが動くと、半端なファイルを読みます。一時的な名前で書き終えてから移動させます。

3つ目は、再送が重複を生むことです。ファイルソースはパスで新しいファイルを判別するので、名前を変えて再送したファイルは、そのまま入ってきます。ユニークなIDで重複を排除し、状態が際限なく大きくならないよう、ウォーターマークを併用します。

実務で本当に大切なこと

次のラボですること

ランディングフォルダーにデバイスイベントのファイルを入れ、availableNowで最初のバッチを実行して、チェックポイントのcommitsと出力フォルダーの_spark_metadataにバッチ0の記録ができるのを確認します。新しいファイルを入れるとそれだけが処理されること、2つのファイルを同時に入れると1つのバッチで処理されることを見ます。名前だけ変えて再送されたファイルが重複として入ってきたevent_idを数えたあと、新しいチェックポイントでdropDuplicatesのストリームを実行して、ユニークな件数だけを残します。最後に、更新時刻を保ってコピーした2つ目のランディングフォルダーで、トリガーごとにファイル1つずつ、10分のウォーターマークと10分のウィンドウ集計を実行して、進捗記録でウォーターマークが動く様子と、40分ほど遅れて届いたイベントが捨てられる様子を見ます。