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

データパイプライン

冪等なupsertパイプライン

TT Labで続きを見る

目標

同じ入力を何回入れても結果が変わらず、内容が実際に変わった行だけが更新される、ロードのパイプラインを作ります。

なぜ重要なのか

パイプラインは必ず失敗します。ネットワークが切れ、デプロイ中に落ち、ソースがデータを遅れて送ってきます。そのたびに、人が「今までどこまで入ったか」を判断しなければならないなら、いつか間違いが起きます。

冪等に設計すれば、その判断が必要なくなります。失敗したなら、そのままもう一度実行すればよいです。核心の仕組みは2つです。自然キーに一意制約をかけて、重複をデータベースに防がせることと、内容ハッシュで実際の変更の有無を判断して、値がそのままの行には触れないことです。

2つ目が特に重要です。条件なしに更新すると、変更のない再実行でも、すべての行の更新時刻が上がります。そうなると、「何が実際に変わったか」がわからなくなり、下流で変更分だけを取得する増分の消費が壊れます。

ステップ

対象はstaging.orders_rawで、正規化ルールは、前のラボと同じです。

  1. orders_final表を作成します。カラムはorder_ref、order_date、amount、status、content_hash、created_at、updated_atの順で、order_refが主キーである必要があります。時刻カラムのデフォルト値は、現在の時刻です。
  2. v_raw_dedupビューを作成します。伝票ごとにraw_idが最も大きい行を1つだけ残し、raw_idカラムを含み、日付と金額とステータスは、正規化された値である必要があります。
  3. v_raw_dedupをorders_finalにロードします。伝票がすでにあれば、値を更新するようにします。ロードしたあとの件数は、異なる伝票の数と同じである必要があります。
  4. content_hashを、md5(coalesce(order_date::text,'') || '|' || coalesce(amount::text,'') || '|' || status)のルールで埋めます。
  5. 同じロードをもう一度実行し、実行前後のorders_finalの件数を、/root/etl/upsert_rerun.txtに、1行ずつ2行で記録します。
  6. staging.orders_rawで、ORD-000100伝票のstatusを別の値に変えたあと、ロードをもう一度実行します。このとき、updated_atがcreated_atより後の行が、ちょうど1件である必要があります。
  7. etl_run_log表を作成します。カラムはrun_id、started_at、inserted_rows、updated_rowsの順で、実行ごとに1行ずつ記録して、最低3行ある必要があります。変更のない実行は、挿入と更新がどちらも0である必要があります。
  8. v_final_checkビューを作成します。カラムはmetric、valueで、total_rows、distinct_refs、changed_rowsの3行を出力します。

参考

自然キーでロード先の表を作る

orders_final表を作成します。カラムはorder_ref、order_date、amount、status、content_hash、created_at、updated_atの順で、order_refが主キーである必要があります。時刻カラムのデフォルト値は、現在の時刻です。

伝票番号を主キーにします。作成時刻と更新時刻のカラムも、一緒に置きます。

重複した伝票を取り除く

v_raw_dedupビューを作成します。伝票ごとにraw_idが最も大きい行を1つだけ残し、raw_idカラムを含み、日付と金額とステータスは、正規化された値である必要があります。

同じ伝票が2回来たなら、あとのものを残します。PostgreSQLのDISTINCT ONが、この作業に合っています。

衝突時に更新するロードを作る

v_raw_dedupをorders_finalにロードします。伝票がすでにあれば、値を更新するようにします。ロードしたあとの件数は、異なる伝票の数と同じである必要があります。

INSERTに衝突処理の句を付けます。ロードの前に、日付とステータスを正規化する必要があります。

内容ハッシュを計算する

content_hashを、md5(coalesce(order_date::text,'') || '|' || coalesce(amount::text,'') || '|' || status)のルールで埋めます。

値を決まった順序でつなげて、ハッシュを作ります。値がなければ、空文字列にします。

変更のない再実行を確認する

同じロードをもう一度実行し、実行前後のorders_finalの件数を、/root/etl/upsert_rerun.txtに、1行ずつ2行で記録します。

同じ入力でもう一度回し、前後の件数を2行で残します。値が同じである必要があります。

遅れて届いた変更を反映する

staging.orders_rawで、ORD-000100伝票のstatusを別の値に変えたあと、ロードをもう一度実行します。このとき、updated_atがcreated_atより後の行が、ちょうど1件である必要があります。

元データの伝票を1つ変えて、もう一度回します。値がそのままの行まで更新されないように、条件をかける必要があります。

実行ログを残す

etl_run_log表を作成します。カラムはrun_id、started_at、inserted_rows、updated_rowsの順で、実行ごとに1行ずつ記録して、最低3行ある必要があります。変更のない実行は、挿入と更新がどちらも0である必要があります。

実行ごとに、挿入件数と更新件数を記録します。何も変わらなかった実行は、両方とも0です。

点検指標のビューを作る

v_final_checkビューを作成します。カラムはmetric、valueで、total_rows、distinct_refs、changed_rowsの3行を出力します。

全体の件数、異なる伝票の数、更新された行数の3つを、名前と値のペアで出力します。