Chapter 14 · ETL vs ELT Architecture: Staging, Raw, Integration, Presentation, and Transform Ownership
Design a Layered Pipeline for the Case Study and Justify Every Persisted Intermediate Dataset
Build and validate AtlasMart landing→raw→staging→integration→presentation locally, justify every persisted dataset, replay from immutable raw bytes, and prove identical output checksums and controls.
Learning outcomes
This lesson turns the chapter contracts into one acceptance run. The fixture includes nine source revision rows: the eight current order lines plus O1002 revision 1 at 200 USD and revision 2 at 190 USD. Integration selects the latest governed revision, so presentation retains Chapter 13's accepted 690-USD state while raw preserves both revisions.
Build the complete AtlasMart landing→raw→staging→integration→presentation pipeline locally from deterministic source revisions.
Justify every persisted dataset and explicitly justify why staging is not persisted in this fixture.
Prove current sales controls remain 8 lines, 5 orders, 10 units, 690 GMV, 425 cost, and 265 gross profit.
Replay from the immutable raw file into an independent database and require equal integration and presentation checksums.
Document cleanup, rollback, lineage, security, and the boundary to Chapter 15 incremental loading.
Chapter 14 preserves the accepted AtlasMart production state established through Chapters 11–13: eight current paid order-line facts, five paid orders, ten units, 690 USD paid GMV, 425 USD cost-at-sale, 265 USD gross profit, and the governed inventory snapshot of 137 units. Chapter 13's Q-prefixed malformed training batch remains isolated quality-test evidence. This chapter changes where transformations execute and which intermediate states are persisted; it does not silently redefine sales grain, correction history, metric formulas, or customer identity.
The mandatory lab is synthetic, local, and free. It uses
Python 3 standard library plus its bundled
sqlite3 module. Generation-time validation ran
with Python 3.13.5 and SQLite 3.46.1; learners should record
their own python --version and
sqlite3.sqlite_version because behavior and
optimizer details can vary. The lab demonstrates layering,
replay, lineage, and deterministic controls—not cloud pricing,
distributed exactly-once guarantees, production durability, or
vendor-specific warehouse performance.
1. Persisted-state justification
| State | Persisted? | Why / acceptance criterion |
|---|---|---|
| Landing CSV | Short-lived | Needed only until byte-preserving raw promotion succeeds and hash is recorded. |
| Raw CSV | Yes | Immutable replay evidence; same batch ID cannot map to different bytes. |
| raw_sales_loaded | Yes inside each run DB | Typed/loaded audit representation with batch/source/load metadata for this local harness. |
| stg_sales | No (TEMP) | Cheap deterministic normalization; no independent consumer or recovery contract. |
| int_sales_current | Yes | Reusable current-revision semantic boundary at order-line grain. |
| mart_sales_daily | Yes | Consumer-facing daily grain with explicit additive controls. |
| pipeline_run + lineage_edge | Yes | Replay checksums and transformation provenance. |
2. Full deterministic lab
Run in an empty working directory. No package installation is
required. The script deletes only its own
atlasmart_ch14_lab directory when resetting.
from __future__ import annotationsimport csv, hashlib, json, os, shutil, sqlite3from pathlib import PathBATCH_ID = "B20260921-001"COLUMNS = [ "order_id","line_no","revision_no","event_ts","recorded_at","customer_id", "product_id","quantity","extended_amount","extended_cost","status"]ROWS = [ ["O1000",1,1,"2026-09-18T15:00:00Z","2026-09-22T08:00:00Z","C001","P200",1,75,45,"paid"], ["O1001",1,1,"2026-09-18T09:15:00Z","2026-09-18T09:16:00Z","C001","P100",2,100,60,"paid"], ["O1001",2,1,"2026-09-18T09:15:00Z","2026-09-18T09:16:00Z","C001","P200",1,25,15,"paid"], ["O1002",1,1,"2026-09-18T11:30:00Z","2026-09-18T11:31:00Z","C002","P300",1,200,125,"paid"], ["O1002",1,2,"2026-09-18T11:30:00Z","2026-09-22T09:30:00Z","C002","P300",1,190,125,"paid"], ["O1003",1,1,"2026-09-19T15:00:00Z","2026-09-19T15:01:00Z","C001","P400",1,100,60,"paid"], ["O1003",2,1,"2026-09-19T15:00:00Z","2026-09-19T15:01:00Z","C001","P200",2,50,30,"paid"], ["O1005",1,1,"2026-09-20T07:10:00Z","2026-09-20T07:11:00Z","C004","P100",1,50,30,"paid"], ["O1005",2,1,"2026-09-20T07:10:00Z","2026-09-20T07:11:00Z","C004","P400",1,100,60,"paid"],]def sha256_bytes(data: bytes) -> str: return hashlib.sha256(data).hexdigest()def canonical_hash(rows) -> str: payload = json.dumps(rows, sort_keys=True, separators=(",", ":"), ensure_ascii=False) return sha256_bytes(payload.encode("utf-8"))def write_landing(path: Path) -> None: path.parent.mkdir(parents=True, exist_ok=True) with path.open("w", encoding="utf-8", newline="") as f: w = csv.writer(f, lineterminator="\n") w.writerow(COLUMNS) w.writerows(ROWS)def promote_raw(landing: Path, raw: Path) -> str: raw.parent.mkdir(parents=True, exist_ok=True) incoming = landing.read_bytes() incoming_hash = sha256_bytes(incoming) if raw.exists(): if sha256_bytes(raw.read_bytes()) != incoming_hash: raise RuntimeError("immutable raw batch already exists with different bytes") else: raw.write_bytes(incoming) return incoming_hashdef load_database(raw: Path, db: Path, batch_id: str, env_name: str) -> dict: if db.exists(): db.unlink() conn = sqlite3.connect(db) conn.executescript(""" PRAGMA foreign_keys = ON; CREATE TABLE raw_sales_loaded ( batch_id TEXT NOT NULL, order_id TEXT NOT NULL, line_no INTEGER NOT NULL, revision_no INTEGER NOT NULL, event_ts TEXT NOT NULL, recorded_at TEXT NOT NULL, customer_id TEXT NOT NULL, product_id TEXT NOT NULL, quantity INTEGER NOT NULL, extended_amount NUMERIC NOT NULL, extended_cost NUMERIC NOT NULL, status TEXT NOT NULL, source_file TEXT NOT NULL, load_ts TEXT NOT NULL, PRIMARY KEY(batch_id,order_id,line_no,revision_no) ); CREATE TABLE int_sales_current ( order_id TEXT NOT NULL, line_no INTEGER NOT NULL, revision_no INTEGER NOT NULL, event_ts TEXT NOT NULL, customer_id TEXT NOT NULL, product_id TEXT NOT NULL, quantity INTEGER NOT NULL, extended_amount NUMERIC NOT NULL, extended_cost NUMERIC NOT NULL, batch_id TEXT NOT NULL, source_file TEXT NOT NULL, PRIMARY KEY(order_id,line_no) ); CREATE TABLE mart_sales_daily ( sales_date TEXT PRIMARY KEY, paid_orders INTEGER NOT NULL, paid_lines INTEGER NOT NULL, units INTEGER NOT NULL, gmv NUMERIC NOT NULL, cost NUMERIC NOT NULL, gross_profit NUMERIC NOT NULL ); CREATE TABLE pipeline_run ( batch_id TEXT PRIMARY KEY, env_name TEXT NOT NULL, raw_sha256 TEXT NOT NULL, integration_sha256 TEXT NOT NULL, presentation_sha256 TEXT NOT NULL ); CREATE TABLE lineage_edge ( upstream TEXT NOT NULL, downstream TEXT NOT NULL, transform_id TEXT NOT NULL, PRIMARY KEY(upstream,downstream,transform_id) ); """) load_ts = "2026-09-21T00:00:00Z" with raw.open("r", encoding="utf-8", newline="") as f: for row in csv.DictReader(f): conn.execute( "INSERT INTO raw_sales_loaded VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (batch_id,row["order_id"],int(row["line_no"]),int(row["revision_no"]), row["event_ts"],row["recorded_at"],row["customer_id"],row["product_id"], int(row["quantity"]),int(row["extended_amount"]),int(row["extended_cost"]), row["status"],raw.name,load_ts) ) # Staging is intentionally TEMP: normalize/validate execution state without persisting another copy. conn.executescript(""" CREATE TEMP TABLE stg_sales AS SELECT batch_id,order_id,line_no,revision_no,event_ts,customer_id,product_id, quantity,extended_amount,extended_cost,lower(trim(status)) AS status,source_file FROM raw_sales_loaded; INSERT INTO int_sales_current SELECT s.order_id,s.line_no,s.revision_no,s.event_ts,s.customer_id,s.product_id, s.quantity,s.extended_amount,s.extended_cost,s.batch_id,s.source_file FROM stg_sales s WHERE s.status='paid' AND NOT EXISTS ( SELECT 1 FROM stg_sales newer WHERE newer.order_id=s.order_id AND newer.line_no=s.line_no AND newer.revision_no>s.revision_no ); INSERT INTO mart_sales_daily SELECT substr(event_ts,1,10) AS sales_date, count(DISTINCT order_id), count(*), sum(quantity), sum(extended_amount), sum(extended_cost), sum(extended_amount-extended_cost) FROM int_sales_current GROUP BY substr(event_ts,1,10) ORDER BY sales_date; """) integration_rows = conn.execute(""" SELECT order_id,line_no,revision_no,event_ts,customer_id,product_id,quantity, extended_amount,extended_cost,batch_id,source_file FROM int_sales_current ORDER BY order_id,line_no """).fetchall() presentation_rows = conn.execute("SELECT * FROM mart_sales_daily ORDER BY sales_date").fetchall() integration_sha = canonical_hash(integration_rows) presentation_sha = canonical_hash(presentation_rows) raw_sha = sha256_bytes(raw.read_bytes()) conn.executemany("INSERT INTO lineage_edge VALUES (?,?,?)", [ (f"landing/{raw.name}",f"raw/{raw.name}","copy-byte-preserving-v1"), (f"raw/{raw.name}","temp.stg_sales","parse-normalize-v1"), ("temp.stg_sales","int_sales_current","latest-revision-paid-v1"), ("int_sales_current","mart_sales_daily","daily-aggregate-v1"), ]) conn.execute("INSERT INTO pipeline_run VALUES (?,?,?,?,?)", (batch_id,env_name,raw_sha,integration_sha,presentation_sha)) controls = conn.execute(""" SELECT count(*),count(DISTINCT order_id),sum(quantity),sum(extended_amount), sum(extended_cost),sum(extended_amount-extended_cost) FROM int_sales_current """).fetchone() assert controls == (8,5,10,690,425,265), controls assert sum(r[4] for r in presentation_rows) == 690 conn.commit() conn.close() return {"raw_sha256":raw_sha,"integration_sha256":integration_sha, "presentation_sha256":presentation_sha,"controls":controls, "presentation_rows":presentation_rows}def main() -> None: root = Path(os.environ.get("ATLASMART_LAB_HOME", "atlasmart_ch14_lab")) if root.exists(): shutil.rmtree(root) landing = root / "landing" / f"erp_sales_{BATCH_ID}.csv" raw = root / "raw" / f"erp_sales_{BATCH_ID}.csv" write_landing(landing) raw_hash = promote_raw(landing, raw) first = load_database(raw, root / "warehouse.db", BATCH_ID, "dev") second = load_database(raw, root / "warehouse_replay.db", BATCH_ID, "replay") assert first["integration_sha256"] == second["integration_sha256"] assert first["presentation_sha256"] == second["presentation_sha256"] assert sha256_bytes(raw.read_bytes()) == raw_hash manifest = { "batch_id": BATCH_ID, "raw_sha256": raw_hash, "integration_sha256": first["integration_sha256"], "presentation_sha256": first["presentation_sha256"], "controls": list(first["controls"]), "layers": { "landing": "receipt buffer; removable after raw promotion", "raw": "immutable byte-preserving source evidence", "staging": "SQLite TEMP table; execution-scoped", "integration": "persisted governed current sales semantics", "presentation": "persisted daily analytical output" } } (root / "lineage_manifest.json").write_text(json.dumps(manifest, indent=2), encoding="utf-8") print("batch:", BATCH_ID) print("raw_sha256:", raw_hash) print("integration_sha256:", first["integration_sha256"]) print("presentation_sha256:", first["presentation_sha256"]) print("controls(lines,orders,units,gmv,cost,gp):", first["controls"]) print("daily:", first["presentation_rows"]) print("replay_equal: True") print("cleanup: remove", root)if __name__ == "__main__": main()
3. Expected output and what it proves
batch: B20260921-001raw_sha256: e59a09c3d89855c024bd6f2cf283c9b1a65aa769f8a0656e57e76e1f88dd716bintegration_sha256: 548c4361d9071b4d5e85e2a507ec4d86b42c6c1e90bb7f2d92dedf5a63826dcepresentation_sha256: c2516669eed256c512ad207494c0cf3086d5d2aa2262e28ae09d4795e3b73690controls(lines,orders,units,gmv,cost,gp): (8, 5, 10, 690, 425, 265)daily: [('2026-09-18', 3, 4, 5, 390, 245, 145), ('2026-09-19', 1, 2, 3, 150, 90, 60), ('2026-09-20', 1, 2, 2, 150, 90, 60)]replay_equal: Truecleanup: remove atlasmart_ch14_lab
The result proves that, for this fixture and canonicalization, two independent databases built from the same immutable raw bytes produce identical integration and presentation state. It also proves the Chapter 13 controls are preserved. It does not prove distributed transactional guarantees, performance superiority, source accuracy, or that these exact persistence choices are optimal at production scale.
4. Verification checklist
-
Run
python --versionand printsqlite3.sqlite_version; record differences from generation-time validation. - Confirm raw CSV contains 9 revision rows while integration contains exactly 8 current rows.
- Confirm O1002 revision 2 contributes 190 USD and revision 1 remains in raw evidence.
-
Confirm integration controls equal
(8,5,10,690,425,265). - Confirm daily GMV sums to 690 and all three presentation rows reconcile to integration.
- Confirm replay integration and presentation SHA-256 values equal the first run.
- Edit the existing raw file under the same batch name and rerun promotion separately: expect a rejection, not silent overwrite.
- Confirm staging is TEMP and not present after the SQLite connection closes.
-
Inspect
pipeline_run,lineage_edge, andlineage_manifest.jsonfor batch/hash/provenance evidence.
5. Controlled failure: “the output matches, so the architecture is correct”
A pipeline can produce 690 once while still being unsafe: raw may be mutable, reruns may duplicate rows, environment paths may be hard-coded, lineage may be absent, and a correction may be unreproducible. The acceptance test therefore includes input identity, independent replay, state ownership, controls, and cleanup—not only one dashboard number.
6. Cleanup/reset and production migration
# Linux/macOS/Git Bashrm -rf atlasmart_ch14_lab# PowerShell# Remove-Item -Recurse -Force .\atlasmart_ch14_lab
In production, replace the local filesystem/SQLite choices only after mapping equivalent guarantees: immutable or versioned raw evidence, transactional/atomic publishing boundaries, batch/run metadata, permissions, lineage, replay capacity, and metric reconciliation. Cloud object stores, warehouses, and orchestrators differ in semantics and cost; re-check current platform documentation before turning this local pattern into product-specific commands.
7. Bridge to Chapter 15
Chapter 14 rebuilt one complete deterministic batch. Chapter 15 introduces the harder operational question: how to ingest only what changed without missing, duplicating, or misordering updates. Watermarks, high-water marks, log-based CDC, overlap windows, idempotent merges, duplicate delivery, late data, deletes, and restart behavior all depend on the raw identity, layer ownership, and replay boundaries established here.
Knowledge check
Check your understanding
- Why does raw contain 9 rows while integration has 8?
- Why is staging TEMP in the lab?
- What exact production controls must remain unchanged?
- What does replay_equal: True mean?
- What new problem begins in Chapter 15?
Review the answers
1. Raw preserves both O1002 revisions; integration selects one current governed revision per order line.
2. It has no independent consumer/recovery requirement and is cheap to deterministically rebuild from raw.
3. 8 current paid lines, 5 orders, 10 units, 690 GMV, 425 cost, and 265 gross profit.
4. The same immutable raw bytes and transform logic produced identical canonical integration and presentation states in two local databases.
5. Incremental/CDC state: deciding what changed and proving restart, duplicate, late, delete, and reprocessing correctness.
Authoritative references
- Python documentation — csvStandard-library CSV reader/writer used for the deterministic landing/raw fixture.
- Python documentation — hashlibSHA-256 fingerprints used as byte and semantic reproducibility evidence; hashes do not establish business correctness by themselves.
- Python documentation — sqlite3DB-API interface used by the local ELT-style execution harness.
- SQLite — CREATE TABLETable and constraint semantics used by persisted raw/integration/presentation state.
- SQLite — TEMP tablesSupports the chapter's execution-scoped staging example without implying that every staging layer must be temporary in production.
- Kimball Group — Dimensional Modeling TechniquesPrimary reference for the dimensional semantics that the presentation layer must preserve regardless of ETL/ELT execution style.
- OpenLineage documentationOptional reference for standardized lineage concepts; OpenLineage is not required by the local lab.
- dbt Developer HubOptional later-course tooling reference for modular SQL/dependency/promotion ideas; dbt is not a prerequisite for this chapter.