It includes the full SQL curriculum, advanced SQL, data modeling, system design, PySpark, dbt and behavioral preparation, along with the 12-week study system.
There’s also a 30-day guarantee if you work through the playbook and don’t feel measurably sharper at SQL and Data Engineer interviews.
One final thought
The interviewer isn’t going to give you the exact question you practiced last night.
They’re going to change the wording.
Add a constraint.
Introduce duplicates.
Add late-arriving data.
Change the scale.
Ask what happens when the pipeline fails.
Then ask you to explain why your solution works.
So don’t prepare only for questions.
Prepare for patterns.
Because the real interview skill isn’t:
“I’ve seen this exact question before.”
It’s:
“I haven’t seen this exact question before, but I recognize what it’s testing.”
A ten-thousand-row scraping job died two hours in, wallet already charged for the first six thousand rows, when a regional storefront started quietly returning empty product grids instead of a clean 403 I could at least catch. A request from a cloud IP gets fingerprinted as automation by anything that matters — pricing pages, ad verification targets, review aggregators — before retry logic even notices something’s wrong. The fix isn’t a longer timeout. It’s a different exit point.
That’s the problem 2extract is actually solving, and it’s worth being precise about what changes and what doesn’t. It’s a residential and mobile proxy network, billed pay-as-you-go, where every targeting decision — country, city, ZIP, carrier, how sticky a session should be — is a parameter you build into a proxy username rather than a checkbox in a vendor portal. That’s a different operating model from the flat “buy a country, get an IP list” products a lot of data teams have already tried and half-abandoned. The same account is also reachable through an MCP (Model Context Protocol) server at mcp.2extract.com, so an agent — Claude, Cursor, Codex, anything that speaks MCP — can check the balance, look up a geo slug, create a proxy, and cap its spend without anyone hand-writing the client shown below.
I found that out the expensive way, twenty minutes into my first real run. I’d copied a proxy username from the dashboard, then pasted -country-us onto the end of it for a quick test — except the copied string already carried a targeting suffix from an earlier session I’d forgotten to clean up. Two suffixes, one string, an instant 407 that read exactly like a wrong password. It wasn’t. It was a syntax problem wearing an authentication error’s clothes, and it’s the first thing this article will save you from repeating.
TL;DR
2extract is a residential and mobile proxy network, not a scraper — targeting lives entirely in the proxy username, not headers or dashboard toggles. This walkthrough starts with 2extract’s MCP server (mcp.2extract.com), which puts balance, proxy, and geo-targeting controls in front of any MCP-capable agent, then covers the same account driven by a Python package (twextract) that builds three real datasets — geo/address, eCommerce pricing, and market research. Real code, real output, real failure modes from both.
One pricing detail worth knowing
Most proxy vendors sell a country as a separate product. 2extract sells one residential proxy and treats geography as a parameter on the request, so the bill doesn’t multiply by how many markets get targeted. A 407 also narrows to one of three specific causes on this stack (see the FAQ below), not a generic dead end — useful to know before you’re debugging it at 2am against a real budget.
What 2extract actually is — and the one thing it isn’t
2extract does not scrape websites. It doesn’t parse HTML, doesn’t decide which fields belong in a product catalog, and doesn’t store anything you collect. What it does is sit between your program and the public internet and make each outbound request look like it originated from a real residential — or, on the mobile product, a real carrier — connection, optionally narrowed to a country, state, city, or ZIP you choose.
Every stack built on top of it has the same three layers. The collector — an agent calling MCP tools, or Python code with requests and BeautifulSoup — chooses URLs, parses responses, dedupes, and writes files. The gateway, at proxy.2extract.net:5555, authenticates the request, applies targeting, picks an exit IP, and tunnels the traffic through. The target is whatever site or API actually holds the data.
Mobile proxies push the residential idea further: the exit is an actual carrier IP (T-Mobile US shows up under ISP code 310260 in the docs), priced separately, and targeted with -isp instead of -country — the two targeting families are mutually exclusive on a single username, and the gateway will reject a request that tries to mix them.
None of this is a license to do anything you couldn’t already justify without a proxy. The Acceptable Use Policy rules out unauthorized access, account farming, scalping, spam, and ad fraud in plain terms, and a residential exit doesn’t override robots.txt, a site’s terms of service, or copyright law.
The proxy username is the actual control plane
Most proxy vendors let you pick a country from a dashboard dropdown or set it in an HTTP header. 2extract does neither. After creating a proxy, the dashboard hands you a base username shaped like 2xt-customer-[CLIENT_ID]-proxy-[PROXY_NAME], and that base is immutable for the life of the proxy. Everything situational gets appended afterward as a hyphenated suffix, built fresh by your code on every request.
@dataclass(frozen=True)
class Targeting:
country: str | None = None
state: str | None = None
city: str | None = None
zip_code: str | None = None
asn: str | None = None
isp: str | None = None
session: str | None = None
time_minutes: int | None = None
const: bool = False
def validate(self) -> None:
t = self.normalized()
geo = any([t.country, t.state, t.city, t.zip_code])
net = any([t.asn, t.isp])
if geo and net:
raise TargetingError(
"Geographic targeting (country/state/city/zip) and network "
"targeting (asn/isp) are mutually exclusive. Use one group only."
)
if (t.state or t.city or t.zip_code) and not t.country:
raise TargetingError("city, state and zip require country in the same username.")
def build_username(base_username: str, targeting: Targeting | None = None) -> str:
base = base_username.strip()
if targeting is None:
return base
parts = targeting.suffix_parts()
return f"{base}-{'-'.join(parts)}" if parts else base
What the residential product is actually sold for
Use case
Why residential matters here
Session strategy
eCommerce & price monitoring
Stores change currency, tax, and stock by country; a datacenter IP gets a generic or blocked page
Sticky for the length of a category walk
Market research
Discussion APIs and open product databases still rate-limit aggressively at volume
Rotating, one exit per batch
Geo / localization
Regional prices, SERPs, and ads only render correctly to an exit that actually looks local
Rotating, unique IP per market sample
Ad verification & SEO / SERP
Same mechanics underneath, but SERPs are heavily JavaScript-driven
Usually needs a real browser, not this stack
How a request travels, and what happens the moment it doesn’t
Every request makes the same trip — collector to gateway to target, response back the same path. A residential peer can simply go offline mid-session. A proxy left in Inactive status returns a 407 even against a perfectly correct password. Gateway-level errors surface in an X-2extract-Error or X-Proxy-Error response header — except that header goes dark over HTTPS CONNECT tunneling.
def _gateway_407(self, targeting, exc):
header = self._read_gateway_header_if_needed(targeting)
hint = header or (
"the proxy is Inactive, the password is wrong, or the wallet is empty. "
"Open My Proxies and set Status to Active, then rerun."
)
if header and "inactive" in header.lower():
hint = f"{header}. Open My Proxies and set Status to Active (Actions menu)."
return GatewayError(f"2extract rejected the proxy login (407): {hint}", status_code=407)
Setting up an account — the part no library can do for you
Create an account and open the dashboard. Unverified accounts carry a $50 lifetime deposit cap.
Top up the Pay-As-You-Go wallet. The minimum is $10, and residential and mobile traffic bill from the same balance.
If the job will ever need more than $50 deposited, verify identity first.
Create an API key scoped to at least geo:read, proxies:read, proxies:write, and balance:read. It’s shown exactly once.
Create one residential proxy, and a second dedicated proxy only if mobile ISP targeting is actually needed.
Set a monthly spend cap on each proxy.
Copy the base username and password somewhere safe — a .env file if you’re driving this with Python, the agent’s own credential store if it’s MCP — never the targeted version with a suffix already attached.
Running the same account through MCP
Everything above assumes a person writing Python. 2extract also ships an MCP server at mcp.2extract.com, exposing that exact same account — balance, proxies, geo catalog, traffic history — as tools any MCP-capable agent can call directly. Nothing about the underlying product changes: it’s still one residential proxy, still billed pay-as-you-go, still targeted through a username suffix. What changes is who assembles that username — the agent, from a plain-language request, instead of you, from the Targeting dataclass above.
fig. 4 — the round trip behind every tool call: agent → MCP server → account, and back
I connected it as a custom MCP server inside an agent’s tool panel and had 19 tools and 19 resources enabled with nothing more than the server URL and an API key — no separate SDK, no client to install.
📷 the 2extract MCP server connected inside an agent’s tool panel — createProxyResource, deactivateProxyResource, deleteProxyResource, and getAccountBalance visible among the 19 enabled tools
fig. 5 — the six steps below, in one continuous session against the same account
Step 1 — balance, in plain language
The simplest sanity check skips the dashboard tab entirely. Asking the agent “What is my 2extract balance?” triggers one tool call — getAccountBalance — and comes back in about 22 seconds with the wallet figure precise to six decimal places, not just the rounded number the UI shows:
{
"balance_display": "$4.75",
"balance_exact": "$4.746951",
"note": "enough to create a new proxy"
}
📷 getAccountBalance answering a plain-language balance check — the agent surfaces both the rounded and the exact figure
Step 2 — geo lookup, the same targeting table without a docs tab
Fig. 1 above warned that every geographic suffix has to come from 2extract’s own catalog, never guessed — a state name that’s off by one word doesn’t error, it just silently gets ignored as a targeting parameter. searchGeoRegions puts that catalog directly in the agent’s hands. Asking for Germany’s regions returns the exact slugs the proxy username expects:
Region
Username param
Baden-Württemberg
badenwurttemberg
Bavaria
bavaria
North Rhine-Westphalia
northrhinewestphalia
Saxony
saxony
Thuringia
thuringia
📷 searchGeoRegions resolving Germany’s states into the exact -state-… slugs the proxy username expects
Notice what the agent volunteered on its own, unprompted: “do not guess them.” That’s the same warning this article gives around fig. 1 — now enforced by the tool call itself instead of a paragraph the reader has to remember.
Step 3 — creating the proxy the lookup was for
With the slug in hand, the same session calls createProxyResource — the tool sitting right next to getAccountBalance in the panel above — with targeting baked into the proxy itself rather than appended per request:
The response is the same base-username shape described earlier in this article — 2xt-customer-[CLIENT_ID]-proxy-de-baden-agent-demo — except this proxy now carries Baden-Württemberg targeting as its default, so a plain request with no suffix at all still exits from that region.
Step 4 — a spend cap, without leaving the chat
Step 6 of the account-setup checklist above says to cap spend on every proxy from the dashboard. Over MCP that’s a second tool call on the proxy just created, rather than a separate trip to a settings page:
At a $4.75 balance that cap doesn’t do much for a demo account, but it’s the same guardrail every production proxy behind this article’s three pipelines runs under — just set through the agent instead of the dashboard.
I deactivated and deleted this demo proxy right after testing it — deactivateProxyResource and deleteProxyResource sit right next to createProxyResource in the same tools panel from the setup screenshot above — so a five-minute demo doesn’t sit around inflating the account’s proxy count. That’s also why the dashboard in the next step still shows exactly one active proxy, not two.
Step 5 — checked against the dashboard, not instead of it
MCP doesn’t read from a separate ledger. The balance, the proxy that now exists, and its cap all show up identically in 2extract’s own web dashboard, because both surfaces read and write the same account state. Asking the agent for a wider view confirms it: a 30-day traffic-and-spend rollup pulls the same wallet history the dashboard’s usage tab shows, broken down by the day it actually happened:
📷 a 30-day usage rollup pulled through MCP — $0.25 spent, 60.5 MB of traffic, all through one proxy — the same figures the dashboard’s usage tab reports for that account and window
And the dashboard itself, checked separately right after, confirms it — balance, active-proxy count, and 30-day traffic all lining up with what the agent had just reported:
📷 the account’s own dashboard — 60.49 MB over 30 days and one active proxy, matching what the agent had just reported
Balance, active-proxy count, and traffic all line up across both views — $4.75 agent-side versus $4.74 on the dashboard is a rounding difference, not a discrepancy, and one active proxy shows on both because the demo proxy was deleted right after use. Same account, two windows onto it — not a shortcut around the dashboard, a second way to reach it.
Step 6 — one request actually routed through the result
Everything so far provisions a proxy; it doesn’t yet prove a byte moved through it. Closing the loop means one ordinary HTTP request, made the usual way, through the exact username MCP just created:
That IP resolves to a Baden-Württemberg ISP block, confirming that the region MCP looked up, and the proxy MCP created, is the region a request actually exits from — the same exit.* field the Python pipelines below log on every row, for exactly this reason.
Three pipelines, three genuinely different datasets
Everything above ran through an agent calling MCP tools. The rest of this guide is the same account, driven by code instead — three pipelines built on one client, for anyone who wants to own the collection logic rather than hand it to an agent.
Geo and address: proof of exit location
def _rotating_geo_targeting(country: str) -> Targeting:
"""A fresh -session id forces 2extract to assign a new exit IP for this hop."""
return Targeting(country=country, session=uuid.uuid4().hex[:12], time_minutes=1)
def scrape_steam_price(client, *, country, app_id="730", targeting=None):
url = f"https://store.steampowered.com/app/{app_id}/"
response = client.get(url, targeting=targeting or Targeting(country=country),
params={"cc": country}, timeout=45)
soup = BeautifulSoup(response.text, "lxml")
price = soup.select_one("div.game_purchase_price") or soup.select_one("[data-price-final]")
return {"app_id": app_id, "price": price.get_text(strip=True) if price else None,
"http_status": response.status_code}
def parse_price(text):
raw = " ".join(str(text).split())
currency = "GBP" if "£" in raw else "EUR" if "€" in raw else "USD" if "$" in raw else None
match = re.search(r"[\d]+\.[\d]+|[\d]+", raw.replace(",", ""))
value = float(match.group(0)) if match else None
return raw, value, currency
Output lands in data/output/eCommerce.xlsx. Sessions stay sticky for the length of a category walk.
Market research
Pulls from Hacker News through the Algolia API and from Open Food Facts. Its “sentiment” field is an engagement heuristic derived from points and comment counts — not a human-coded label and not an LLM classifier.
From collected rows to a spreadsheet someone will actually open
Every pipeline writes a JSON checkpoint every hundred rows before touching Excel at all — a habit that’s saved a multi-hour run from a crashed openpyxl write more than once. The JSON-to-Excel step deliberately does no new computation; it only formats data that’s already been collected and validated.
How the package is actually organized
fig. 3 — three pipelines share one client; checkpoints guarantee reproducible dataset generation
None of the three pipeline scripts talk to requests directly. They all go through a single ExtractClient that owns authentication, retry behavior, and username construction, and that boundary is the entire reason adding the market research pipeline took an afternoon instead of a week.
Residential, mobile, or datacenter — the real tradeoffs
Datacenter
Residential
Mobile
Exit looks like
Cloud ASN (AWS, GCP)
Home ISP connection
Carrier IP
Block / CAPTCHA rate
Highest
Low
Lowest
Typical cost
Cheapest
Mid
Highest
Targeting granularity
Data center location only
Country / state / city / ZIP
Carrier (ISP code) only
Best fit
Low-stakes internal jobs
eCommerce, geo, market research
Mobile-specific SERPs, app traffic
Where this actually breaks
Country targeting is a request sent to the gateway, not a legal guarantee that a specific city lands on every hop — the only trustworthy record of where a request actually exited is the exit.* data that comes back, and it belongs in every row you store. JSON APIs like DummyJSON and Open Food Facts still return a bare 403 without a properly set User-Agent, proxy or not. Matching “the same product” across two catalogs by name and brand is a heuristic, not an identity. And none of this changes the baseline fact that a proxy alters your odds of getting a response, not your obligation to respect what a site’s terms of service actually say.
How this was tested
Every pipeline described above was run end-to-end against a live 2extract account — real wallet balance, real proxy status states, real 407s along the way — not just read from documentation. The MCP walkthrough was run against that same account in the same session; the balance, geo slugs, and traffic figures shown are what the account actually returned, not stand-ins. The full source is on GitHub under an MIT-style layout anyone can clone and rerun against their own account.
Glossary
Term
Meaning
Exit IP
The public address the target website actually sees
Gateway
proxy.2extract.net:5555
Base username
2xt-customer-…-proxy-name, with no targeting suffix appended
Sticky session
Same exit IP held across many requests via -session (+ optional -time)
407
Gateway rejected authentication or proxy status — rarely the password itself
PAYG
Wallet billed per traffic used, not an unlimited monthly crawl allowance
MCP server
mcp.2extract.com — exposes balance, proxies, and the geo catalog as tools an agent can call
Questions worth asking before rolling this into production
Does 2extract scrape websites for me?
It doesn’t — it’s a proxy network that routes requests through residential or mobile exit IPs. Parsing, deduplication, and storage stay entirely the collector’s job.
Does 2extract have an MCP server for agents like Claude, Cursor, or Codex?
Yes — mcp.2extract.com exposes the same account (balance, proxies, geo catalog, traffic history) as tools any MCP-capable agent can call directly, as shown in the walkthrough above. It’s the same wallet and the same billing as the Python package; MCP just hands the control plane to the agent instead of a human writing the client.
Why does a request get a 407 even when the password is definitely right?
In practice that almost always means the proxy is set to Inactive in the dashboard, the wallet balance has run out, or the username carries malformed or conflicting targeting parameters.
Can country targeting and mobile ISP targeting live on the same proxy?
They can’t. Geographic targeting and network targeting are mutually exclusive on a single username, and mobile ISP targeting needs its own dedicated proxy credential besides.
Is a separate proxy needed for every country being targeted?
No — one residential proxy covers every supported country. Targeting is a parameter appended per request, not a separate product purchased per market.
Does a residential proxy make scraping any site legally safe?
It doesn’t. A proxy changes what IP a target sees; it has no bearing on robots.txt, a site’s terms of service, or copyright law.
Pricing, account limits, and interface details in this article reflect 2extracts dashboard as of September 2026, captured while writing and testing the code above end-to-end against a live account. Treat this as a technical guide to the architecture and integration pattern rather than a live pricing reference vendors update rates and UI labels independently of when an article like this is published, so confirm current numbers on the dashboard before budgeting a production run.
When the consumer at the end of a pipeline is a language model rather than a dashboard, the expensive step moves from the transform to the last hop, and every row you push through it costs money. This article covers the pipeline shape we keep returning to, a four-question test for choosing batch, event-driven or request-path inference, and the three cost levers that matter, in the order they matter.
Photo: Derrick Coetzee, “Front of server racks at NERSC”, Wikimedia Commons, CC0 1.0 public domain dedication
Most teams adding an AI feature to an existing product start by asking which model to use, or whether they need a streaming platform. Both matter far less to the bill than a duller question.
How many model calls does one business event cause, and how many of them could have waited, or not happened at all?
A pipeline built around that question stays predictable. A pipeline that treats the model as one more sink, like a reporting table, tends to produce the invoice that gets the feature switched off.
The Consumer Changed, the Pipeline Did Not
A classic analytics pipeline ends in a dashboard. The transform is the heavy step, the output is an aggregate, and reprocessing a day of data is cheap because warehouse compute is cheap per row. Staleness of a few hours is usually fine.
A pipeline that feeds an AI feature inverts most of that:
The expensive step is the last one. Warehouse SQL over a few hundred thousand rows costs little. Sending those same rows to a model provider is billed per token, per row.
The output is per record, not aggregated. A product description, a ticket category, a summary for one report. There is no “roll it up” shortcut.
Reprocessing is no longer free. A backfill that is harmless for a dashboard can be the most expensive job of the month when a model sits at the end of it.
Idempotency becomes a cost control, not just a correctness property. A retry that re-sends a row is a second charge.
The practical consequence: design the pipeline so the model sees as few rows as possible, each as small as possible, and never the same unchanged row twice.
The Pipeline Shape We Keep Coming Back To
Across the AI retrofits we have shipped into products that were already in production, from ticket triage on a marketplace to drafting product copy and summarising reports for internal reviewers, the architecture has settled into four layers.
The operational store stays the write path, the warehouse prepares a model input contract, and a worker only calls the model for rows whose input actually changed.
Figure 1: The pipeline shape. The operational store stays the write path, the warehouse prepares a model input contract, and a worker only calls the model for rows whose input actually changed.
Our examples use BigQuery and a Node.js worker, because that is where most of this work runs. The shape maps directly onto AWS: object storage plus Athena or Redshift for the warehouse layer, and a scheduled container task for the worker.
Layer 1: Ingestion, Kept Deliberately Boring
The operational database stays the write path. The relevant tables are exported, by a managed streaming export or a scheduled one, into append-only raw tables in the warehouse.
This is the same additive pattern we use for any workload that outgrows the operational database: keep writes where they are and build the new read path elsewhere. The AI feature can then be removed without touching the product’s core data.
Layer 2: Transformation Produces a Model Input Contract
The most useful artifact in the whole pipeline is a single view that defines exactly what the model is allowed to see, in exactly the shape the prompt reads. Nothing else is sent.
-- One row per record the model may see, in the shape the prompt reads.
CREATE OR REPLACE VIEW ai.product_copy_input AS
SELECT
p.product_id,
p.title,
p.category,
p.origin,
p.unit,
-- Hash only the fields the prompt uses. Price and stock are excluded on purpose.
TO_HEX(SHA256(TO_JSON_STRING(STRUCT(p.title, p.category, p.origin, p.unit)))) AS input_hash
FROM raw.products_latest AS p
WHERE p.status = 'unpublished';
Three things happen in that view. Fields are trimmed to the handful the prompt needs. Anything that must not leave your infrastructure is removed or replaced with a placeholder here, in SQL you can review, not in application code scattered across services. And input_hash gives every row a content-based identity the next layer uses to skip work.
On a feature that read health-related customer records, replacing a full record with a hand-picked field set cut prompt tokens by more than half, and reduced what left our infrastructure at all. That one change beat every model-pricing decision on the same feature.
Layer 3: The Inference Worker Is a Pipeline Stage
The model call belongs in a worker that behaves like any other pipeline stage: it selects its inputs, processes them in batches, validates its outputs, and records what it did.
interface ModelInput { productId: string; inputHash: string; prompt: string }
async function runNightlyDrafts(cap: SpendCap): Promise<void> {
// Only rows whose input hash has never been sent to the model
const rows: ModelInput[] = await warehouse.query(`
SELECT i.* FROM ai.product_copy_input AS i
LEFT JOIN ai.processed_inputs AS d USING (product_id, input_hash)
WHERE d.product_id IS NULL`);
for (const batch of chunk(rows, 50)) {
if (!cap.allows(estimateTokens(batch))) break; // stop quietly, the app keeps its fallback
const results = await provider.generate(batch);
const valid = results.filter(isValidDraft); // schema check, invalid output is dropped
await sidecar.upsert(valid); // keyed by product_id and input_hash
await warehouse.insert('ai.processed_inputs', results.map(toLedgerRow)); // failures too, or they are re-sent nightly
cap.record(results);
}
}
The output never overwrites the product’s own fields. It lands in a sidecar table keyed by record ID and input hash, and the application reads it when a valid row exists.
When none exists, because the worker has not run, the cap was hit or validation failed, the application renders what it rendered before the feature existed. Removing the feature means deleting a table, not reversing a migration.
Batch, Event-Driven or Request Path: A Four-Question Test
Most “batch versus streaming” debates about AI features are really a question about who is waiting. We run every feature through four questions, in order, and stop at the first one that gives a clear answer.
Figure 2: The four-question test. Each question is an exit; most features leave at question one or two.
1. Can the Output Be Computed Before Anyone Asks for It?
If yes, it is a scheduled batch job, full stop. Latency is free, batch pricing from providers applies, and the input hash means the nightly run only touches rows that changed.
Product description drafting is our clearest example: retailers submit products during the day, drafts appear overnight in a draft field, and a person approves before anything is published.
2. Does a Person Wait on Screen for the Result?
If nobody is watching, it is event-triggered and asynchronous: a post-write trigger puts the record on a queue, a worker calls the model, and the result lands later.
Support ticket triage works this way. The ticket is saved first, the customer sees no added latency, and the suggested category appears for the ops team a few seconds later. If the call fails or exceeds its timeout, the ticket goes to the default queue as it always did.
For most AI features, this is what “streaming” means: a trigger and a queue. A dedicated event streaming platform earns its place when many independent consumers need the same event history, which one AI feature rarely justifies.
3. Did the Person Explicitly Ask for It?
If a user clicked “tidy up this description”, it is a user-triggered call. Seconds are acceptable because they asked and are watching a loading state. It still needs a cancel path and the original content preserved.
4. Can the Page Render Acceptably If the Call Is Skipped?
Only now do we consider the request path, and only with a hard timeout well under the page’s existing latency budget and a deterministic fallback that renders the pre-AI experience.
If the page cannot render without the model’s answer, the feature is not ready for that path. Precompute it, or do not ship it there.
Where the Bill Actually Comes From
The monthly cost of an AI feature is roughly:
calls per business event × tokens per call × event volume
The price per token is the factor people argue about, and it is the one we touch last.
Across our retrofits, per-call cost varied by two orders of magnitude between features, and the levers that moved it were, in order of impact:
Trim the input. Covered in layer 2. Fewer fields, fewer tokens, less data leaving your systems.
Key the cache on content, not identity. See below.
Switch models, but only with an eval set. A smaller model is cheaper per call, but without a fixed set of real inputs and assertions to prove quality held, a model switch is a guess with a saving attached.
Key the Cache on Content, Not Identity
An expensive mistake we have made ourselves is deciding whether to regenerate based on the record: its ID plus an updated_at timestamp. Records change constantly for reasons the model does not care about.
Figure 3: An illustrative sequence of changes to one product. Keyed on identity, every change triggers a model call. Keyed on a hash of the fields the prompt reads, only the changes the model would notice do.
Keyed on identity, every change triggers a model call. Keyed on a hash of the fields the prompt reads, only the changes the model would notice do.
On a marketplace feature that generates copy once per product version and serves it thousands of times, moving the cache key from the product ID to a hash of the normalised attribute set meant regeneration only happened when attributes actually changed.
Monthly spend on that feature dropped to a fraction of its launch figure, with no change to the model or the prompt.
Put the Spend Cap in the Pipeline
Provider dashboards can alert you, but they cannot make your product degrade gracefully.
We write a hard ceiling into the worker itself, as in the snippet above. When it is reached, the worker stops and the application falls back, so an overspend becomes a quiet degradation someone reviews in the morning rather than an invoice discovered at month end.
Where Managed Services Stop Being Worth It
The pitch for a managed ETL connector, a distributed compute cluster, or a dedicated vector database is usually made as if the data volume were the hard part. For most AI features inside an existing product, it is not.
The volume that matters is bounded by business events: products listed, tickets opened, reports produced.
Our rules of thumb:
Managed connectors earn their fee for SaaS sources you do not control and whose APIs change under you. For your own operational database, a native export into the warehouse is usually simpler and cheaper.
Distributed compute such as Spark earns its place when the transform itself is the heavy step: preparing very large corpora, non-SQL processing at scale, or feeding self-hosted models. When the transform fits in warehouse SQL and the expensive step is a rate-limited API call, a cluster adds operational weight without shortening the part that is slow.
Serverless functions suit short triggers. Long batch runs belong in a container runtime with no execution time ceiling.
A dedicated vector store is worth it once the corpus and query volume outgrow what your existing database or warehouse can serve. For a staff-only internal search over a modest document set, it is often one more system to secure and keep in sync.
The trade-off we accept is a pipeline that looks unimpressive on an architecture slide. In return, fewer systems hold copies of the data, which matters when a deletion request has to reach every copy.
When a Pipeline Is the Wrong Answer
Not every AI feature needs one.
If the feature is user-triggered, operates on content already on screen, and runs a few times a day per user, a direct call with a timeout and a preserved original is simpler and cheaper than any pipeline.
A pipeline also cannot fix missing rules. We once scoped automatic routing of support messages against categories that existed in a dropdown, while the real routing logic lived in one person’s head and contradicted it.
No amount of data engineering helps there. If the rules cannot be written down before work starts, write them down first.
What to Do on Monday
Pick your most expensive AI feature and write down three numbers:
Model calls per business event.
Average input tokens per call.
How many of last month’s calls processed an input identical to one already processed.
Then add an input hash to the view that feeds it.
If that third number is not close to zero, the hash will pay for itself before you touch the model.
Enterprise data platforms often begin with a simple objective: move data from operational systems into a place where it can be analyzed. Over time, however, the number of data sources, consumers, and business requirements grows. A pipeline originally created for one dashboard becomes useful to another team. A transformation developed for a reporting workload is recreated for an application. An AI team builds yet another pipeline because the existing data was not structured for its requirements.
The problem is not that organizations have too few pipelines. In many mature environments, they have too many pipelines performing overlapping work.
This creates a different challenge for data engineering: how do we build data infrastructure that can be reused across analytics, applications, and AI without turning every new requirement into another independent pipeline?
One answer is to move from thinking primarily about ETL pipelines toward thinking about data products.
A data product is not simply a table in a warehouse or a dataset stored in a lake. It is a reusable data asset with defined meaning, ownership, quality expectations, metadata, lineage, and consumers. The objective is to make the data useful beyond the specific pipeline that originally produced it.
The Problem with Use-Case-Specific Pipelines
Traditional ETL architectures are often organized around downstream requirements. A team receives a request for a report, builds an extraction and transformation process, and produces the required dataset. Another team later needs similar information and creates another pipeline because its requirements are slightly different.
At first, this approach is reasonable. The system is small, the requirements are clear, and the fastest solution is often to build exactly what is needed.
The difficulty appears as the organization grows.
Multiple pipelines may independently extract the same source data, apply similar business rules, and create slightly different versions of the same business entity. One pipeline may define an active customer differently from another. One dashboard may calculate revenue using one transformation while another application uses a different version.
Eventually, the organization has a collection of pipelines that individually work but collectively create a difficult data environment.
The goal of a data product approach is not to eliminate pipelines. Pipelines remain essential. The change is in what the pipeline is designed to produce.
Instead of building a pipeline exclusively for one downstream consumer, the pipeline can contribute to a reusable data asset with clearly defined characteristics.
Figure 1. The evolution from use-case-specific ETL pipelines toward reusable data products.
From Pipelines to Data Products
The distinction is subtle but important.
A pipeline describes how data moves and changes.
A data product describes what trusted data is made available for others to use.
For example, an organization may have customer, order, inventory, or product information arriving from multiple operational systems. Instead of creating separate transformations for every consumer, the platform can produce a curated data product representing a well defined business concept.
That product should answer basic questions before another team consumes it:
What does this data represent?
Who owns it?
How frequently is it updated?
What quality expectations does it have?
What transformations have been applied?
Where did the data originate?
Which downstream systems depend on it?
How should consumers interpret important fields?
This turns the dataset from an anonymous technical output into something that other teams can confidently build upon.
The distinction becomes particularly important when AI systems enter the architecture. AI applications need access to enterprise information, but simply exposing more raw data does not necessarily produce better results. The data must have consistent meaning, appropriate granularity, and enough context for the consuming system to use it correctly.
Designing the Architecture for Reuse
A reusable data architecture does not require one enormous centralized pipeline. Instead, it separates concerns while establishing clear interfaces between layers.
A typical architecture can begin with operational databases, APIs, files, event streams, and other enterprise sources. An ingestion layer brings that information into the platform, where raw data can be preserved before transformation.
Transformation and quality processes then produce curated datasets. The important difference is that these curated datasets are designed as reusable products rather than temporary outputs for one report.
Figure 2. A reusable data product architecture separates ingestion, transformation, quality, governance, and consumption.
The architecture can support multiple consumers from the same trusted data product.
Analytics teams may use it for dashboards and reporting. Applications may consume it through APIs or services. Data scientists may use it for machine learning workflows. AI systems may use it as part of retrieval, contextualization, or decision support workflows.
This does not mean every consumer receives exactly the same representation. Different consumers may require different interfaces or derived views. The important principle is that core business logic should not be unnecessarily duplicated.
Data Contracts Make Reuse Possible
Reusability becomes difficult when consumers do not know what they can rely on.
A data contract provides an explicit agreement between data producers and consumers about the expected characteristics of a data asset. At the simplest level, this can include schema and data types. In a mature environment, the contract can go further.
It can define expected semantics, ownership, freshness, acceptable values, compatibility expectations, and changes that require communication.
Consider a field called status.
From a technical perspective, a string is a perfectly valid datatype. But what does the string mean?
Does active mean an account is currently usable? Does it mean the customer has purchased something recently? Does it mean a subscription is paid?
Schema validation cannot answer that question.
For reusable data products, semantic consistency is as important as structural consistency.
A contract therefore becomes a mechanism for protecting consumers from unexpected changes while giving producers a clear responsibility for maintaining the data they publish.
Quality Is Part of the Product
Data quality should not be treated as a final step performed after a pipeline has been built.
If a dataset is intended to become a reusable data product, quality is part of the product itself.
Different products will require different checks, but common considerations include completeness, validity, uniqueness, consistency, and freshness.
For example, a product containing transactional information might need to detect duplicate records. A product supporting operational decisions might require strict freshness expectations. A product used for historical analysis may tolerate delayed updates but require strong consistency over time.
The important point is that quality expectations should be explicit and measurable.
This also changes how data engineers think about failures. Instead of asking only whether a pipeline completed successfully, engineers can ask whether the resulting data product continues to meet its defined expectations.
Metadata and Lineage Are Not Optional Extras
When organizations have hundreds of datasets, discovering what a dataset means can become as difficult as producing it.
Metadata helps answer questions such as where a dataset came from, what its fields represent, how frequently it changes, and who is responsible for it.
Lineage provides another important dimension: understanding how data moved through the system and which upstream sources contributed to the final product.
This becomes especially valuable when something changes.
If a source field is modified, engineers should be able to determine which transformations and downstream consumers may be affected. Without lineage, that investigation can become a manual search across pipelines and documentation.
For data products to remain reusable, discoverability and explainability need to be designed alongside the data itself.
One Data Product, Multiple Consumers
A major advantage of this approach is that the same trusted data foundation can support different types of workloads.
Figure 3. A reusable data product can support analytics, applications, and AI workloads without duplicating core transformation logic.
Consider a curated product representing a business entity such as a product, customer, transaction, or inventory position.
An analytics team might use it to create operational dashboards. An application might use the same information to support a workflow. An AI system might use it to provide context to an agent or model.
The consumers are different, but the underlying business definitions do not need to be reinvented each time.
This is where the concept becomes particularly powerful for enterprise AI.
Making Data Products Useful for AI
AI systems introduce a new category of data consumer.
Traditional analytical workloads often operate through structured queries and predefined metrics. AI applications may need to retrieve information dynamically, combine multiple pieces of context, interpret relationships, and use that information as part of an inference or action.
That places additional demands on the underlying data.
AI systems benefit from data that is:
semantically consistent
sufficiently granular
appropriately contextualized
discoverable
governed
fresh enough for the intended use case
accessible through reliable interfaces
This does not mean every data product needs to be redesigned specifically for AI.
Instead, organizations should build reusable data foundations that can support AI as one of several consumers.
That distinction helps prevent a common architectural mistake: creating an entirely separate data ecosystem every time a new AI initiative appears.
Avoiding the “One Pipeline Per Use Case” Trap
The answer is not to centralize every transformation into one massive pipeline.
Over centralization can create its own problems. A change made for one consumer can unexpectedly affect many others. Teams may also become dependent on a central group for every modification.
A better approach is to identify which data assets and transformations are genuinely reusable.
Common business entities and shared definitions are strong candidates for reusable products. Highly specialized analytical logic may remain closer to the consuming workload.
The architectural question should therefore be:
What should be shared, and what should remain specific to the consumer?
Good data engineering is not about maximizing reuse at any cost. It is about finding the right boundaries.
Practical Principles for Building Reusable Data Infrastructure
Organizations beginning this transition can start with a few practical principles.
Build once, consume many times. When multiple teams repeatedly implement the same business logic, investigate whether the underlying data should become a reusable product.
Define ownership early. A reusable dataset without clear ownership eventually becomes nobody’s responsibility.
Treat metadata as part of the product. Documentation, definitions, lineage, and discoverability are not administrative additions. They determine whether another team can actually use the data.
Make quality measurable. Define expectations around freshness, completeness, validity, and other characteristics that matter to the product’s consumers.
Design for change. Schemas, business rules, and upstream systems will evolve. Data products should have clear compatibility and change management practices.
Separate shared data from consumer specific logic. Not every transformation needs to be centralized. Reuse the parts that represent stable, broadly useful business concepts while allowing downstream teams to build specialized views.
Conclusion
The evolution from ETL pipelines to data products is not about replacing one technology with another. It is a shift in how organizations think about the outputs of data engineering.
A pipeline can successfully move data from one system to another and still create little long term value if every downstream consumer must interpret, validate, and transform that data independently.
A data product takes a different approach. It treats trusted data as a reusable enterprise capability with defined meaning, quality expectations, ownership, metadata, and lineage.
That approach becomes increasingly important as organizations add AI systems to their technology landscape. AI does not eliminate the need for sound data infrastructure. It increases the number of ways that trusted enterprise data can be consumed.
The mature data platform, therefore, is not simply a collection of pipelines.
It is an ecosystem of reliable data products that allows analytics, applications, and AI systems to build on the same trusted foundation.
As a data engineer, you’ve likely been in a design review where someone says at the very end, “We should just do ELT for everything!”. You are a data engineer, and you’ve probably been in a design review where you heard someone say at the end, “We should just do ELT for everything!”. Or you’ve inherited a package in an old 10-year-old version of SSIS that you would have to pick apart row-by-row in a painful manner, and the business is wondering why the “modern” warehouse team can’t simply replace it overnight. The industry has morphed ETL, ELT and now Reverse ETL into tribal identities: pick one, fight about it in slack threads, ship. But that’s backwards. The pattern is not the personality, it is a means to an engineering end.
The reality, though, is that ETL, ELT, and Reverse ETL are not mutually exclusive approaches — they are three different kinds of pipelines that are tackling three different problems and most production data pipelines require all three to be running concurrently. Don’t identify a winner. It’s to align the pattern with the limitations of the workload: data sensitivity, complexity of transformation, latency requirements, and the most cost-efficient place to execute the workload.
The Problem Deep Dive
All three patterns use the same three verbs—extract, load, transform—but in different orders: that’s why it’s a bit of a muddle.
ETL (Extract, Transform, Load): ETL is the process of moving data into a staging area, where it is transformed in some way outside of the target system before it is loaded into the clean data in the target system. This was the standard practice for decades as target systems or data warehouses, particularly transactional databases, would be costly to compute on and required to be buffered from raw and dirty data. This pattern has become the backbone of careers such as those of tools like SSIS, Informatica, and Talend.
ELT (Extract, Load, Transform) first loads raw data into the destination, and then transforms it-while-stored using compute functions on the destination. It really only became dominant when cloud warehouses (Snowflake, BigQuery, Redshift, Azure Synapse) separated storage from compute and started to offer cheap ways to store raw data and transform it later, incrementally, using tools like dbt.
Reverse ETL moves the already transformed, already modeled data out of the warehouse and into operational systems for business teams to take action: Salesforce, HubSpot, Braze (an ad platform). It’s more of a third point on the same spectrum; it’s the vehicle for the other two, which is a problem that neither ETL nor ELT were intended to solve—getting warehouse truth into the tools where other things and humans are happening.
The sticking point is that teams find themselves using these as “short cuts” rather than as decisions based on the workload:
· Failure of compliance due to ELT by-default. A healthcare team loads raw PII into a warehouse without masking or tokenization as that is the way modern data stacks work.A healthcare or fintech team places the raw PII in a warehouse without masking or tokenization first. Now, without the mask, SSNs or PHI are stored in raw schemas that are accessible to half the analytics org and the cost of the compliance retrofit outweighs the transform step. This is the same for ETL’s pre-load transformation; scrub before, don’t scrub after it lands.
· Runaway warehouse costs due to misusing ELT. A team pushes a full nightly extract of a 500-million-row transactional table into Snowflake, and uses dbt models to re-scan the data on every run, rather than incremental models. The warehouse bill expands due to the compute moving from a dedicated ETL server to the meterized cloud credits, without any one looking at the meter.
· Operational data that has been stale due to the lack of Reverse ETL. A customer lifetime value model is created in the marketing team’s warehouse, but not piped back to the CRM. The warehouse model has no way to get back out of Salesforce, so the Sales reps still see the raw purchase counts. The insight is there, but it isn’t there where the person needs it at decision time.
· ETL (Row-by-row) instead of ELT (Set-based). Some classic SSIS or legacy on-prem pipelines where they were changing records one by one inside the pipeline, and the same transformation can be written as a simple set-based SQL statement in the target warehouse and run in a fraction of the time.
All these are the correct tools for the wrong jobs — but not bad tools!
The Solution: A Decision Framework
Ask 4 questions about each workload, instead of “which pattern do we standardize on?”. In reality, all three patterns are implemented together on various pipelines on most platforms.
1. Is the data required to be scrubbed, masked, or filtered before it reaches any place it can be queried? If yes — PII, PHI, cardholder data, anything under GDPR/HIPAA/PCI scope — transform before load. This remains ETL’s main reason to be, regardless of the ELT fashions. Mask/tokenize in the extraction layer (Azure Data Factory data flows, or a simple Python/SQL Server SSIS step), meaning that any raw sensitive values never reach the raw schema of the warehouse. For those on the Microsoft stack, this is the typical best case scenario for maintaining ADF or SSIS in the mix in an otherwise ELT-focused Azure Synapse or Fabric pipeline, instead of removing it from the mix because dbt is cool.
2. Does the transformation require a lot of steps and iterations and needs to be versioned, tested, re-run by analysts? If yes, use ELT. Bring in raw (or lightly scrubbed) data to the warehouse and leave the work of in-place modeling to dbt, stored procedures or Synapse/Databricks notebooks. You have version-controlled transformation logic, built-in automated testing, lineage graphs, and can rebuild history without re-extracting from a fragile source system when someone discovers a bug in the transformation.
The most important decision when creating a model that causes cost blowouts is whether to run it incrementally (only the new or changed rows since the last run) or not. The key for teams transitioning from SQL Server/SSIS to a dbt-style ELT mindset is to get this concept in their heads early: a transformation is no longer “a step in a pipeline,” it’s a “materialized view” that needs to be refreshed, and it’s that refresh strategy where most of the warehouse bill lives or dies.
3. What is the latency requirement and is it variable with each hop? Batch analytics (nightly board reporting) does not require ELT’s load then transform lag. Extracting and initial transforming typically occur within a streaming layer (Kafka/Event Hubs + stream processing) before anything enters a warehouse, or in the context of ETL, transform early.
4. Is there a need for this insight to act upon outside of the warehouse for a human or downstream system(s)? Whether the data pipeline was ETL or ELT, without Reverse ETL the value of the data is lost once it’s returned. Tools such as Census, Hightouch, or a scheduled Azure Function or Logic App that read from a warehouse view and write to an API fill in the gap. Model it once, sync it wherever it needs to act it out. What this typically involves in practice is creating a single, well-governed model in the warehouse, e.g. a “customer health” or “customer segment” view, which categorizes each customer as “at-risk”, “high-value” or “standard” according to recency and lifetime spend, and then letting a Reverse ETL tool match that segment field with a custom field in the CRM, following a set schedule. Never again will any engineer have to create a CSV to Salesforce, and a sales rep’s view of the segment is always up to date with the warehouse model that powers it.
In reality: streaming or batch extraction with PII scrubbed at the source (ETL) → raw but safe data deposited in the warehouse (ELT) → curated marts synced back to operational tools (Reverse ETL). Three patterns, one pipeline, for each one of them it is really good at doing.
Proof: What This Looks Like in Practice
The initial design on a retail insight app I was working on was ELT only – raw order and customer data was coming straight from Salesforce and the transactional database into Snowflake and all masking and transformation was being done downstream in dbt. Until a compliance audit brought up the fact that raw customer PII was accessible to any analytics use case with access to the warehouse, which also accessed some fields that never saw use in any analytics use case.
But it wasn’t about giving up on ELT. It was putting in a thin ETL step at extraction: an Azure Data Factory data flow that was stripping out PII fields before putting data into the raw schema, and everything else flowed directly through to dbt for modeling. Access to detokenized values was restricted to a few service accounts.
The measurable outcome: compliance gap was closed without any action on the 40+ dbt models that were already deployed, as the business logic that needed to be changed didn’t need to be moved. Compute cost was not impacted — the masking step did not cost anything in the warehouse compute, it did on the extraction layer. The team then integrated an hour-by-hour (via Hightouch) Reverse ETL sync to move the customer health segment view from the marketing platform into the CRM – removing a weekly manual CSV export performed by a marketing analyst every Monday. None of these three required giving up the other changes, and each pattern was used precisely where the compromises were warranted.
The Close
ETL, ELT, and Reverse ETL are not competing architectures for your loyalty, but three tools to address three different questions: what must be cleaned before it is put into the warehouse, what is less expensive to transform when compute resides there, and what must come out of the warehouse to have an impact? The ones who are burned are those who choose one pattern and apply all of their workloads through it.
The next time you are thinking of setting up a pipeline, don’t ask “are we an ETL shop or an ELT shop?”. For each workload: Does it need scrubbing before it lands? Is the transform complex enough to benefit from version control and incremental materialization? Does it need to be materialized at extraction time (latency) given that output needs to walk back out of the warehouse to do any good? When answering the four questions truthfully for each data flow, the correct pattern — typically multiple patterns — emerges spontaneously.
Every data engineering leader has sat through the same meeting: the platform work is done, the pipelines are stable, but defending its value to the board is difficult because it rarely appears on a single P&L line. Data platform budgets are often treated as discretionary IT spend, meaning they get cut first during a downturn. Getting ROI measurement right is what keeps your next initiative fundable.
Why Data Engineering ROI Is Hard to Measure (and Why It Matters)
Boards fund outcomes. Data engineering delivers infrastructure. That mismatch is the root of the problem.
A board can evaluate a sales tool by pipeline generated, or a marketing spend by cost-per-acquisition. Data platform investments don’t map that cleanly pipeline reliability, schema governance, and data quality improvements are foundational, meaning they enable other initiatives rather than generating value on their own. An Al model that improves fraud detection accuracy gets the credit; the governed, clean, well-lineaged data pipeline underneath it, without which the model wouldn’t have worked, gets none.
This isn’t just a communication problem. Left unaddressed, it becomes a funding problem. Data platform budgets get treated as discretionary IT spend, get cut first in a downturn, and then get blamed when the next Al initiative underperforms because the data underneath it was never solid. Getting ROI measurement right isn’t an exercise in optics, it’s what keeps the next initiative fundable.
Metrics That Actually Translate to Business Impact
The fix starts with picking metrics a board member without a data engineering background can actually interpret. A few that consistently translate well:
Data downtime cost avoided: Every hour a critical pipeline is down or serving bad data has a real cost delayed reporting, blocked decisions, or in regulated industries, compliance exposure. Tracking incidents avoided (or their reduced frequency after a platform investment) turns an abstract reliability improvement into a dollar figure.
Time-to-insight reduction: How long does it take from “we need this data” to “here’s the answer”? If that cycle shrinks from days to hours after a platform investment, that’s a directly measurable efficiency gain that maps to faster business decisions.
Engineering hours reclaimed from firefighting: A mature platform investment shows up as a shift in how engineers spend their time less time patching broken pipelines and chasing data quality issues, more time building new capabilities. That ratio, tracked before and after, is one of the cleanest ROI signals available.
Data quality incident rate: Fewer downstream errors caused by bad data, wrong numbers in a report, a broken dashboard, a flawed model input is a leading indicator of platform health that’s easy to track and easy to explain.
Cost-per-query or compute efficiency: For teams on modern cloud data stacks, tracking compute spend against query volume or data processed shows whether platform investments are actually improving unit economics, not just adding capability.
None of these require exotic instrumentation. Most are extractable from existing observability and cost-monitoring tools already in place. The work is in deciding which ones matter for a given business and tracking them consistently.
Connecting Data Initiatives to Business Outcomes
Metrics alone don’t make the case they need to be tied to a specific business decision or outcome of the platform investment enabled or unblocked.
The strongest version of this argument doesn’t say “we modernized our data stack.” It says: “faster, more reliable data pipelines cut our fraud review time from four hours to forty minutes,” or “consolidating our data sources let underwriting make decisions same-day instead of next-day.” Specific, traceable, and tied to something the board already understands the value of.
This only works if a baseline exists before the investment. Teams that skip measuring the “before” state lose the ability to prove improvement later, a gap worth closing at the start of any platform initiative, not after the fact when the board asks for numbers. A structured data-readiness assessment before a major platform investment is one of the more reliable ways to establish that baseline, since it forces a documented starting point across data quality, infrastructure, and governance maturity that the post-investment numbers can be measured against.
Framing matters too. An investment task built around “we need to modernize our data infrastructure” competes with every other infrastructure request in the budget cycle. An investment task built around “this unblocks same-day underwriting decisions” competes on the same terms as revenue-generating initiatives and tends to win more often.
Making the Case to the Board
When it’s time to present, resist the instinct to show everything. A board conversation isn’t the place for a full metrics dashboard, it’s the place for three or four numbers, chosen because they answer the two questions every board member is actually asking: why now, and what happens if we don’t.
“Why now” is answered by connecting the investment to a business pressure the board already recognizes regulatory deadlines, a competitor’s faster decision cycles, or a growth plan that the current data infrastructure can’t support. “What happens if we don’t” is answered by quantifying the cost of inaction: the downtime already being absorbed, the compliance exposure already being carried out, the engineering hours already being spent on maintenance instead of building.
This is a distinction we see play out constantly at Samta.ai, working with BFSI and regulated clients across Singapore. The teams that get board sign-off aren’t necessarily running the most technically impressive platforms, they’re the ones who walked into the room with a baseline, a business outcome, and a dollar figure attached to inaction.
A recent IDC-backed business value study on enterprise data platform investments found that organizations with mature data discovery and governance infrastructure consistently recovered platform costs through reduced analyst search time and fewer duplicate data efforts alone before counting any downstream Al or analytics gains. That’s the kind of framing that resonates with a board: cost recovery that doesn’t depend on a speculative future win.
Processing millions of food records taught me that data quality is rarely one clever cleaning function. It is a chain of small contracts: what a number means, which identifier owns a record, what a missing value means, and what an incremental update is allowed to change.
I learned this while building DietlyAPI, a nutrition API backed by more than 4.2 million indexed food records. Much of the catalog originates from Open Food Facts, a valuable worldwide crowdsourced database. That scale and openness are useful, but they also expose every awkward case a data pipeline eventually encounters: mixed units, incomplete labels, placeholder barcodes, duplicate products, implausible nutrition values, and partial updates.
This article explains the patterns that made the pipeline safer. The examples are simplified, but the failure modes are real.
Validation happens at more than one boundary in this pipeline — ingestion protects structure, serving enforces trust, converted here from the original mermaid flowchart.
TL;DR
→ Treating data import as one boolean is_valid check throws away useful information; structural validity, plausibility, and publish-readiness are separate questions that deserve separate gates.
→ Store nutrient values in one consistent unit contract (per 100g) and keep serving size as separate metadata, so no consumer has to guess what a number means.
→ Null, zero, and “omitted from this update” are three different states — collapsing them into one makes incomplete records look complete and produces false-confidence calculations downstream.
→ Relational checks (does sugar exceed total carbs? does the stated calorie count match the macro math?) catch bad records that individually pass simple range checks.
→ Partial updates are more dangerous than full imports: a naive upsert that blindly assigns every incoming field can silently null out good data the source simply didn’t include that day.
→ At 4.2 million records, rare edge cases stop being rare — the highest-value tests target invariants like “an omitted delta nutrient preserves the stored value,” not just a successful job exit code.
The Pipeline Is a Series of Trust Boundaries
My first mistake was thinking about the import as one operation:
Read a source file, clean each row, and insert it.
In production, there are several separate decisions:
Can the source record be parsed?
Can its fields be mapped to a stable internal schema?
Is the record identifiable across future imports?
Are its values plausible enough to store?
Is it complete and trustworthy enough to rank highly or publish?
Can a later partial update safely modify it?
Treating those questions as one boolean is_valid check loses useful information. A record with a name and barcode but no calories may still be worth retaining. A record with impossible calories should not appear in a “popular foods” response. A partial daily update should not delete fields simply because the source omitted them.
The important detail is that validation happens at more than one boundary. Ingestion protects the database from malformed structure. Serving and publishing apply stricter quality rules appropriate to their users.
Lesson 1: Define the Unit Contract Before Writing Transformations
Nutrition sources frequently mix:
values per 100 grams;
values per serving;
grams, milligrams, and micrograms;
kilocalories and kilojoules;
numbers and human-readable strings such as 1 cup (240 g).
If these representations leak into the application layer, every consumer must guess what each number means. That guarantees inconsistent calculations.
Dietly’s internal contract stores nutrient values per 100 grams. Serving information is separate metadata:
calories_kcal = energy per 100 g
protein_g = protein per 100 g
sodium_mg = sodium per 100 g
serving_size_g = optional weight of one stated serving
serving_desc = original display text, such as "1 cup (240 g)"
That separation matters. A serving_size_g of 30 does not mean the stored calories are already scaled to 30 grams. Consumers can calculate a serving explicitly:
Unit conversion should occur once, close to ingestion:
def to_milligrams(value, unit):
if value is None:
return None
normalized = (unit or "g").lower()
if normalized == "mg":
return value
if normalized in {"µg", "mcg", "ug"}:
return value / 1000
if normalized in {"g", ""}:
return value * 1000
# Unknown is not the same as zero.
return None
The final line is deliberately conservative. Silently guessing an unfamiliar unit creates a valid-looking wrong value, which is harder to detect than a null.
The same rule applies to failed parsing:
def parse_number(raw):
if raw is None or not str(raw).strip():
return None
try:
return float(raw)
except ValueError:
return None
In nutrition data, zero is a claim. Null means “not known.” Converting missing values to zero makes incomplete products appear complete and can produce dangerously confident downstream calculations.
Lesson 2: Validate Relationships, Not Only Individual Columns
A schema can confirm that calories are numeric, but it cannot tell you whether 6,000 kcal per 100 grams is credible. Simple range checks catch many broken rows:
Field
Plausible range
calories_kcal
0 – 900
protein_g
0 – 100
fat_g
0 – 100
carbs_g
0 – 100
fiber_g
0 – 100
sugar_g
0 – 100
serving_size_g
0 – 2000
RANGES = {
"calories_kcal": (0, 900),
"protein_g": (0, 100),
"fat_g": (0, 100),
"carbs_g": (0, 100),
"fiber_g": (0, 100),
"sugar_g": (0, 100),
"serving_size_g": (0, 2000),
}
def outside_range(record):
failures = []
for field, (low, high) in RANGES.items():
value = record.get(field)
if value is not None and not low <= value <= high:
failures.append(f"{field}:outside_range")
return failures
However, many bad records contain values that are individually believable but mutually inconsistent. Relational checks are more powerful:
def plausibility_failures(food):
failures = []
if (
food.get("sugar_g") is not None
and food.get("carbs_g") is not None
and food["sugar_g"] > food["carbs_g"] + 0.5
):
failures.append("sugar_exceeds_carbohydrate")
if (
food.get("saturated_fat_g") is not None
and food.get("fat_g") is not None
and food["saturated_fat_g"] > food["fat_g"] + 0.5
):
failures.append("saturated_fat_exceeds_total_fat")
macros = ("protein_g", "carbs_g", "fat_g")
if food.get("calories_kcal") and all(food.get(x) is not None for x in macros):
estimated = (
food["protein_g"] * 4
+ food["carbs_g"] * 4
+ food["fat_g"] * 9
)
stated = food["calories_kcal"]
if abs(estimated - stated) > max(120, stated * 0.5):
failures.append("energy_macro_mismatch")
return failures
These tolerances are intentionally broad. Food labels round values, fiber and alcohol complicate energy calculations, and source conventions vary. The purpose is not to “correct” every label mathematically. It is to catch extreme contradictions before they are promoted, summarized, or used to generate authoritative-looking content.
That led to another useful distinction:
Hard structural checks decide whether a record can enter storage.
Quality gates decide whether it can appear in high-trust surfaces.
Ranking signals decide which acceptable record should appear first.
A sparse record may remain searchable without being selected for a featured-food endpoint. Keeping these policies separate avoids throwing away potentially useful data.
Lesson 3: Identity Is Not the Same as a Barcode
It is tempting to use a barcode as the universal product key. In real data, that fails for several reasons:
some records have no barcode;
scanner noise and hand-entered placeholders exist;
the same source record may be updated while keeping its source identifier;
different sources can use different identifiers for the same food;
similar products are not necessarily the same product.
Dietly uses source provenance as the idempotency key:
CREATE UNIQUE INDEX idx_foods_source_id
ON foods (source, source_id);
This answers a narrow but essential question: “Have I already imported this exact source record?” It does not pretend to solve global entity resolution.
Known placeholder barcode patterns are removed rather than used for lookups. Returning no barcode match is safer than returning a confidently wrong product.
Product deduplication then becomes a separate serving-layer concern. Search candidates can be grouped using a normalized name key — case-folding, punctuation removal, and collapsing repeated words — and the best row can be selected using signals such as:
presence of an image;
serving information;
complete core macros;
realistic ranges;
number of populated nutrient fields;
source confidence.
This approach does not claim that all duplicates disappear. Instead, it prevents weak duplicates from dominating common queries while preserving the original rows and their provenance.
Lesson 4: Partial Updates Are More Dangerous Than Full Imports
The most instructive failure appeared in the incremental pipeline. During one update, eight already-published food pages lost their calorie values and dropped out of the page build. The import had completed successfully; the data had still become worse.
A full export contains a broad set of fields. A daily delta may contain only the fields that are currently present upstream. If an upsert blindly assigns every incoming field, an omitted value becomes SQL NULL and can erase good data already stored.
The unsafe version looks reasonable:
ON CONFLICT (source, source_id) DO UPDATE SET
calories_kcal = EXCLUDED.calories_kcal,
protein_g = EXCLUDED.protein_g,
fat_g = EXCLUDED.fat_g;
But it treats “not included in this update” as “delete the existing value.”
The safer policy for Dietly’s source is to preserve stored nutrition when a delta omits it:
ON CONFLICT (source, source_id) DO UPDATE SET
name = EXCLUDED.name,
calories_kcal = COALESCE(EXCLUDED.calories_kcal, foods.calories_kcal),
protein_g = COALESCE(EXCLUDED.protein_g, foods.protein_g),
fat_g = COALESCE(EXCLUDED.fat_g, foods.fat_g),
carbs_g = COALESCE(EXCLUDED.carbs_g, foods.carbs_g),
image_url = COALESCE(EXCLUDED.image_url, foods.image_url),
updated_at = NOW()
WHERE foods.name IS DISTINCT FROM EXCLUDED.name
OR foods.calories_kcal IS DISTINCT FROM
COALESCE(EXCLUDED.calories_kcal, foods.calories_kcal)
OR foods.protein_g IS DISTINCT FROM
COALESCE(EXCLUDED.protein_g, foods.protein_g)
OR foods.fat_g IS DISTINCT FROM
COALESCE(EXCLUDED.fat_g, foods.fat_g)
OR foods.carbs_g IS DISTINCT FROM
COALESCE(EXCLUDED.carbs_g, foods.carbs_g);
There are two protections here.
First, COALESCE encodes the meaning of a missing delta field. This policy is source-specific: if an upstream system supports explicit deletion, it should send a deletion marker rather than relying on null.
Second, IS DISTINCT FROM avoids rewriting unchanged rows. At millions of records, unnecessary updates create write-ahead-log traffic, dead tuples, index churn, and disk pressure. Idempotency is an operational feature, not only a correctness property.
The delta cursor is committed after each successfully processed file. If a job stops halfway through a series, it resumes from the last committed file instead of replaying the entire history or skipping uncommitted work.
Lesson 5: Preserve Provenance All the Way to the API
Once several sources share one table, it becomes easy to flatten away where a value came from. That makes later debugging and trust decisions much harder.
Each Dietly row retains fields such as:
source
source_id
confidence
created_at
updated_at
Provenance supports practical questions:
Which source produced this suspicious value?
Can the record be re-imported deterministically?
Should one source rank above another for this query?
Which rows were affected by yesterday’s delta?
What attribution or license applies downstream?
Confidence is best treated as a ranking input, not proof that a value is correct. A high-confidence source can still contain an error, while an incomplete crowdsourced record can still be useful.
Open Food Facts data is available under the Open Database License, so attribution and downstream license obligations also need to survive the journey from source to product.
Lesson 6: Test the Failure Policy, Not Just the Happy Path
Row counts and successful job exits are weak evidence of pipeline health. A pipeline can finish successfully after replacing thousands of values with null.
The highest-value tests in this system target invariants:
importing the same record twice does not create a duplicate;
an omitted delta nutrient preserves the stored value;
a changed nutrient updates the stored value;
an unchanged record is not rewritten;
placeholder barcodes cannot produce a false lookup;
values outside realistic ranges cannot enter high-trust responses;
public response fields remain backward-compatible.
For SQL generation, even a focused regression test can prevent a repeat:
def test_partial_updates_preserve_nutrition():
sql = UPSERT_SQL.upper()
for column in ("CALORIES_KCAL", "PROTEIN_G", "FAT_G", "CARBS_G"):
assert f"COALESCE(EXCLUDED.{column}" in sql
In addition, record rejection or suppression reasons as categories rather than a single invalid count:
Their trends reveal upstream schema changes faster than inspecting random rows. A sudden jump in invalid_number, for example, may indicate a delimiter or unit change rather than a genuine decline in data quality.
A Practical Checklist
Before calling a large ingestion pipeline reliable, I now ask:
Does every numeric field have a documented unit and reference basis?
Are null, zero, deletion, and omission distinct states?
Is the idempotency key tied to source identity?
Are structural validation, quality gating, and ranking separate?
Do checks cover relationships between fields?
Can partial updates erase existing values?
Do unchanged upserts avoid physical rewrites?
Is source provenance retained in storage and responses?
Can interrupted incremental jobs resume safely?
Do tests reproduce the pipeline’s previous failures?
At 4.2 million records, rare edge cases stop being rare. A one-in-a-million parsing issue is no longer hypothetical, and a harmless-looking upsert can become millions of unnecessary writes.
The central lesson was simple: reliable data quality does not mean making every source row perfect. It means making uncertainty explicit, containing bad values, preserving what is already known, and ensuring that retries produce the same result.
That is less glamorous than the word “pipeline” sometimes suggests. It is also what makes the pipeline dependable.
Someone on the backend team renamed order_total to order_amount. Clean name. Makes total sense for their domain model. They shipped it on a Thursday afternoon. By Friday morning, your revenue dashboard was showing zero. Not wrong numbers. Zero. Because your Snowflake pipeline was still selecting order_total from the events table, and the column simply wasn’t there anymore.
You found out from a Slack message. From a director. At 9 AM.
This is the most common production incident in data engineering in 2026, and it’s almost never caused by bad code. It’s caused by the absence of a formal agreement between the team producing data and the team consuming it. That agreement has a name: a data contract. And most data teams still don’t have one.
The excuse is usually some version of “we move too fast.” The reality is that the teams who move fastest are the ones with contracts, because they stop discovering breaking changes from directors on Friday mornings and start catching them in CI on Thursday afternoons, before anything ships.
TL;DR
→ A data contract is a formal specification — schema, semantics, SLAs, ownership — between a data producer and its consumers. Not documentation. Enforcement.
→ Most data incidents don’t start with missing data or broken code. They start with a well-intentioned upstream change that silently invalidated an assumption someone downstream was relying on.
→ Contracts have three parts: schema (structure and types), semantics (what fields actually mean), and SLAs (freshness, completeness, availability). Schema-only contracts miss most real breakages.
→ The dual-write pattern is the only safe migration path for breaking changes: keep old field + add new field → both populated during transition → deprecation notice with a hard date → removal at v2. Each phase takes at minimum 30 days. Skipping phases causes incidents.
→ 90 days minimum notice for breaking changes. Data pipelines have long release cycles; consumers need time to update downstream logic, tests, and dashboards.
→ A contract not enforced in CI is just documentation. The ODCS (Open Data Contract Standard) YAML spec plus `datacontract-cli` gives you executable, version-controlled contracts in about 30 minutes per dataset.
→ dbt integration: map contract checks to dbt tests. Require a version bump plus consumer sign-off on breaking changes before merge. After one month of this, most teams report significantly fewer schema surprises.
→ The worst gotcha: contracts that only cover schema, not semantics. A field that changes meaning without changing type is undetectable to automated checks — and it’s how revenue figures silently drift for weeks.
Why schemas break and who owns the blame
Schema evolution sits between two teams that don’t talk to each other on the same cadence. The producer team — usually a backend or platform engineering team — is shipping product features, often weekly, and treats every field they emit as their own. The consumer team — your data engineering team — is running pipelines that depend on those fields staying stable, and finds out about breaking changes the same way archaeologists find ruins: by digging through wreckage.
The producer isn’t wrong for evolving their schema. The consumer isn’t wrong for depending on it. The incident happens because there was no shared definition of what “a safe change” means, no process for communicating it, and no tooling to enforce the agreement. The blame falls on the process, not the person. Which means the fix is a process change, not a person change.
Schema evolution is the load-bearing problem in data engineering in 2026, and it’s the problem most teams handle the worst. The good teams treat upstream schemas as contracts and run checks against those contracts on every pipeline run. The teams that lose stakeholder trust treat upstream schemas as suggestions and find out about every breaking change from a Slack message that starts “hey, the dashboard looks weird.”
That Slack message is always sent on a Friday. It is always sent to a director.
What a data contract actually contains
The mistake most teams make when they start with data contracts is writing schema-only contracts. Field names, data types, nullability. It feels rigorous. It catches a specific class of errors — column removed, type changed — but misses most real incidents.
Real breakages happen at the semantics layer. The producer changes order_total from gross to net revenue. Same field name. Same FLOAT type. No schema violation. But your revenue dashboard is now off by 23%, silently, because the number means something different than it did last week. A schema validator cannot catch this. Only a semantic contract can — one that documents what a field means, how it should be used, and what constitutes a valid business interpretation of its values.
A complete data contract has three layers. Schema: field names, data types, nullability, constraints (no negative values in a price field, for example). Semantics: what each field means in business terms, how it maps to domain concepts, what transformations are applied before it reaches the consumer. SLAs: freshness guarantees (this dataset is refreshed within 15 minutes of source update), completeness thresholds (at least 99.5% of expected rows must be present), availability targets, and a named owner with actual contact information — not “data team.”
The Open Data Contract Standard and the YAML spec
The good news for teams starting in 2026 is that there’s a growing standard: ODCS (Open Data Contract Standard), a YAML-based specification that defines schema, quality rules, SLAs, and ownership in a single document. It’s human-readable, version-controllable in git, and machine-parseable by tools like `datacontract-cli`, which can validate contracts, run compatibility checks, and generate reports.
A minimal ODCS contract for an orders dataset looks like:
dataContractSpecification: 0.9.3
id: orders-v1
info:
title: Orders
version: 1.0.0
owner: [email protected]
servers:
production:
type: snowflake
database: PROD_DB
schema: PUBLIC
table: orders
models:
orders:
fields:
order_id:
type: string
required: true
description: Unique identifier for the order
order_amount:
type: number
required: true
description: Net revenue after discounts and returns, in USD
minimum: 0
created_at:
type: timestamp
required: true
servicelevels:
freshness:
description: Data refreshed within 15 minutes of source update
threshold: PT15M
completeness:
description: At least 99.5% of expected rows present
threshold: "99.5%"
This is not documentation theater. This YAML file is executable. `datacontract-cli test` validates your actual Snowflake table against this contract. It checks types, required fields, minimum values, and can be wired into CI so that any schema change that would violate the contract fails the PR before it merges.
The only safe migration path for breaking changes
When a producer needs to make a breaking change — remove a field, rename it, change its type, change its semantics — the contract provides a coordination mechanism. There’s a specific pattern that works, and teams that skip steps in it pay for it.
Day 0: Announce. The producer creates a deprecation notice in the contract YAML, updates the changelog, and notifies consumers via a designated channel. Critically, this notification includes a hard date for removal — not “eventually” or “when everyone has migrated.” Deprecated without a date is just a polite rumor. A field can sit in limbo for eighteen months while producers assume nobody uses it and consumers assume it will live forever.
Days 0–60: Dual-write. The producer populates both the old field and the new field simultaneously. Consumers can migrate on their own schedule during this window. The producer monitors usage of the old field (this is easy with Snowflake’s QUERY_HISTORY and column-level access tracking) to know when all consumers have switched.
Day 60: Deprecation notice with hard date. Consumers who haven’t migrated get a 30-day final warning. This is the reminder that actually motivates stragglers. The hard date is non-negotiable.
Day 90+: Removal at v2. The old field is gone. The contract version bumps to 2.0.0. This is a semantic major version — it breaks backward compatibility — and that bump is what triggers automated alerts to any consumer still on v1.
No drama. No guessing. No 2 AM rollback. Give consumers at least 90 days notice for breaking changes. This seems long, but data pipelines have long release cycles, and consumers need time to update downstream logic, tests, and dashboards.
Making it executable: CI enforcement that actually works
The critical architectural decision with data contracts is this: a contract not enforced in CI is just documentation, and documentation drifts. Within six months, the contract YAML and the actual schema diverge, nobody updates the contract when they ship features, and you’re back to tribal knowledge with extra steps.
The enforcement pattern that works:
1. Compatibility check on PR. Before any schema change merges, run `datacontract-cli diff` against the current production contract. Breaking changes fail the PR automatically. Non-breaking changes (adding a nullable field, loosening a constraint) pass. The definition of “breaking” is explicit in the contract spec, not up to whoever reviews the PR.
2. Consumer sign-off for breaking changes. If a breaking change is intentional (the producer knows and has planned for it), the PR requires explicit approval from all registered consumers of that dataset. This is enforced via GitHub CODEOWNERS or equivalent. Producers can’t ship breaking changes unilaterally.
3. dbt test integration. Map contract quality rules to dbt tests. Freshness SLAs become `dbt source freshness` checks. Completeness thresholds become row count assertions. Not-null requirements become `not_null` tests. These run on every dbt build, so violations are caught before models complete — not after reports are wrong.
4. Runtime validation at ingestion. Before data loads into your Silver or Gold layers, validate incoming records against the contract. Rows that violate constraints get quarantined in a dead-letter queue, not silently loaded as nulls. This catches semantic drift that schema validation misses: an order_amount field that’s suddenly returning negative values because someone upstream changed the sign convention.
The gotchas that sink most implementations
Exposing raw transactional schemas as data products. This is the most common structural mistake. When your data contract directly mirrors your application’s OLTP schema, every application refactor becomes a consumer’s problem. The fix is a stable abstraction layer — expose only what consumers need, not the underlying operational detail. Schema changes to the application layer should be absorbed by your ingestion layer, not propagated downstream.
Brittle contracts that break more than they prevent. Strict attribute lengths, tightly constrained enums, or hyper-specific format requirements feel like good quality controls. In practice, they make schemas so rigid that producers constantly need change approvals for minor operational updates that have no downstream impact. Design contracts around semantic guarantees and business invariants, not implementation details. amount > 0 is a semantic guarantee. DECIMAL(18,4) is an implementation detail that will change.
Unclear ownership is the silent killer. Data contracts fail most often not because of tooling gaps, but because accountability is unclear. When something breaks, teams scramble to diagnose issues that fall between ownership boundaries. Every contract needs a named owner with actual incident-response obligations. Not a team. Not a Slack channel. A person whose name is in the contract and who gets paged when a contract violation is detected at runtime.
Semantic changes that look like no-ops. Changing what a field means without changing its name, type, or schema is the hardest class of breakage to catch. order_amount switching from gross to net. A user_id changing from internal to external identifiers. These require semantic versioning (a major version bump) and human review, not just automated compatibility checks. Your CI can catch structural breakage; only your team can catch semantic breakage.
Contracts that cover batch but ignore streaming. If you have a Kafka-based event pipeline feeding your Snowflake tables, the schema contract lives in the Kafka topic, not in the table. Changes to the Kafka Avro schema — registered in Confluent Schema Registry or AWS Glue — need the same versioning and deprecation discipline as your warehouse schemas. Most teams only contract the warehouse side and get burned by streaming schema changes that propagate silently into their pipeline.
The real cost math
Data engineering incidents from schema breakage are expensive in ways that don’t show up on warehouse bills. A typical schema incident at a mid-sized company looks like: 3–4 hours of two engineers debugging, 1 hour of a data analyst investigating wrong numbers, a director review, and a post-mortem. Call that 10 person-hours, at a blended rate of $150/hour. That’s $1,500 per incident.
Teams that experience two schema incidents a month — which is conservative for a team without contracts — are burning $3,000/month, or $36,000/year, on incidents alone. That doesn’t count the cost of wrong decisions made from bad data before the incident was even discovered. One revenue calculation running off a silent semantic change for three weeks is often worth more than a year of incident cost.
The tooling investment for data contracts — `datacontract-cli`, ODCS YAML per dataset, CI integration — is a few days of engineering time. The 90-day discipline is a process change, not a tooling cost. The math is not close.
Where to start (not where everyone starts)
Everyone says “start with your most critical datasets.” That’s correct but useless. More specifically: identify the three datasets that caused production incidents in the last 90 days. Start with those. Not your biggest datasets. Not your most complex. The ones that already broke something.
For each: write the ODCS YAML (schema + semantics + SLAs + owner). Add `datacontract-cli` compatibility checks to the PR workflow for that dataset. Map the quality rules to dbt tests. That’s the first sprint. After one month of this on three datasets, you’ll have a template, a workflow, and enough muscle memory to expand to the rest of the catalog without it feeling like a governance initiative nobody asked for.
The one principle
Change is inevitable. Unmanaged change is expensive. A data contract is the agreement that makes change boring instead of dangerous. The goal isn’t to prevent schemas from evolving — schemas should evolve as the business evolves. The goal is to make every evolution visible, deliberate, and announced far enough in advance that nobody finds out about it from a director on a Friday morning.