An Idempotent Upsert Pipeline
Goal
You build a load pipeline whose result does not change no matter how many times you feed in the same input, and that updates only the rows whose content actually changed.
Why it matters
Pipelines certainly fail. The network drops, they die during deployment, and the source sends data late. If a person had to judge "how far did it get in so far" every time, sooner or later a mistake would happen.
If you design it to be idempotent, that judgment is no longer needed. If it failed, you just run it again. There are two key devices. One is to apply a unique constraint to the natural key so that the database blocks duplicates, and the other is to judge whether something actually changed by a content hash and leave rows whose values are unchanged alone.
The second is especially important. If you update unconditionally, even in a rerun with no changes, the update time of every row goes up. Then you can no longer tell "what actually changed," and incremental consumption, where downstream takes only the changes, breaks.
Steps
The target is staging.orders_raw, and the normalization rules are the same as in the previous lab.
- Create the
orders_finaltable. The columns are, in order,order_ref,order_date,amount,status,content_hash,created_at, andupdated_at, andorder_refmust be the primary key. The default of the time columns is the current time. - Create the
v_raw_dedupview. For each voucher, keep only the one row with the largestraw_id, include theraw_idcolumn, and the date, amount, and status must be normalized values. - Load
v_raw_dedupintoorders_final. If the voucher already exists, have it update the values. The count after loading must equal the number of distinct vouchers. - Fill
content_hashusing the rulemd5(coalesce(order_date::text,'') || '|' || coalesce(amount::text,'') || '|' || status). - Run the same load once more, and record the
orders_finalcount before and after the run in/root/etl/upsert_rerun.txtas two lines, one value per line. - In
staging.orders_raw, change thestatusof voucherORD-000100to a different value and then run the load again. At that point there must be exactly 1 row whoseupdated_atis later than itscreated_at. - Create the
etl_run_logtable. The columns are, in order,run_id,started_at,inserted_rows, andupdated_rows, and you record one row per run so that there are at least 3 rows. A run with no changes must have both inserts and updates at 0. - Create the
v_final_checkview. The columns aremetricandvalue, and it produces three rows:total_rows,distinct_refs, andchanged_rows.
Reference
- Update on conflict:
INSERT ... ON CONFLICT (order_ref) DO UPDATE SET ... WHERE 대상.content_hash IS DISTINCT FROM EXCLUDED.content_hash(the placeholder stands for the target table) - Distinguishing insert from update: receive
RETURNING (xmax = 0) AS insertedin a CTE and count. - Latest row per voucher:
SELECT DISTINCT ON (order_ref) * FROM ... ORDER BY order_ref, raw_id DESC - Common mistake 1: if you do not put a condition on
DO UPDATE, the update time of all rows goes up even in a rerun with no changes. - Common mistake 2: if you load without deduplication, the same voucher appears twice within one batch and an error occurs.
Create the target table with a natural key
Create the orders_final table. The columns are, in order, order_ref, order_date, amount, status, content_hash, created_at, and updated_at, and order_ref must be the primary key. The default of the time columns is the current time.
Make the voucher number the primary key. Also include the creation time and update time columns.
Remove duplicate vouchers
Create the v_raw_dedup view. For each voucher, keep only the one row with the largest raw_id, include the raw_id column, and the date, amount, and status must be normalized values.
If the same voucher came twice, keep the later one. PostgreSQL's DISTINCT ON fits this job.
Build a load that updates on conflict
Load v_raw_dedup into orders_final. If the voucher already exists, have it update the values. The count after loading must equal the number of distinct vouchers.
Attach a conflict-handling clause to the INSERT. You must normalize the date and status before loading.
Compute the content hash
Fill content_hash using the rule md5(coalesce(order_date::text,'') || '|' || coalesce(amount::text,'') || '|' || status).
Concatenate the values in a fixed order to make the hash. If a value is missing, use an empty string.
Confirm a rerun with no changes
Run the same load once more, and record the orders_final count before and after the run in /root/etl/upsert_rerun.txt as two lines, one value per line.
Run once more with the same input and leave the counts before and after as two lines. The values must be equal.
Apply a late-arriving change
In staging.orders_raw, change the status of voucher ORD-000100 to a different value and then run the load again. At that point there must be exactly 1 row whose updated_at is later than its created_at.
Change one voucher in the source and run again. You must put a condition so that even rows whose values are unchanged are not updated.
Leave a run log
Create the etl_run_log table. The columns are, in order, run_id, started_at, inserted_rows, and updated_rows, and you record one row per run so that there are at least 3 rows. A run with no changes must have both inserts and updates at 0.
Record the number of inserts and updates for each run. A run where nothing changed has both at 0.
Create a check-metrics view
Create the v_final_check view. The columns are metric and value, and it produces three rows: total_rows, distinct_refs, and changed_rows.
Produce three things as name and value pairs: the total count, the number of distinct vouchers, and the number of updated rows.