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

データパイプライン

バッチ取り込みとウォーターマーク

TT Labで続きを見る

目標

抽出、ロード、ウォーターマークの記録、増分処理、再実行の検証まで、バッチパイプラインの1サイクルを、手で回してみます。

なぜ重要なのか

バッチパイプラインで事故が起きる地点は、ほぼ決まっています。境界条件と再実行です。

境界条件は、>と>=の1つの違いで、実行のたびに1件が重複したり欠落したりします。1日に1回動くバッチなら、1年に365件がずれ、その間、誰も気づきません。そのため、抽出の直後に、件数を元のデータと突き合わせる検証ステップを、パイプラインの中に入れるのがよいです。

再実行はもっと重要です。パイプラインは必ず失敗し、失敗したらもう一度実行します。このとき、結果が変わると、データが汚染されます。そのため、ロード先の表に一意制約をかけ、衝突したときは無視するか更新するように設計する必要があります。このラボでは、同じ増分作業を2回実行して、件数が変わらないかを、自分で確認します。

ステップ

作業ディレクトリは/root/etlです。

  1. ordered_atが2025年1月の注文のid、ordered_at、status、total_amountを、/root/etl/orders_2025_01.csvに抽出します。ヘッダーを含めます。基準時刻は、TIMESTAMPTZ '2025-01-01 00:00:00+09'以上、TIMESTAMPTZ '2025-02-01 00:00:00+09'未満です。
  2. ファイルの行数が、該当期間の注文件数にヘッダー1行を足した値と等しいかを確認します。
  3. orders_archive表を作成します。カラムはid、ordered_at、status、total_amountの順で、idに主キーまたは一意制約がある必要があります。
  4. ステップ1で作成したCSVを、orders_archiveにロードします。
  5. etl_watermark表を作成します。カラムはjob_name、last_ordered_atで、job_nameがorders_archiveの行に、ロードしたデータの最大のordered_atの値を記録します。
  6. ウォーターマーク以降から、TIMESTAMPTZ '2025-03-01 00:00:00+09'未満までの注文を増分ロードし、ウォーターマークを更新します。ロードしたあと、orders_archiveは、1月と2月の注文をすべて含んでいる必要があります。
  7. 同じ増分作業を、もう一度実行します。実行前後のorders_archiveの件数を、/root/etl/rerun.txtに、1行ずつ2行で記録します。2つの値が同じで、重複行がない必要があります。
  8. orders_archiveのステータスごとの件数を、/root/etl/report.tsvに保存します。ステータスと件数をタブで区切り、ヘッダーは入れません。

参考

1月の注文をCSVに抽出する

ordered_atが2025年1月の注文のid、ordered_at、status、total_amountを、/root/etl/orders_2025_01.csvに抽出します。ヘッダーを含めます。基準時刻は、TIMESTAMPTZ '2025-01-01 00:00:00+09'以上、TIMESTAMPTZ '2025-02-01 00:00:00+09'未満です。

psqlの\copyは、クライアント側のファイルに書き出します。ヘッダーを含めるオプションがあります。

抽出件数を検証する

ファイルの行数が、該当期間の注文件数にヘッダー1行を足した値と等しいかを確認します。

ファイルの行数は、データ行数にヘッダー1行を足した値である必要があります。境界条件を片側だけ含めているかを確認してください。

ロード先の表を作成する

orders_archive表を作成します。カラムはid、ordered_at、status、total_amountの順で、idに主キーまたは一意制約がある必要があります。

再実行の安全性は、ここで決まります。同じ注文が2回入れないようにする制約が必要です。

CSVを表にロードする

ステップ1で作成したCSVを、orders_archiveにロードします。

\copyは、反対の方向にも使います。ヘッダー行をスキップするオプションを忘れないでください。

ウォーターマークを記録する

etl_watermark表を作成します。カラムはjob_name、last_ordered_atで、job_nameがorders_archiveの行に、ロードしたデータの最大のordered_atの値を記録します。

どこまで処理したかを表に残します。値は、ロードしたデータの最後の時刻である必要があります。

2月分を増分ロードする

ウォーターマーク以降から、TIMESTAMPTZ '2025-03-01 00:00:00+09'未満までの注文を増分ロードし、ウォーターマークを更新します。ロードしたあと、orders_archiveは、1月と2月の注文をすべて含んでいる必要があります。

ウォーターマーク以降のデータだけを選んで入れ、終わったらウォーターマークを更新します。

2回回しても同じか確認する

同じ増分作業を、もう一度実行します。実行前後のorders_archiveの件数を、/root/etl/rerun.txtに、1行ずつ2行で記録します。2つの値が同じで、重複行がない必要があります。

同じ増分作業をもう一度実行し、前後の件数をファイルに2行で残します。衝突時に無視する句が必要です。

ロード結果のレポートを作る

orders_archiveのステータスごとの件数を、/root/etl/report.tsvに保存します。ステータスと件数をタブで区切り、ヘッダーは入れません。

ステータスごとの件数を、タブで区切って保存します。ヘッダーなしで、値だけを残してください。