Eight posts ago you were reading a CSV somebody handed you. Today the pipeline reads itself: it pulls the city's 311 feed every morning, checks every record against a contract it can actually run, quarantines what fails, stores what passes, and proves — with arithmetic, not vibes — that every record is accounted for. This is the milestone the whole stage was building toward.

Assumes: Posts 1–8. Installs: pip install requests pydantic pandas. Everything below ran against the real 311 open-data feed; every number shown is computed, not asserted.

Monday, 8:05 AM. Maria: "I need the 311 intake to run itself every morning, and I need to trust it without reading it."

"Trust it without reading it" is the whole milestone in one sentence. She doesn't want a script — she wants a report: a few lines she can glance at over coffee that prove the night's pull worked, and a guarantee that anything weird got quarantined instead of silently dropped. If the pipeline can't produce that proof, it isn't done.

The toolkit is everything you've built: Post 6's fetching, Post 5's UTC instants, Post 4's cleaning, Post 3's SQL, Post 2's boundary instinct. Today they stop being separate lessons and become one machine.

Before you code: clarify the ask

You: "Trust it without reading it — so what does the morning report need to say for you to believe it?"

Maria: "How many records came in, how many landed in the database, and what happened to the rest. If those three numbers don't add up, I want to know before the ops review does."

You: "And when the feed changes shape overnight — new field, renamed column — what should the pipeline do?"

Maria: "Not quietly absorb it. Quarantine first, ask questions later."

You: "And reruns? If the 6 AM run fails halfway and I rerun at 6:20?"

Maria: "It should just... work. No duplicates."

Input: the city's 311 open-data feed, every morning
Output: a SQLite database Maria's reports read from, plus a reconciliation report: fetched vs stored vs quarantined, all adding up
Deadline: running unattended by next Monday

Three requirements fell out of that conversation, and they're the spec: reconciliation (every record accounted for), quarantine on drift (reject what you don't understand), idempotent reruns (running twice changes nothing).

The minimal concept

A data pipeline has five stages, and each one earns Maria's trust differently:

StageDoesTrust mechanism
ExtractPull the feed (Post 6: sessions, timeouts)Half-open date window — no record counted twice, none skipped
ValidateCheck every record against the contractThe contract is runnable: schema drift is rejected by code, not by hope
CleanNormalize (Post 4: strip, case, UTC instants)Cosmetic only — cleaning never rescues a record the contract rejected
LoadInsert into SQLite (Post 3)Dedupe key as PRIMARY KEY — reruns are no-ops by construction
ReconcileCount everything, hash everythingfetched == stored + quarantined + duplicates, and the hashes match

Four ideas carry the milestone:

The data contract is a runnable artifact. A schema document in a wiki is a wish. A Pydantic model in the pipeline is a bouncer: every record is validated before it touches the database, and anything the model doesn't understand goes to the dead-letter queue. When the vendor renames a field overnight, you find out from the quarantine count, not from a broken dashboard three weeks later.

The DLQ is not a trash can. Quarantined records are stored with their raw payload, the validation error, and a timestamp. Nothing is silently dropped — "we threw away 16 rows" is a reportable fact with evidence attached, which is exactly what Maria asked for.

Idempotency is a design choice, not a retry policy. The dedupe key is the PRIMARY KEY. Inserting the same record twice is a no-op at the database level, so a failed 6 AM run can simply be rerun at 6:20. You don't need distributed locks; you need a key.

The audit log is the pipeline's memory. Every run writes one row: window, counts, hashes, timestamp. When Maria asks "what happened Tuesday?", the answer is a query, not a archaeology expedition.

Build it, part 1: extract and the v1 contract

