TT Lab
Get started
Learn Learning paths Courses

Data Pipelines

Idempotency — The Pipeline Will Certainly Run Again

Continue in TT Lab

One-line summary

Pipelines fail, you rerun them when they fail, and if the result differs at that point the data gets contaminated, so rerun safety is not an option but a premise of design.

Why this was needed

Suppose a batch failed in the early morning. You rerun it in the morning. But if the failure point was in the middle of the load, some of it is already in. If you insert it again as is, duplicates arise, and if you delete everything and insert again, other data that came in during that time is lost too.

If you leave it to a person to judge this situation every time, sooner or later a mistake happens. It is better to make it so that the result is the same no matter how many times you insert in the first place.

How it works

The starting point of idempotency is the natural key. There must be a value that uniquely identifies each row so that you can recognize "what has already come in." Things like a voucher number, an order number, or an event ID. If you apply a unique constraint to this value, the database blocks duplicates for you.

Next is the upsert.

INSERT INTO orders_final (order_ref, amount, status, content_hash)
SELECT ...
ON CONFLICT (order_ref) DO UPDATE
  SET amount = EXCLUDED.amount,
      status = EXCLUDED.status,
      content_hash = EXCLUDED.content_hash,
      updated_at = now()
  WHERE orders_final.content_hash IS DISTINCT FROM EXCLUDED.content_hash;

The final WHERE clause is important. Without it, even in a rerun where not a single thing changed, the updated_at of every row is updated. Then you can no longer tell "what actually changed," and incremental consumption, where downstream takes only the changes, breaks too.

Having a content hash makes the comparison simple. You concatenate the values in a fixed order, compute a hash, and update only when that value differs. Even as the columns grow, the comparison logic stays one line.

Deduplication is also needed. If the same voucher was loaded twice, you must set a rule for which one to keep. Usually you keep the one that came in later. If the source has an increasing key that indicates load order, you just pick the row with the maximum value per voucher.

Swapping at the partition level

There are cases where a row-level upsert does not work. When the source gives a full day's data again in its entirety, or when there is no way to find out which rows were deleted. In this case you take the overwrite unit to be the partition.

BEGIN;
  -- 1) 새 데이터를 임시 표에 적재한다
  CREATE TEMP TABLE stage_20260906 (LIKE orders_final INCLUDING ALL);
  COPY stage_20260906 FROM ...;

  -- 2) 그날 파티션만 통째로 바꿔 끼운다
  ALTER TABLE orders_final DETACH PARTITION orders_20260906;
  ALTER TABLE orders_final ATTACH PARTITION stage_20260906
        FOR VALUES FROM ('2026-09-06') TO ('2026-09-07');
COMMIT;

The key is that it is not "delete and insert" but "prepare and swap in." Between deleting and inserting, there is a period with no data, and someone who queries then sees an empty result. A swap finishes atomically inside a transaction, so there is no such period.

This approach is safe for reruns too. No matter how many times you run it, the final state of that partition is just the result of the last run.

How to recognize unsafe operations

Operations that are not idempotent have something in common. They read the current value and write based on it.

Unsafe Safe
UPDATE t SET n = n + 1 UPDATE t SET n = <계산된 절대값> (the placeholder is a computed absolute value)
INSERT (no constraint) INSERT ... ON CONFLICT DO UPDATE
Append to a file Write the file anew and rename
Publish a message to a queue Publish, but the consumer filters duplicates
Call an external API (payment, etc.) Send an idempotency key with it

The last row is important. When a pipeline calls an external system, it has to rely on that side's idempotency. Most payment APIs accept an Idempotency-Key header. If you call twice with the same key, the second call returns the result of the first as is. Because the key must be the same across retries, you must not create a new one for each request, and you must derive it deterministically from the source data (for example, sha256(주문번호 + 금액), that is, the order number plus the amount).

The watermark and the reprocessing range

In incremental loads you record "how far have I processed." The principle is to update this value after processing. If you update it first, then when it dies midway you skip that interval forever.

읽기 → 변환 → 적재 → (성공 시에만) 워터마크 갱신

And leave some slack in the watermark. If the source allows data that arrives late by event time (late arrival), set the watermark not to the last timestamp but to the last timestamp − a grace period. Then you read the overlapping interval again every time, but you have made it idempotent, so it is not a problem. With idempotency, overlapping reads become free.

What it looks like in the field

Leaving a run log helps a great deal. If you record the number of inserts and updates for each run, a run where nothing changed is left with both at 0, and a run where a late-arriving change was applied shows a higher update count. Just by looking at these two numbers, you can judge whether the pipeline is healthy or something is wrong at the source.

And inside retry code you must put only database operations. If sending mail or an external API call is mixed in, that side effect repeats on every retry. If an external call is needed, that side too must be one that accepts an idempotency key.

What to do in the next lab

You apply a primary key to the natural key, remove duplicate vouchers, and build an upsert that judges whether something changed with a content hash. Then you check how a rerun where nothing changed and a rerun where only one row changed are each recorded.