ノートでは動いたのにポッドで死ぬ
目標
大きなファイルをメモリに載せずに扱うツールstream.pyを作り、メモリの上限をかけて、それが本当にストリーミングであることを終了コードで証明します。集計・上位N・ユニークな値の計数・外部ソート・検証を、すべて手に持つものの個数で設計します。
なぜ重要なのか
ファイルをまるごと読み込むコードは、小さなファイルではより短く、より速いです。そのため、最初に書くときは疑う理由がなく、ファイルが大きくなるかどうかは、私たちではなく顧客が決めます。コンテナがメモリの上限を超えると、カーネルがプロセスを殺しますが、その死はPythonの例外として捕まえられず、標準出力にも何も残りません。 設計の基準は、ファイルサイズではなく、どの時点でも手に持っているものの個数です。集計は1行、上位N件はN個で足ります。ところが、ユニークな値の計数は、ストリーミングでも、メモリが値の種類数に比例します。この違いを知らなければ、「ストリーミングで書いたのに、なぜ落ちるのですか」ということになります。 そして、「動かしてみたらできましたよ」は証明ではありません。そのときそのファイルでできたという意味にすぎません。上限をかけて通って初めて証明であり、その結果が文章ではなく終了コードとして残れば、次の人が聞き直しません。 採点ツールは、提出された文言を信じません。一時ディレクトリに、採点ツールが作ったログを用意し、作成したツールをメモリの上限の下で実際に実行して、答えを照合します。件数と金額は、実行ごとに変わります。
ステップ
- /root/stream/gen_events.pyを作成して実行し、events.csvを作成してください(保存先: /root/stream/data/events.csv)。30万行です。
- /root/stream/stream.pyに
agg <파일>を作成して(プレースホルダーはファイルです)、1行ずつ読んで件数・合計・平均・最小・最大を出力させてください。 top <파일> <N>を追加して(プレースホルダーはファイルです)、ヒープで上位N件だけを保持させてください。distinct <파일> <칼럼>を追加して(プレースホルダーはファイルとカラムです)、ユニークな値を数え、そのときの最大メモリを一緒に出力させ、種類数が少ないカラムと多いカラムを比較して書いてください(保存先: /root/stream/distinct.json)。- /root/stream/cap.pyと、/root/stream/slurp.pyを作成し、メモリの上限64MiBの下で2つの版を並べて実行して、結果を書いてください(保存先: /root/stream/limit.json)。
sortmerge <파일> <칼럼> <청크행수>を追加して(プレースホルダーはファイル、カラム、チャンクの行数です)、分割・ソート・マージで外部ソートさせてください。verify <파일> <칼럼>を追加して(プレースホルダーはファイルとカラムです)、ソートされているかどうかと、順序に依存しないフィンガープリントを出力させ、原本とソート済みのファイルを照合して書いてください(保存先: /root/stream/verify.json)。- 全体を1枚にまとめて、/root/stream/stream_report.jsonと、/root/stream/stream_report.mdを作成してください。
参考
- 実行契約:
python3 /root/stream/stream.py <명령> <파일> [인자](プレースホルダーは順に、コマンド、ファイル、引数です)。コマンドはagg・top・distinct・sortmerge・verifyの5つです。成功すれば終了コード0、ファイルがなければ3、コマンドや引数の数が間違っていれば2です。 aggの応答:rows・sum_amount・mean_amount・min_amount・max_amount・peak_kib。平均は小数第2位で四捨五入します。peak_kibは、そのプロセスが使った最大メモリで、resource.getrusage(resource.RUSAGE_SELF).ru_maxrssで得ます。Linuxでの単位はKiBです。topの応答:{"n": 정수, "rows": [[event_id, amount], ...], "peak_kib": 정수}(プレースホルダーは整数です)。金額の降順で、金額が同じなら、event_idが大きいほうが上です。distinctの応答:column・rows・exact・peak_kib。sortmergeの応答:rows・chunks・chunk_rows・output・peak_kib。チャンクは、入力ファイルがあるディレクトリのchunks/に書き、結果は、同じディレクトリのsorted.csvに、ヘッダーとともに書きます。verifyの応答:rows・ordered・sum_amount・digest・peak_kib。digestは、ヘッダーを除いた行ごとに、sha256(줄 전체 바이트)(韓国語の部分は「行全体のバイト列」という意味です)の先頭8バイトを整数として読んで、すべて足した値を、2の64乗で割った余りで、16桁の小文字の16進数で出力します。順序が違っても同じ値になる必要があります。- ソートキーは、値が整数の形なら
(정수, 원래 문자열)、そうでなければ(0, 원래 문자열)とみなします(プレースホルダーは整数と元の文字列です)。このルールはこのラボの前提です。 python3 /root/stream/cap.py <MiB> <명령> [인자...](プレースホルダーはコマンドと引数です)は、そのコマンドをアドレス空間の上限の下で実行して、{"limit_mib": 정수, "argv": [...], "exit_code": 정수, "ok": true|false}(プレースホルダーは整数です)を出力します。cap.py自身は、子が死んでも0で終わります。python3 /root/stream/slurp.py <파일>(プレースホルダーはファイルです)は、ファイルをまるごと読んでリストに入れる版で、上限の下で落ちるように作る対照群です。- 公式ドキュメント: heapq・resource・hashlib・POSIX sort
- よくあるミス:
read()やreadlines()で始めること、全体をソートしてから先頭を切ること、検証のために原本と結果を両方リストに載せること、チャンクファイルを消さないこと。 - Podのリソースは、CPU 2コア・メモリ2Gi・一時ディスク6Giです。これより大きなファイルを作らないでください。
一度に読めないファイルを作る
/root/stream/gen_events.pyを作成して実行し、events.csvを作成してください(保存先: /root/stream/data/events.csv)。ヘッダーはevent_id,shop_id,kind,amount,tsで、30万行です。
作るときもリストに集めず、1行ずつファイルに書いてください。shop_idは種類数を少なく、event_idは行ごとに違うように作らなければ、あとでユニークな値を数えるときのメモリの違いが見えません。
1行ずつ読んで集計する
/root/stream/stream.pyにagg <파일>を作成し(プレースホルダーはファイルです)、rows・sum_amount・mean_amount・min_amount・max_amount・peak_kibを出力させてください。
ファイルオブジェクトをそのままfor文に入れれば、1行ずつ読まれます。read()やreadlines()を呼んだ瞬間、その性質は失われます。最小・最大は、いま読んだ値と比べるだけでよく、平均は、合計と件数から最後に出します。
上位N件だけを手に持つ
top <파일> <N>を追加し(プレースホルダーはファイルです)、金額の上位N件を[[event_id, amount], ...]として出力させてください。金額の降順で、金額が同じなら、event_idが大きいほうが上です。
全体をソートして先頭を切ると、全体がメモリに載ります。サイズNの最小ヒープを置き、ヒープがN個になったら、新しい値が一番下より大きいときだけ押し込んでください。heapqのheappushpopが一度に行ってくれます。ヒープに入れる値を「金額とevent_id」の組にしておけば、同点の処理まで一緒にできます。
ユニークな値を数えるメモリは何に比例するか
distinct <파일> <칼럼>を追加し(プレースホルダーはファイルとカラムです)、column・rows・exact・peak_kibを出力させてください。そして、shop_idとevent_idの2つのカラムの結果を、rows・low・highとして書いてください(保存先: /root/stream/distinct.json)。
1行ずつ読んでも、すでに見た値を覚えておく必要があるので、メモリが値の種類数に比例します。種類数が少ないカラムと、行ごとに違うカラムを、それぞれ測ってみれば、その違いが数字で見えます。lowには種類数が少ないほう、highには多いほうの結果を入れてください。
上限をかけて証明する
/root/stream/cap.pyと、/root/stream/slurp.pyを作成し、上限64MiBの下でstream.py aggとslurp.pyをそれぞれ実行して、limit_mib・file_bytes・rows・stream_exit・slurp_exitを書いてください(保存先: /root/stream/limit.json)。
Linuxには、プロセスが確保できるアドレス空間の上限があり、Pythonではresourceモジュールで設定できます。subprocessで子を起動するときは、子の側で上限をかけなければ、親も一緒に死んでしまいます。cap.py自身は、子が死んでも0で終わり、子の終了コードをJSONで知らせるだけです。
分割・ソート・マージで外部ソートする
sortmerge <파일> <칼럼> <청크행수>を追加し(プレースホルダーはファイル、カラム、チャンクの行数です)、チャンクごとにソートしてchunks/に書き、まとめて同じディレクトリのsorted.csvにヘッダーとともに書かせてください。応答は、rows・chunks・chunk_rows・output・peak_kibです。
チャンクの行数ぶんだけ集めてソートし、ファイルに書き出してバッファーを空にすることを繰り返します。マージするときは、チャンクファイルをすべて開いておき、各ファイルから1行ずつだけ取り出せば足り、heapq.mergeがソートキーを受け取って、その仕事をしてくれます。前に作ったチャンクファイルは、開始時に消してください。消さないと、次の実行で結果が混ざります。
中間成果物なしで検証する
verify <파일> <칼럼>を追加し(プレースホルダーはファイルとカラムです)、rows・ordered・sum_amount・digest・peak_kibを出力させ、原本とソート済みのファイルを照合して、column・source・sorted・same_multisetを書いてください(保存先: /root/stream/verify.json)。
原本と結果を両方リストに載せて比べると、ストリーミングでソートした意味がなくなります。行ごとにハッシュを出してすべて足せば、順序が違っても同じ値が出て、行が1つでも欠けたり重なったりすると違う値になります。ソートされているかどうかは、直前の行のキーだけを覚えておけばわかります。
証明を1枚に残す
file_bytes・rows・sum_amount・chunks・limit_mib・stream_exit・slurp_exit・digest_matchを書き(保存先: /root/stream/stream_report.json)、/root/stream/stream_report.mdに、## 무엇을 받았나、## 어떻게 처리했나、## 정말로 스트리밍인가、## 남은 한계の4つの節で書いてください(韓国語の見出しは順に「受け取ったもの」「どう処理したか」「本当にストリーミングか」「残る限界」という意味です)。
前のステップで残したlimit.jsonとverify.jsonとdistinct.jsonを読んで、まとめれば足ります。レポートには、上限の値と2つの終了コードを数字で書いてください。それが証明です。残る限界の節には、ユニークな値を数えるメモリが種類数に比例するという事実を書いてください。