Extraction is Post 6's pattern with a half-open window — [Sep 20, Sep 23), so a record stamped 23:59:59.500 can't fall through the cracks (Post 5's boundary lesson, applied to queries):

def fetch_311(session, start, end, limit=1000):
    resp = session.get(
        "https://data.cityofnewyork.us/resource/erm2-nwe9.json",
        params={
            "$select": "unique_key,created_date,agency,complaint_type,"
                       "descriptor,borough,incident_zip,city",
            "$where": f"created_date >= '{start}T00:00:00' "
                      f"AND created_date < '{end}T00:00:00'",
            "$limit": str(limit),
        },
        timeout=(5, 30),
    )
    resp.raise_for_status()
    return resp.json()   # 1000 real rows for the 3-day window

Then the contract — v1, the straightforward version. Timestamps go through Post 5's parse_ts (feed stamps NYC wall time; we store UTC), and the feed's field names get aliased to ours:

from pydantic import BaseModel, ConfigDict, Field, field_validator
from typing import Optional

class ServiceRequestV1(BaseModel):
    model_config = ConfigDict(populate_by_name=True)

    unique_key: str
    created_at: datetime = Field(alias="created_date")  # feed name != our name
    agency: str
    complaint_type: str
    descriptor: Optional[str] = None
    borough: Optional[str] = None
    incident_zip: str = Field(min_length=5, max_length=5)
    city: Optional[str] = None

    @field_validator("created_at", mode="before")
    @classmethod
    def _ts(cls, v):
        try:
            return parse_ts(v)   # naive in -> UTC instant out; garbage raises
        except (ValueError, TypeError) as e:
            raise ValueError(f"bad timestamp {v!r}") from e

The validation gate — Post 2's quarantine instinct, now as a function. Note what it catches: only to add context (the DLQ entry) — the error itself propagates into the quarantine record:

def validate(rows, model):
    """Contract gate: valid models out, DLQ entries for the rest."""
    valid, dlq = [], []
    for raw in rows:
        try:
            valid.append(model(**raw))
        except ValidationError as e:
            dlq.append({"raw": raw, "error": e.errors()[0]["msg"]})
    return valid, dlq

Cleaning is Post 4's toolkit doing what it's for — cosmetic normalization, nothing structural:

def clean(valid):
    df = pd.DataFrame([r.model_dump() for r in valid])
    for col in ["agency", "complaint_type", "descriptor", "borough", "city"]:
        df[col] = df[col].str.strip()
    df["borough"] = df["borough"].str.upper()
    df["created_at"] = pd.to_datetime(df["created_at"], utc=True)
    # ...back to models, NaN -> None (the contract speaks None, not NaN)

Loading is Post 3's SQL with the dedupe key as PRIMARY KEY. v1 dedupes on unique_key — the feed's own id, the obvious choice:

CREATE TABLE service_requests(
  fingerprint TEXT PRIMARY KEY,   -- v1: this holds unique_key; v2: content hash
  unique_key TEXT NOT NULL,
  created_at TEXT NOT NULL,        -- UTC ISO instant
  agency TEXT NOT NULL,
  complaint_type TEXT NOT NULL,
  descriptor TEXT, borough TEXT, incident_zip TEXT, city TEXT,
  ingested_at TEXT NOT NULL
);
-- INSERT OR IGNORE: the PRIMARY KEY makes reruns no-ops by construction

First run, v1, against the real feed — 1000 rows in, plus 5 deliberately poisoned rows (injected to demonstrate the DLQ without waiting for the feed to misbehave):

fetched: 1005 | valid: 989 | quarantined: 16 | stored: 989
adds_up: true   (989 + 16 = 1005)
hash_match: true

The arithmetic holds. But look at the DLQ: 16 quarantined, and only 5 were my poison. Eleven real rows failed the contract. Before you read on — that's the pipeline working correctly and the contract being wrong. The question is why.

Dev, mid-build, over your shoulder: "Bad news — the ZIP codes in the feed are unreliable. Some are missing, some are wrong."

There it is. Eleven real 311 requests — a DOB plumbing complaint in Manhattan among them — quarantined because their incident_zip was missing, and v1 demanded a 5-digit string. The data was fine; the contract was stricter than reality. And Dev drops the second shoe: when the feed "corrects" a record, it re-emits it under a new unique_key — so v1's dedupe key would double-count every correction.

Two fixes, both in the contract: stop treating the zip as load-bearing, and stop treating unique_key as identity.

Build it, part 2: the v2 contract

The updated model. incident_zip becomes Optional with a normalizing validator — a 5-digit zip is kept, anything else becomes NULL instead of failing the row. The row survives; only the unreliable value is quarantined, in place:

class ServiceRequest(BaseModel):
    # ... same fields, except:
    incident_zip: Optional[str] = None  # unreliable: normalize, never reject

    @field_validator("incident_zip", mode="before")
    @classmethod
    def _zip(cls, v):
        if v is None:
            return None
        z = str(v).strip()
        return z if (len(z) == 5 and z.isdigit()) else None

And identity moves from the feed's key to a content fingerprint — a hash of the stable fields, deliberately excluding both unique_key (re-emissions get new ones) and incident_zip (can't be trusted):

