Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
届くファイルを一度ずつだけ処理し、遅れたデータを捨てる
目標
ランディングフォルダーにファイルが届くたびにStructured Streamingで処理してParquetシンクに書き込み、チェックポイントとシンクログが「どこまで処理したか」をどう覚えているかを確認します。名前だけが違う再送ファイルが重複を生むのを見て、dropDuplicatesで防いだあと、ウォーターマークが遅れて届いたデータを捨てる様子を、ウィンドウ集計で見ます。
なぜ重要なのか
ストリーミングの難しさは、計算ではなく記憶です。ジョブが死んで再び起動したとき、何をすでに行い、何をまだ行っていないかを知る必要があります。Structured Streamingは、これをチェックポイントフォルダーに書きます。バッチを始める前に、読む範囲をoffsetsに書き、シンクに結果を書き終えたあとでcommitsに書きます。ファイルシンクは、自分のフォルダーの_spark_metadataに、バッチごとに書いたファイルの一覧を残し、読む側は、その一覧にあるファイルだけを結果として見ます。
その記憶には限界があります。ファイルソースはファイル名で処理したかどうかを覚えているので、パートナーが同じ内容を別の名前で再送すると、新しいデータとして処理します。内容で重複を消すには、見たことのあるキーを状態として保持する必要があります。
イベント時刻でウィンドウを集計すると、また別の問題が生じます。ウィンドウをいつ閉じるかです。ウォーターマークは「これまでに見た最も遅い時刻から許容遅延を引いた値」で、それより古いデータは遅延データとして捨て、終端がウォーターマークを過ぎたウィンドウだけを結果として出します(appendモード)。ウォーターマークはバッチの間でしか動かないので、データが何個のバッチに分かれて入ってくるかが、結果を変えます。
ステップ
- /root/spk/stream/in(ランディングフォルダー)を作成し、
/data/stream/batch-01.jsonlをコピーしてください。 - /root/spk/stream/stream.py(アプリ
spk-stream-run)で、ランディングフォルダーを読み、/root/spk/stream/out/eventsにParquetで書き込むストリームを、trigger(availableNow=True)で1回実行してください。チェックポイントは、/root/spk/stream/ckpt/eventsです。 batch-02.jsonlをランディングフォルダーに追加して、同じストリームをもう一度実行してください。新しいファイルだけが新しいバッチになる必要があります。batch-03.jsonlとbatch-04.jsonlを一緒に追加して、もう一度実行してください。2つのファイルが1つのバッチで処理される必要があります。- 再送ファイル
batch-06.jsonl(内容はbatch-03と同じ)を追加してもう一度実行したあと、シンクに2回入ったevent_idの数を、/root/spk/stream/out/dups.txtに整数で記入してください。 - /root/spk/stream/dedup.py(アプリ
spk-stream-dedup)で、ランディングフォルダーをdropDuplicates(["event_id"])して、/root/spk/stream/out/dedupに書き込む新しいストリーム(チェックポイントは、/root/spk/stream/ckpt/dedup)を実行してください。 - /root/spk/stream/late_in(2つ目のランディングフォルダー)に
batch-01–batch-05をcp -pでコピーし、/root/spk/stream/window.py(アプリspk-stream-window)で、maxFilesPerTrigger=1、ウォーターマーク10分、10分ウィンドウの件数集計を、appendモードで、/root/spk/stream/out/windowsに書き込んでください(チェックポイントは、/root/spk/stream/ckpt/windows)。進捗記録を、/root/spk/stream/out/progress.jsonに、遅延データの分析を、/root/spk/stream/out/late.jsonに書き込んでください。 - /root/spk/stream/report.mdに、
## 한 번씩만 처리하기・## 재전송과 중복・## 늦은 자료の3つの節を書いてください(見出しは韓国語で、順に「1回ずつだけ処理する」「再送と重複」「遅延データ」を意味します)。2つ目の節にステップ5の重複数を、3つ目の節にステップ7の遅延イベント数を入れてください。
参考
- ファイル:
/data/stream/batch-01.jsonl…batch-06.jsonl(各400行、カラム: event_id, device, event_time, value)。batch-05には40分ほど遅れて届いたイベントが混ざっていて、batch-06はbatch-03の再送です。 - スキーマを必ず指定してください:
event_id STRING, device STRING, event_time TIMESTAMP, value INT。ストリーミングのファイルソースは、デフォルトではスキーマ推論をしません。 - ファイルソースは更新時刻の順にファイルを拾います。ステップ7で
cp -pで元の時刻(ファイルごとに1分間隔)を保たないと、順序がファイル番号と同じになりません。 - シンクログ
out/<싱크>/_spark_metadata/<배치 번호>(プレースホルダーは順にシンク、バッチ番号です)は、最初の行v1のあとに、ファイルごとにJSON 1行が続きます。コミットされた結果は、この一覧にあるファイルだけです。 - よくあるミス: チェックポイントフォルダーを消して実行し直し、最初から処理し直すこと、2つのストリームが1つのチェックポイントを共有すること、ステップ7を1つのバッチで処理してウォーターマークが動かないこと。
- 公式ドキュメント: Structured Streaming Programming Guide・Getting Started・APIs on DataFrames and Datasets
ランディングフォルダーに最初のファイルを置く
/root/spk/stream/inをランディングフォルダーとして作成し、/data/stream/batch-01.jsonlをその中にコピーしてください(コピー先のパス: /root/spk/stream/in/batch-01.jsonl)。
ストリーミングのファイルソースは、フォルダーを見張り、新しく現れたファイルを次のバッチで拾います。ファイルは完成した状態で一度に現れる必要があります。書き込み中のファイルを拾うと、半端な状態が処理されます(そのため、通常は別の場所に書いてから移動します)。
最初のバッチ(チェックポイントとシンクログ)
/root/spk/stream/stream.pyを、アプリ名spk-stream-runで作成し、/root/spk/stream/inをスキーマ(event_id STRING, device STRING, event_time TIMESTAMP, value INT)を指定して読み込み、/root/spk/stream/out/eventsにParquetで書き込むストリームを、checkpointLocation=/root/spk/stream/ckpt/events、trigger(availableNow=True)で開始して、終わるまで待ってください。
availableNowは「今あるものをすべて処理して止まる」です。終わると、チェックポイントにoffsets/0とcommits/0が、シンクフォルダーの_spark_metadata/0にこのバッチが書いたファイルの一覧が残ります。採点ツールは、その一覧のファイルを読んで、batch-01のevent_idと同じかを確認します。
新しいファイルだけが次のバッチになる
/data/stream/batch-02.jsonlをランディングフォルダーに追加して、ステップ2のストリームをそのままもう一度実行してください。新しいバッチには、batch-02のイベントだけが入る必要があります。
チェックポイントがすでにbatch-01を処理したと覚えているので、読み直しません。チェックポイントを消すと記憶も消えて、最初から処理し直し、シンクに同じデータがもう1セット積み上がります。採点ツールは、シンクログでbatch-02のidだけを含むバッチを探します。
2つのファイルが1つのバッチに
batch-03.jsonlとbatch-04.jsonlを一緒にランディングフォルダーに追加して、ストリームをもう一度実行してください。2つのファイルのイベントが1つのバッチで処理される必要があります。
バッチ1つが何個のファイルを拾うかは、maxFilesPerTriggerが決め、指定しなければ、そのときにある新しいファイルをすべて拾います。バッチの境界は遅延とウォーターマークに影響するので、ステップ7でまた出てきます。
名前だけ違う再送が重複を生む
再送ファイル/data/stream/batch-06.jsonl(内容はbatch-03と同じです)をランディングフォルダーに追加して、ストリームをもう一度実行してください。そのあと、シンクのコミットされたファイル(各バッチの_spark_metadataの一覧)を読んで、2回以上出てきたevent_idの数を、/root/spk/stream/out/dups.txtに整数で記入してください。
ファイルソースの記憶は、ファイル名です。名前が新しいので、新しいデータとして処理します。シンクフォルダーをそのまま読んでもかまいませんが(Sparkは_spark_metadataを見て読みます)、ディレクトリのpartファイルを直接走査するツールは、失敗したバッチの残骸まで拾ってしまうことがある点に注意してください。
dropDuplicates(見たことのあるidを覚える)
/root/spk/stream/dedup.pyを、アプリ名spk-stream-dedupで作成し、同じランディングフォルダーを読んでdropDuplicates(["event_id"])した結果を、/root/spk/stream/out/dedupに書き込む新しいストリームを、checkpointLocation=/root/spk/stream/ckpt/dedup、availableNowで実行してください。結果のevent_idは、すべて1回ずつである必要があります。
重複排除は、状態を持つ処理です。見たidを状態ストアに保持しておき、同じidが来たら捨てます。ウォーターマークなしで行うと状態が際限なく大きくなるので、本番ではwithWatermarkと一緒に使うか、dropDuplicatesWithinWatermarkを使います。チェックポイントはストリームごとに別々です。
ウォーターマーク(遅れたデータは捨て、閉じたウィンドウだけを出力する)
/root/spk/stream/late_inを作成して、batch-01–batch-05をcp -pでコピーし、/root/spk/stream/window.pyを、アプリ名spk-stream-windowで作成して、maxFilesPerTrigger=1で読み込み、withWatermark("event_time", "10 minutes")のあと、10分ウィンドウで件数を数えて、カラムstart・end・countで、/root/spk/stream/out/windowsにappendモードで書き込んでください(checkpointLocation=/root/spk/stream/ckpt/windows、availableNow)。終わったあと、query.recentProgressから、バッチごとにbatch・rows・watermark・dropped(状態演算子のnumRowsDroppedByWatermarkの合計)を、/root/spk/stream/out/progress.jsonに一覧として書き込み、batch-05が入ってくるときのウォーターマークより早いbatch-05のイベント数を、/root/spk/stream/out/late.jsonに{"watermark": "yyyy-MM-dd HH:mm:ss", "late_events": 정수, "dropped_rows_metric": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。
ウォーターマークは、バッチが終わるときに「見たものの中で最も遅い時刻から10分を引いた値」に上がり、次のバッチから使われます。そのため、ファイル1つずつをバッチとして入れて初めて、batch-05が入ってくるときに、ウォーターマークがすでに09:29ごろに来ています。遅れたイベントは数十件なのにdropped指標が1であることも見てください。状態演算子の前で、部分集計がすでにウィンドウごとに1行にまとめているからです。
記憶・重複・遅延を数値で残す
/root/spk/stream/report.mdに、## 한 번씩만 처리하기・## 재전송과 중복・## 늦은 자료の3つの節を書いてください(見出しは韓国語で、順に「1回ずつだけ処理する」「再送と重複」「遅延データ」を意味します)。2つ目の節にステップ5の重複数を、3つ目の節にステップ7のlate_eventsとdropped_rows_metricを入れてください。
最初の節には、チェックポイントのoffsets・commitsとシンクログが、それぞれ何を記憶しているか、2つ目の節には、なぜ名前だけ変わったファイルが重複になるのか、3つ目の節には、遅延イベント数と指標がなぜ違うのかを書いてください。