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

データパイプライン

昨日の数字が今日変わった — ウォーターマークと遅れて届いたデータ

TT Labで続きを見る

目標

イベント時間で集計するツールwm.pyを作ります。処理時間とイベント時間の開きを分布で測り、ウォーターマークを到着順に動かし、ウィンドウが閉じたあとに届いたデータを分け、捨てるポリシーと反映するポリシーの結果を比べ、閉じた瞬間の値とそのあとの訂正を別に残します。

なぜ重要なのか

このコースの前のラボで、ウォーターマークを一度扱いました。ロードしたデータの最大時刻を表に書いて、次の実行がその後だけを読むようにするもので、そのウォーターマークが答える問いは、どこまで読んだかの1つだけでした。ここでは、その次を扱います。モバイルアプリは、電波を失って数分後に、機内モードだった端末は数時間後に、イベントを送ります。そのイベントが発生した時刻は、すでに過ぎ去った時刻で、私たちはその時間帯の集計を、すでに送り出しています。ウォーターマークは、そのデータがもう届かないとみなすという約束であり、事実ではありません。そのため、決めることが3つあります。許容遅延をいくつにするか、約束を破って届いたデータを捨てるか遅延更新として反映するか、そして、反映したなら、昨日出した数字と今日出した数字の違いをどう説明するかです。時間を測るラボではありません。イベントに書かれたevent_timeとingest_timeという2つの欄で計算するラボで、プログラムが実際に何秒かかるかは、何の関係もありません。採点ツールは、提出された文言を信じません。一時ファイルに、採点ツールが作ったイベントの流れを用意して、自分で作ったツールを実際に動かし、ウィンドウサイズと許容遅延を変えながら、答えを突き合わせます。件数と遅延の分布は、実行ごとに変わります。

ステップ

  1. /root/wmark/gen_stream.pyを作成して実行し、/root/wmark/work/stream.jsonlを作ってください。
  2. /root/wmark/wm.pyにskewを作り、遅延の分布と、順序がずれた件数を出力させてください。
  3. watermarkを追加して、到着順にウォーターマークを動かしてください。
  4. windowsを追加して、イベント時間でウィンドウを分けて集計させてください。
  5. lateを追加して、ウィンドウが閉じたあとに届いたデータを分けさせてください。
  6. aggを追加して、捨てるポリシーと反映するポリシーの結果を、それぞれ出力させてください。
  7. closeを追加して、閉じた瞬間の値と、そのあとの訂正を別に出力させてください。
  8. /root/wmark/work/watermark_report.jsonと/root/wmark/work/watermark_report.mdを書いてください。

参考

遅れて届くイベントを手に握る

/root/wmark/gen_stream.pyを作成して実行し、/root/wmark/work/stream.jsonlを作ってください。60件以上で、ファイルの順序がingest_timeの昇順で、すべての行でingest_timeがevent_time以上で、遅延が120秒以上の行が5件以上、イベント時間の幅が1200秒以上、keyは3種類以上である必要があります。

遅延を1種類の分布だけで作ると、あとで見るものがありません。大半は数秒以内に入ってきて、一部は数分、ごく一部は数十分後に押し寄せるように、3つの系統にばらまいてください。作ったあとで、ingest_timeでソートしてファイルに書けば、それが到着順になります。シードを固定しておかなければ、許容遅延を変えながら比べる間に、データが揺れます。

どれだけ開くかをまず測る

/root/wmark/wm.pyにskew <파일>(プレースホルダーはファイルです)を作り、件数・順序がずれた件数・遅延の最小・最大・中央値・95パーセンタイルを、JSONで出力させてください。

遅延はingest_time - event_timeです。分位数は最近傍順位でとり、補間しないでください。整数で落ちてこそ、採点も会議も揺れません。順序がずれた件数は、到着順にたどって、これまでの最大イベント時間より早いものを数えればよいです。

ウォーターマークを到着順に動かす

