TT Lab
Get started
Learn Learning paths Courses

Data Pipelines

Batch Loading and Watermarks

Continue in TT Lab

Goal

You run by hand one cycle of a batch pipeline: extract, load, watermark recording, incremental processing, and rerun verification.

Why it matters

The points where incidents happen in a batch pipeline are almost fixed. They are boundary conditions and reruns.

For boundary conditions, a one-character difference between > and >= makes one row duplicated or missed on every run. For a batch that runs once a day, 365 rows go wrong in a year, and nobody knows in the meantime. So it is good to put a verification step inside the pipeline that compares the count right after extraction against the source.

Reruns are even more important. A pipeline will certainly fail, and when it fails you run it again. If the result changes at that point, the data becomes contaminated. So you must design it so that a unique constraint is applied to the load target table and conflicts are ignored or updated. In this lab you run the same incremental job twice and directly check that the count does not change.

Steps

The working directory is /root/etl.

  1. Extract the id, ordered_at, status, and total_amount of orders whose ordered_at is in January 2025 to /root/etl/orders_2025_01.csv. Include the header. The reference times are greater than or equal to TIMESTAMPTZ '2025-01-01 00:00:00+09' and less than TIMESTAMPTZ '2025-02-01 00:00:00+09'.
  2. Check that the number of lines in the file equals the number of orders in that period plus one header line.
  3. Create the orders_archive table. The columns are, in order, id, ordered_at, status, and total_amount, and id must have a primary key or a unique constraint.
  4. Load the CSV you created in step 1 into orders_archive.
  5. Create the etl_watermark table. The columns are job_name and last_ordered_at, and in the row where job_name is orders_archive, record the maximum ordered_at of the data you loaded.
  6. Incrementally load the orders from after the watermark up to less than TIMESTAMPTZ '2025-03-01 00:00:00+09' and update the watermark. After loading, orders_archive must contain all the orders of both January and February.
  7. Run the same incremental job once more. Record the count of orders_archive before and after the run in /root/etl/rerun.txt as two lines, one value per line. The two values must be equal and there must be no duplicate rows.
  8. Save the count per status of orders_archive in /root/etl/report.tsv. Separate the status and the count with a tab, and do not include a header.

Reference

Extract the January orders to a CSV

Extract the id, ordered_at, status, and total_amount of orders whose ordered_at is in January 2025 to /root/etl/orders_2025_01.csv. Include the header. The reference times are greater than or equal to TIMESTAMPTZ '2025-01-01 00:00:00+09' and less than TIMESTAMPTZ '2025-02-01 00:00:00+09'.

psql's \copy writes to a client-side file. There is an option that includes the header.

Verify the extracted count

Check that the number of lines in the file equals the number of orders in that period plus one header line.

The number of lines in the file must be the number of data rows plus one header line. Check whether you included only one side of the boundary condition.

Create the load target table

Create the orders_archive table. The columns are, in order, id, ordered_at, status, and total_amount, and id must have a primary key or a unique constraint.

Rerun safety is decided here. You need a constraint that makes it impossible for the same order to come in twice.

Load the CSV into the table

Load the CSV you created in step 1 into orders_archive.

\copy works in the opposite direction too. Do not forget the option that skips the header line.

Record the watermark

Create the etl_watermark table. The columns are job_name and last_ordered_at, and in the row where job_name is orders_archive, record the maximum ordered_at of the data you loaded.

You leave in a table how far you have processed. The value must be the last timestamp of the data you loaded.

Incrementally load the February data

Incrementally load the orders from after the watermark up to less than TIMESTAMPTZ '2025-03-01 00:00:00+09' and update the watermark. After loading, orders_archive must contain all the orders of both January and February.

Pick and insert only the data after the watermark, and update the watermark when finished.

Confirm it is the same when run twice

Run the same incremental job once more. Record the count of orders_archive before and after the run in /root/etl/rerun.txt as two lines, one value per line. The two values must be equal and there must be no duplicate rows.

Run the same incremental job once more and leave the counts before and after in the file as two lines. You need a clause that ignores conflicts.

Create the load result report

Save the count per status of orders_archive in /root/etl/report.tsv. Separate the status and the count with a tab, and do not include a header.

Save the count per status separated by a tab. Leave only the values, without a header.