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.

8 min read • View on GitHub • More from dkgupta2510

An editorial illustration of a warehouse conveyor line where an inspector stamps boxes REJECTED and sets them on a side shelf, while the same rejected boxes also continue down the main belt and get stacked onto a loading pallet. It depicts a quality-control station that works perfectly and still ships the rejected goods, mirroring a validation stage whose output is never read.
The inspector rejects the box and shelves it. The box reaches the loading dock anyway. That is the whole architecture of ETL-Data-Pipeline in one drawing.
Key Takeaways

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.

Step a single record through all five stages and watch its own safeguard fail to fire. Toggle the broken edge to see what the corrected pipeline would have done.

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.

A close-up editorial illustration of a single wall clock face carrying two sets of hands offset by five and a half hours, one pair drawn in solid ink and the other in lighter crosshatching, with two paper timecards taped below showing the same shift start in two different notations. It visualizes how a naive timestamp treated as UTC and converted to IST moves every hour of the day.
One clock, two time systems. The solid hands are what the generator wrote. The hatched hands are what the transform believes it read.

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 codeWhat it looks like it doesWhat 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 numericConflates "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 timestampsDrops 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 missingSwallows 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 saysCode actually does
Five-stage pipeline: generate, validate, deduplicate, load, queryThe 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 recordsTransform 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.

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.