冪等なupsertパイプライン
目標
同じ入力を何回入れても結果が変わらず、内容が実際に変わった行だけが更新される、ロードのパイプラインを作ります。
なぜ重要なのか
パイプラインは必ず失敗します。ネットワークが切れ、デプロイ中に落ち、ソースがデータを遅れて送ってきます。そのたびに、人が「今までどこまで入ったか」を判断しなければならないなら、いつか間違いが起きます。
冪等に設計すれば、その判断が必要なくなります。失敗したなら、そのままもう一度実行すればよいです。核心の仕組みは2つです。自然キーに一意制約をかけて、重複をデータベースに防がせることと、内容ハッシュで実際の変更の有無を判断して、値がそのままの行には触れないことです。
2つ目が特に重要です。条件なしに更新すると、変更のない再実行でも、すべての行の更新時刻が上がります。そうなると、「何が実際に変わったか」がわからなくなり、下流で変更分だけを取得する増分の消費が壊れます。
ステップ
対象はstaging.orders_rawで、正規化ルールは、前のラボと同じです。
orders_final表を作成します。カラムはorder_ref、order_date、amount、status、content_hash、created_at、updated_atの順で、order_refが主キーである必要があります。時刻カラムのデフォルト値は、現在の時刻です。v_raw_dedupビューを作成します。伝票ごとにraw_idが最も大きい行を1つだけ残し、raw_idカラムを含み、日付と金額とステータスは、正規化された値である必要があります。v_raw_dedupをorders_finalにロードします。伝票がすでにあれば、値を更新するようにします。ロードしたあとの件数は、異なる伝票の数と同じである必要があります。content_hashを、md5(coalesce(order_date::text,'') || '|' || coalesce(amount::text,'') || '|' || status)のルールで埋めます。- 同じロードをもう一度実行し、実行前後の
orders_finalの件数を、/root/etl/upsert_rerun.txtに、1行ずつ2行で記録します。 staging.orders_rawで、ORD-000100伝票のstatusを別の値に変えたあと、ロードをもう一度実行します。このとき、updated_atがcreated_atより後の行が、ちょうど1件である必要があります。etl_run_log表を作成します。カラムはrun_id、started_at、inserted_rows、updated_rowsの順で、実行ごとに1行ずつ記録して、最低3行ある必要があります。変更のない実行は、挿入と更新がどちらも0である必要があります。v_final_checkビューを作成します。カラムはmetric、valueで、total_rows、distinct_refs、changed_rowsの3行を出力します。
参考
- 衝突時に更新:
INSERT ... ON CONFLICT (order_ref) DO UPDATE SET ... WHERE 대상.content_hash IS DISTINCT FROM EXCLUDED.content_hash(プレースホルダーは対象テーブル名です) - 挿入と更新の区別:
RETURNING (xmax = 0) AS insertedをCTEで受けて数えればよいです。 - 伝票ごとの最新行:
SELECT DISTINCT ON (order_ref) * FROM ... ORDER BY order_ref, raw_id DESC - よくある間違い1:
DO UPDATEに条件をかけないと、変更のない再実行でも、全行の更新時刻が上がります。 - よくある間違い2: 重複排除なしにロードすると、同じ伝票が1つのバッチの中で2回出てきて、エラーになります。
自然キーでロード先の表を作る
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つを、名前と値のペアで出力します。