The Validation Step That Never Runs: Inside ETL-Data-Pipeline
A tidy five-stage Python pipeline generates messy call-center data, quarantines the bad rows, and logs exactly how many it rejected. Then it loads them anyway.
- ETL-Data-Pipeline validates its records, quarantines the bad ones, and logs the rejection count, then rereads the raw file and loads the rejected rows anyway.
- Every bug in this repo shares one trait: it looks handled, which is exactly why the audit log and the database contradict each other.
- A pipeline that runs cleanly once proves the happy path, not the contracts between stages.
- The strongest design choice in the codebase is a single flag captured before imputation, and it is the pattern worth stealing.
A Log That Says 75 Records Were Rejected
Run the pipeline once and it finishes clean. Five scripts, no errors, a database you can query. The ingestion_log table records a row for the run, and one of its columns is rejected_count. It says 75.
Then open the calls table. The rejected records are in it.
Not similar records. The same rows, the ones the validator flagged for a missing required field, sitting in the serving table with NULLs or coerced values where the missing data should be. Two artifacts from the same execution disagree about what happened, and both are correct about their own narrow job. The log truthfully reports what the validator decided. The table truthfully reports what the loader inserted. Nothing in between enforced that the second should follow the first.
That gap is the story. It is not a crash, not a stack trace, not a wrong number in a report. It is a pipeline that generates its own audit trail and then quietly ignores it.
The Validator Nobody Calls
The README describes five stages in order: generate, validate, deduplicate, load, query. The code delivers all five. The problem is the seam between stage two and stage three.
pipeline.py reads raw_calls.json and walks each record through an ordered chain of checks: missing required field, duplicate call_id, wrong type, unparseable timestamp. Records that fail any check get a reject_reason and land in rejected_log.json. Records that pass land in clean_valid_records.json. This is a real validator doing real work, and it produces a correct quarantine file.
Then deduplicate.py runs, and its first line of input is this:
# pipeline.py
raw = pd.read_json("raw_calls.json")
# ... validates, writes clean_valid_records.json and rejected_log.json
# deduplicate.py
df = pd.read_json("raw_calls.json") # <- the raw file again
# ... transforms, writes transformed_calls.json
Both scripts open the same file. clean_valid_records.json is written and never read. It is a dead-end branch of the dependency graph, and because it is a file on disk rather than a function call, nothing complains.
Trace one record and the consequence is concrete. A call arrives missing its retry_flag. The validator flags it, writes it to the reject log, and excludes it from the clean file. The transform, reading the raw file, sees the record again, fills the missing flag with NaN, and passes it downstream. The loader then runs bool(row["retry_flag"]), and bool(float('nan')) evaluates to True, so the record enters the database marked as a retry. The rejected record was rejected, logged as rejected, and loaded.
The File Names Point the Wrong Way
The reason nobody caught the orphaned file is sitting in the filenames. The script that validates is called pipeline.py. The script that transforms is called deduplicate.py.
Read that way, the README's run order sounds reasonable. A file named pipeline.py sounds like a driver or a core stage. A file named deduplicate.py sounds like a cleanup step that runs after. So a reader follows the instructions, sees five scripts execute in the listed order, and assumes the chain is intact.
The actual responsibilities are inverted. pipeline.py is a validator that ends in a dead-end branch. deduplicate.py is the real transform, and it starts from raw input. Because the names describe an earlier design, the documented flow and the real dependency graph run in different directions, and the names made the difference invisible.
This is a generalizable failure, not a one-off. File and function names are the cheapest documentation a codebase has, and they are also the most trusted, because readers use them to decide what to read closely. A name that describes a previous version of the code does not just mislead a reader. It tells the reader they can stop looking.
Every Timestamp Is 5.5 Hours Off
The second silent bug never touches the record count. It touches every record.
generate.py builds timestamps with datetime(2024, 1, 1) + timedelta(...) and writes them out with .isoformat(). No timezone offset, no Z, no zone name. They look like local business hours, and they are generated between 08:00 and 20:59.
The transform then calls pd.to_datetime(df["start_time"], utc=True). Passing utc=True on a naive series tells pandas to assume those values are already UTC, and it attaches the offset accordingly. The next line converts them to Asia/Kolkata, which is UTC+05:30.
So an 08:15 call becomes 13:45 IST. The generated business-hours window of 08:00 to 20:59 shifts to roughly 13:00 to 02:00. Every downstream time feature inherits the shift: call_hour, is_weekend (which reads the IST date), and the peak-callback-hour query in samplequeries.txt that the whole demo is built to showcase. The pipeline is internally consistent and semantically wrong, and nothing in the code, the README, or the (declared but unwritten) test suite catches it.
Four Failures That Look Like Handling
The orphaned validator and the timezone shift are the headline findings. Four smaller ones share the same shape: each line reads as if the author thought about the failure mode and dealt with it. Each one actually converts a loud failure into a quiet one.
| Line of code | What it looks like it does | What it actually does |
|---|---|---|
| <code>bool(row["retry_flag"])</code> in <code>sqlload.py</code> | Casts the flag to a clean boolean for storage | <code>bool(float('nan'))</code> is <code>True</code>, so a missing flag loads as a positive one |
| <code>df["amount_promised"].fillna(0)</code> in <code>deduplicate.py</code> | Handles nulls so the column is always numeric | Conflates "no amount promised" with "amount unknown", and query #3 averages those zeros |
| <code>errors="coerce"</code> then <code>dropna</code> in <code>deduplicate.py</code> | Safely discards malformed timestamps | Drops the rows without writing them to <code>rejected_log.json</code>, so the audit trail is incomplete |
| <code>except:</code> around the reject-log read in <code>sqlload.py</code> | Keeps the load running if the log is missing | Swallows every error including a corrupt file, and the run reports success |
None of these are exotic. Coercing, filling, dropping, and catching are the four most common verbs in data code, and each has a legitimate use. The trap is that all four default to silence. A coercion that fails loudly gets fixed. A coercion that succeeds quietly ships.
What the Author Got Right
A critique that only lists bugs is not a fair read, and this codebase earns some credit. The validator quarantines rejected records with a reason attached instead of dropping them. That is the right instinct for auditability, and it is why the contradiction in this article is even visible: the reject log exists, so you can compare it against the table. A pipeline that silently dropped bad rows would have hidden the same bug behind a smaller record count.
The strongest single choice is this pair of lines in deduplicate.py:
df["is_amount_imputed"] = df["amount_promised"].isnull()
df["amount_promised"] = df["amount_promised"].fillna(0)
The imputation flag is computed before the fill. That preserves the information a downstream analyst needs to tell a real zero from a substituted one, which is exactly the distinction that fillna(0) alone destroys. It is a small pattern and it is the one worth copying into other codebases.
The generator is also seeded with random.seed(42) and injects faults at known rates: 15 percent missing fields, 5 percent duplicates, 3 percent malformed timestamps. That makes the whole pipeline testable in principle, since you know roughly how many records should be rejected and why. The project does not exploit this yet, but the design invites it.
The self-assessment is honest about storage and honest about load semantics. What it leaves out is the row-by-row iterrows() loop with a per-row execute(), and the fact that dedup-by-sort loads the entire dataset into memory. Those are the first things that break at volume, before the choice of database matters.
Idempotent, Except It Crashes on the Second Run
The README advertises idempotent inserts. That word is doing more work than the code supports.
For the calls table, the pattern is DROP TABLE IF EXISTS calls, then CREATE TABLE, then INSERT OR REPLACE per row. Drop-and-recreate does make that table reproducible: run it twice and you get the same contents. The idempotency comes from the drop, not from the insert semantics, which is a fine distinction but not the real problem.
The real problem is one table over. ingestion_log is created with a plain CREATE TABLE and is never dropped. On the second run, the statement raises sqlite3.OperationalError: table ingestion_log already exists and the load stops. The fix is one word.
# first run: fine. second run: OperationalError
conn.execute("CREATE TABLE ingestion_log (...)")
# the fix
conn.execute("CREATE TABLE IF NOT EXISTS ingestion_log (...)")
The lesson is not about SQLite. It is that idempotency is a property you verify by running the thing twice, and a label you apply to code that has not been run twice is a claim, not a guarantee. The README's word and the code's behavior diverged, and only execution would have shown it.
Claim vs. Code
Pulled together, the README and the code tell two different stories about the same five scripts.
| README says | Code actually does |
|---|---|
| Five-stage pipeline: generate, validate, deduplicate, load, query | The validate stage runs but its output is never read; <code>deduplicate.py</code> rereads <code>raw_calls.json</code> |
| Run order: <code>generate.py</code> then <code>pipeline.py</code> then <code>deduplicate.py</code> then <code>sqlload.py</code> then <code>queries.py</code> | The real dependency graph makes <code>pipeline.py</code> a dead-end branch, so the order is documented but not chained |
| Idempotent inserts | <code>DROP TABLE</code> plus <code>INSERT OR REPLACE</code> on <code>calls</code>; plain <code>CREATE TABLE</code> on <code>ingestion_log</code>, so the second run crashes |
| Rejected records are quarantined in <code>rejected_log.json</code> | Rejected records still load into <code>calls</code> with <code>NULL</code>s or coerced values |
| Transform reads validated records | Transform reads <code>raw_calls.json</code> directly |
It is worth being precise about what this repo is. It is a learning artifact, a portfolio piece for someone practicing the shape of an ETL job, and judged that way it does a lot: it has a real validator, a real quarantine log, a real transform, real analytics queries, and an honest list of what would change at scale. It is not competing with Airflow, Dagster, or dlt, and it does not claim to.
But the missing pieces are instructive, because they are exactly the pieces that would have caught the bugs. An orchestrator would have made clean_valid_records.json a declared input to the transform, so an unread output would be a wiring error rather than a silent file on disk. A schema contract would have rejected the missing retry_flag at the load boundary. Upsert semantics with a real key would have made the rerun safe by construction. The repo's bugs live in the space those tools occupy.
The Gap Between Runs and Runs Reliably
The single most useful thing this codebase demonstrates is that a successful run is weak evidence. Everything executed, nothing threw, the database has rows in it. A pipeline can do all of that while quietly violating every contract it appears to honor. What separates a pipeline that runs from one that runs reliably is not the number of stages. It is whether the stages are actually connected.
- Each stage consumes exactly what the previous stage produced. If a file is written and never read, the graph has a dead branch, and a dead branch is a bug the type system does not catch.
- Every rejected row is accounted for in one place, and that place is the same place the load consults. A reject log that nothing reads is a report, not a control.
- A rerun is safe. Not labeled safe, verified safe, by running it twice and watching what happens.
- A default is a decision. Coercing, filling, dropping, and catching all choose an outcome for you, and the honest version of each one writes down what it chose.
None of this requires a framework. It requires knowing which file each script reads, and being suspicious of any output that nothing else opens. That is the whole discipline, and this repo is worth reading because it shows what happens when it slips: a validator that runs, a log that tells the truth, and a database that disagrees with both.