watermark <파일> --lateness=<초>(プレースホルダーは、ファイルと秒です)を追加して、lateness・max_event_time・advances・final_watermarkを出力させてください。advancesは、最大イベント時間を新しく更新した到着の件数です。

ウォーターマークは後ろに戻りません。遅れて届いたイベントが最大イベント時間を下げるのを許すと、すでに閉じたウィンドウが、再び開いて閉じることを繰り返し、どんな値も確定しません。最初の到着は、比べる前のデータがないため、それ自体が1回の前進です。

イベント時間でウィンドウを分ける

windows <파일> --size=<초>(プレースホルダーは、ファイルと秒です)を追加して、ウィンドウごとの件数と金額を出力させてください。ウィンドウ開始はevent_time - (event_time % size)で、このステップでは、遅いか早いかを問わず、すべて数えます。

JSONのキーは文字列である必要があるため、ウィンドウ開始を文字列で書きます。あとでソートするときは、文字列ではなく整数で比べてください。桁数が違うと、文字列のソートは、おかしな順序を出します。

ウィンドウが閉じたあとに届いたデータを選り分ける

late <파일> --size=<초> --lateness=<초>(プレースホルダーは、ファイルと秒です)を追加して、on_time・late・late_by_windowを出力させてください。どのイベントが遅れたかは、それより先に到着したものだけで計算したウォーターマークで判断し、ウォーターマークがそのウィンドウの終わり以上なら、遅れたものです。

遅れは、データの性質ではなく、到着順の性質です。扱いやすくするためにイベント時間で一度ソートしてしまうと、到着順が失われ、遅延データが0件で出ます。問題がないからではなく、見る目をなくしたのです。ファイルに書かれた順序のままたどり、最大イベント時間の更新は、判定を終えたあとで行ってください。最初の到着は、比べる前のデータがないため、遅れません。

捨てるか反映するか

agg <파일> --size --lateness --policy=drop|update(プレースホルダーはファイルです)を追加してください。dropなら、遅延データを除いてdroppedで数え、restatedは空のリストで、updateなら、遅延データも入れ、droppedは0で、遅延データが入ったウィンドウをrestatedに入れます。

捨てるほうを選んでも、捨てた件数を必ず数えてください。数えずに捨てると、あとで合計が合わない理由を説明する根拠がありません。restatedは、ウィンドウ開始を昇順で入れますが、比べるときは整数で比べてください。知らないポリシー名が来たら、終了コード2で終わります。

閉じた瞬間の値と訂正を別に残す

close <파일> --size --lateness(プレースホルダーはファイルです)を追加して、sealed(閉じる前に届いたものだけ)・corrections(ウィンドウ開始の昇順の訂正の一覧)・final(両方を足した値)を出力させてください。遅延データだけがあるウィンドウのsealedは0です。

2つを合わせて1つだけ出すと、昨日の数字が今日変わった理由を説明する方法がありません。別に残せば、変わった量がそのまま答えになります。correctionsには、訂正が実際にあるウィンドウだけを入れ、finalには、すべてのウィンドウを入れます。

変わった数字を説明する1枚

/root/wmark/work/watermark_report.jsonに、size・lateness・events・windows・on_time・late・dropped・restated・max_lag・p95_lag・policyを書き、/root/wmark/work/watermark_report.mdに、## 무엇을 재었나、## 허용 지연을 얼마로 잡았나、## 늦게 온 자료를 어떻게 했나、## 어제 숫자가 바뀐 이유の4つの節で書いてください(韓国語の見出しで、順に「何を測ったか」「許容遅延をいくつにしたか」「遅延データをどうしたか」「昨日の数字が変わった理由」を意味します)。

許容遅延は、1以上、最大遅延未満にし、その値で、遅れて届いたデータが1件以上出る必要があります。出なければ、値を小さくしてください。レポートには、遅れて届いた件数を数字で書いてください。その数字1つが、次の会議で最初に出る質問の答えです。前のステップで作った関数を、そのまま呼べばよいです。