def fingerprint(self) -> str:
    core = "|".join([
        self.created_at.isoformat(), self.agency, self.complaint_type,
        self.descriptor or "", self.borough or "", self.city or "",
    ])
    return hashlib.sha256(core.encode()).hexdigest()

One honest tradeoff, stated plainly: the fingerprint is conservative. Two genuinely distinct requests filed in the same second with identical fields would merge — the run found 10 such pairs in 1000 rows (two HPD pest complaints in the Bronx, same second). unique_key is still stored on every row for forensics, so nothing is unrecoverable. Dedupe keys are always a bet; this one bets that re-emission double-counts are worse than same-second merges.

Same intake, v2 contract:

fetched: 1005 | valid: 1000 | quarantined: 5 | duplicates: 10 | stored: 990
adds_up: true   (1000 + 5 = 1005; 990 + 10 = 1000)
hash_match: true

The 11 zip-less rows are back in — incident_zip NULL, row intact. The DLQ holds exactly the 5 poison rows. The 10 duplicates are the same-second merges. Every record accounted for.

Break it, three ways

A pipeline that only works on good days is a demo. Here's what the bad days look like — all three reproduced against the real feed.

1. The poison rows. The 5 injected rows landed in the DLQ with their raw payloads and real validation errors — this is what "quarantine with evidence" means:

POISON2 -> Value error, bad timestamp 'yesterday-ish'
POISON3 -> Input should be a valid string        # agency: None
POISON4 -> Value error, bad timestamp {'epoch': 1726870000}
POISON5 -> Field required                        # complaint_type missing

Each DLQ row stores the raw JSON, the error, and the timestamp. Maria never has to take "we dropped some rows" on faith — she can read exactly which rows and why.

2. Schema drift. The vendor renames complaint_type to complainttype overnight — no announcement, just a quiet Tuesday. The contract is a bouncer, not a janitor:

fetched: 1000 | valid: 0 | quarantined: 1000 | stored: 0
adds_up: true
rows in service_requests after the drifted batch: 0

Every single row rejected, zero rows polluting the table, and the reconciliation still adds up. The 6 AM alert isn't "the dashboard is wrong" — it's "quarantine spiked to 1000, the feed changed shape." That's the failure mode you want: loud, contained, and reversible. You fix the contract (or call the vendor), rerun, and the PRIMARY KEY dedupe means the rerun only fills the gap.

3. The re-emissions. Dev's second shoe, simulated from three real rows: same request, new unique_key, "corrected" zip. Under v1's key-based dedupe, all three would have stored as new rows — the probe run stored 992 instead of 989, silently inflating the counts Maria reports. Under the fingerprint: caught, all three, zero new rows stored. The correction is absorbed; the double-count never happens.

Productionize: the 6 AM machine

The pipeline now needs to run itself. Three pieces, all from the last two posts' toolkit:

Schedule. One cron line — the pipeline already takes --start/--end, so "yesterday" is just arguments:

30 6 * * * /opt/cityops/venv/bin/python /opt/cityops/intake.py \
  --start yesterday --end today >> /var/log/cityops/intake.log 2>&1

