TT Lab
Get started
Learn Learning paths Courses

Data Pipelines

Operating a Quarantine: Accumulate, Break, Requeue

Continue in TT Lab

Goal

You write expectations as a file, and build a gatekeeper qgate.py that filters a drop at the row level on every run and piles it separately into a sqlite fact table and quarantine table. You cut off a run that is wrong too much in its entirety, put fixed rows back in without double counting, catch aged quarantines, and report run state and data state separately.

Why it matters

Most people go as far as leaving unusable rows in quarantine instead of discarding them. The problem is what comes next. Quarantine piles up again on every run, and when you open it half a year later it holds 40,000 rows, and all that time the aggregate had been going out with that many missing. The pipeline succeeded every day — success means it finished without errors, not that it put in everything that should have been put in. One row being wrong and a whole file being wrong are also different events. If upstream sends the column order changed, nearly every row is off, and if you quarantine at the row level at that point, ten thousand rows go into the quarantine table and the aggregate goes out empty. For such a file, it is better to cut off the run without putting in a single row. The criterion for splitting is a ratio, and that ratio is set not from a round number but from the usual failure ratio. And you need a loop that puts fixed rows back in. Reinsertion is something people run by hand, so it will certainly be run several times, and at that point it is easy to go into the fact table twice, or to close an already closed quarantine again so that "resolved 3" is reported twice. The grader does not trust the text you write down. It sets up the drops and expectations that the grader created in a temporary directory, actually runs your gatekeeper, and then directly opens the created sqlite file and compares the facts, quarantine, and run records with values the grader counted directly. The thresholds, shop names, and amounts change on every run.

Steps

  1. Create and run /root/quarantine/gen_drops.py to create four drops under /root/quarantine/drops and /root/quarantine/expectations.json.
  2. In /root/quarantine/qgate.py, create check so that it only counts against the expectations.
  3. Add run to pile facts, quarantine, and run records into sqlite.
  4. Add run --max-fail-ratio to cut off a run that is wrong too much in its entirety.
  5. Add requeue to put fixed rows back in while keeping the numbers from changing even if you put them in twice.
  6. Add aging to catch aged quarantines.
  7. Add status to report run state and data state separately.
  8. Run your four drops in order to create /root/quarantine/pipeline.db, and write /root/quarantine/status.json and /root/quarantine/quarantine_report.md.

Reference

Keep the drops and expectations as files

Create and run /root/quarantine/gen_drops.py to create four drops under /root/quarantine/drops and /root/quarantine/expectations.json. The expectations must use all five types, and one drop must violate the expectations on more than half of its rows.

If you keep the expectations as a file rather than code, you can talk with upstream over that file, and the name of the violated rule becomes the quarantine reason as is. Mix into the drops lines with a status not in the list, a negative amount, a quantity outside the range, an empty shop, a voucher number that does not follow the format, and a voucher number identical to an earlier line. Three must have fewer than half their lines in violation and one must have more than half.

Only count against the expectations

In /root/quarantine/qgate.py, create check --drop <파일> --expect <파일> (the placeholders stand for the files) so that it produces drop, rows, passed, failed, fail_ratio, and by_rule as JSON.

failed is the number of lines that violate one or more rules and by_rule is the number of lines violating each rule, so the sums differ — because one line can violate two rules. In by_rule, also put in as 0 the rules nobody violated. If a value is empty, treat it as failing every rule except not_null.

Pile facts and quarantine separately

Add run --drop <파일> --db <파일> --expect <파일> (the placeholders stand for the files) so that rows that pass go into facts as an update, rows that violate open in quarantine for each violated rule, and the run is left in runs. Even if the same reason for the same line comes again, quarantine opens only once.

The column names and order of the three tables are in the reference section. The fact table has order_id as the primary key, so put in by update. Before you open quarantine, check whether there is a line still open with the same order_id and rule — if you do not, the same line piles up on every run.

Cut off a run that is wrong too much in its entirety

Add run --max-fail-ratio R to cut off the run if fail_ratio exceeds the threshold. A cut-off run is left in runs as broken, puts not a single line of facts or quarantine, and the exit code is 5.

One row being wrong and a whole file being wrong are different events. Set the threshold not from a round number but from the usual failure ratio — measure the ratios of the four drops with check and then decide. If you put in half and stop when cutting off, the next person has to judge how far it got in.

Put the fixed rows back in

Add requeue --db <파일> --fixes <파일> --expect <파일> (the placeholders stand for the files) to check the fixed rows against the expectations again, and if they pass, put them into the fact table as an update and close only the quarantines of that voucher that are still open. Even if you run the same reinsertion twice, the fact count and the number of closed quarantines must not grow.

Reinsertion is something people run by hand, so it will certainly be run several times. If you put in by insert, the facts double, and if you do not apply an open condition when closing, "resolved 3" gets reported twice. The resolved of the second run must be 0. A row whose voucher number itself is broken cannot be fixed by reinsertion — because if you fix the number, it becomes a different row.

Catch aged quarantines

Add aging --db <파일> [--max-age-runs N] (the placeholder stands for the file) to measure the age of open quarantines as latest_run_id - first_run_id, count those at or above max_age_runs, and produce open, aged, by_rule, oldest, and alert.

If you look only at counts, you cannot see what does not shrink. It is easier to handle age by number of runs than by time — it matches human intuition that if the pipeline does not run, the age does not advance either. oldest is, among the open ones, the one with the smallest first_run_id.

Report "it succeeded" and "it is correct" separately

Add status --db <파일> [--max-age-runs N] (the placeholder stands for the file) to produce run state (runs) and data state (data) separately and answer with the two values runs_ok and data_ok.

That the run succeeded and that the data is correct are different axes. If you lump them into one field, you see neither. runs_ok is whether there are runs and the last run was not cut off, and data_ok is whether both open quarantines and aged quarantines are 0 — days when both are true are rare.

Run a week with your own drops and report

Run your four drops in date order to create /root/quarantine/pipeline.db, leave the answer of status in /root/quarantine/status.json, and write /root/quarantine/quarantine_report.md in four sections. Set the stop threshold by measuring the usual failure ratio.

If you set the threshold from the usual value, only the day whose violations exceed half gets cut off. Do not write status.json by hand; save the answer of the status command as is — the grader directly opens the database and compares. In the report, write the fact count and the open quarantine count as numbers.