CLI + logging. argparse for the window and limit, structured logging with the run_id on every line — so when the 6:20 rerun happens, its log lines don't interleave ambiguously with the 6:00 attempt's. (This is Post 8's packaging-and-logging harness, with the pipeline living inside it.)

The test suite. Post 7's discipline, and it replays the saved feed payload — no network in unit tests:

test_contract_accepts_good_row ............ ok
test_contract_quarantines_poison .......... ok
test_reconciliation_adds_up ............... ok
test_rerun_is_idempotent .................. ok
test_reemission_caught_by_fingerprint ..... ok

5 passed

And the audit log — one row per run, so "what happened Tuesday?" is a query:

CREATE TABLE audit_log(
  run_id TEXT PRIMARY KEY, window_start TEXT, window_end TEXT,
  fetched INT, valid INT, quarantined INT, duplicates INT, stored INT,
  source_hash TEXT, stored_hash TEXT, hash_match INT, finished_at TEXT
);

Explain it to the customer

"Maria — this morning's pull, in four numbers: 1,005 records in, 990 stored, 5 quarantined, 10 duplicates merged. The 5 quarantined are malformed feed rows — they're in the DLQ table with the exact errors, nothing was dropped silently. The hashes match, so the database holds exactly what the valid feed sent. If anything looks off, rerunning is safe — the dedupe key makes reruns no-ops. And one thing the pipeline taught us: the feed's ZIP codes can't be trusted, so they're stored but never used for identity. If a 'correction' re-emits a record under a new key, we recognize it instead of double-counting it."

The pattern, one last time: the numbers, then the caveats, then the rerunnable method. Maria doesn't read the pipeline — she reads the proof the pipeline produces.

Must know

  • The data contract is a runnable artifact: validate every record with code before it touches the database — never after
  • DLQ: quarantine the raw payload with the error and timestamp; nothing is silently dropped, ever
  • Reconciliation is arithmetic: fetched == valid + quarantined, and the source hash must equal the stored hash
  • Dedupe key as PRIMARY KEY — idempotent reruns by construction, not by careful retry logic
  • Half-open windows (>= start AND < end) for scheduled pulls: no gaps, no overlaps
  • Exclude unreliable fields from identity — a key you can't trust is worse than no key

Useful later

  • Great Expectations / dbt tests — when the team grows past one pipeline and one engineer
  • Change-data-capture vs full pulls — when the feed gets too big to re-pull daily
  • Postgres instead of SQLite — when the readers outgrow a file
  • Alerting on quarantine spikes — the DLQ count is a metric; page on it

Don't memorize this

  • Pydantic validator decorator spelling — remember validate at the boundary, look up the syntax
  • Cron field order — remember schedule it, look up the five fields
  • SQLite DDL details — remember key as PRIMARY KEY, look up the rest

Where this lands in CityOps

The intake pipeline is the second spine of CityOps. Milestone 1 gave the project its skeleton; Milestone 2 gives it a heartbeat — a database that refills itself every morning, with proof. Everything downstream reads from service_requests now: the weekly report, the weather enrichment, and in Stage 3, the API itself. The contract file is the most important artifact you produced this stage: it's the executable version of "what we believe about the feed," and every future change to the feed starts by changing it.

And the stage exit check, stated plainly: you can now take "clean up this data and automate it" from a vague request to a tested, documented, scheduled, reconciled pipeline. That's the whole of Stage 2 in one sentence — and it's the job.

Post 4's principle was if you can't rerun it, you didn't clean it. Post 6's was don't trust the happy path. Milestone 2's: every record accounted for — fetched, stored, quarantined, or duplicate, the arithmetic always closes.

Field check

  1. The morning report says: fetched 1,005, stored 990, quarantined 5. Do the numbers add up? What's missing?
  2. The vendor renames complaint_type to complainttype overnight. Walk through the 6 AM run: what happens, in order?
  3. Why is incident_zip excluded from the fingerprint — and what did that exclusion cost?
  4. The 6 AM run fails halfway; you rerun at 6:20 and it stores 0 new rows. Failure or success? Why?
  5. Maria asks: "Can we just delete the DLQ rows? They're bad data." What's your answer?
What good answers look like

1. Not quite: 990 + 5 = 995, not 1,005. The missing category is duplicates — 10 same-second merges. The full accounting is fetched (1,005) = valid (1,000) + quarantined (5), and valid (1,000) = stored (990) + duplicates (10). A reconciliation report that omits a category is a report that can't be checked — show all four. 2. Extract pulls 1,000 rows; the contract rejects all 1,000 (complaint_type is required, the field is gone); all 1,000 land in the DLQ with "Field required"; 0 rows reach the table; the report reads fetched 1000 / valid 0 / quarantined 1000, adds_up true. The DB is untouched — you fix the contract (or call the vendor), rerun, and the PRIMARY KEY dedupe fills exactly the gap. The alert fires on the quarantine spike, not on a broken dashboard. 3. Dev revealed ZIPs are unreliable — missing or wrong — so a fingerprint built on zip would treat corrections as new records and miss real dupes. The cost: the fingerprint is conservative, so two genuinely identical requests in the same second merge (10 in 1,000 rows here). unique_key is still stored for forensics, so nothing is unrecoverable. 4. Success — that's idempotency working. The PRIMARY KEY dedupe makes the rerun a no-op by construction: duplicates 1,000, stored_this_run 0, hashes still match. The scary outcome would be 1,000 new rows. 5. No — and this is the hill to die on. The DLQ is evidence, not garbage: the reconciliation arithmetic depends on it (delete the rows and fetched no longer equals valid + quarantined), and it's the only record of what the feed sent that you refused. If the rows are truly useless, fix the contract or fix the feed — don't delete the witnesses.