Implement state transitions that can be replayed and reconciled.
Implement Dimensional Models, SCDs/Snapshots, Incremental/CDC Loads, Data Quality, Orchestration, and Semantic Metrics
Implement and verify the AtlasMart dimensional slice with SCD history, snapshots, CDC/restart semantics, quality gates, orchestration, and governed metrics.
Learning outcomes
Implement the sales star at the declared atomic grain and preserve conformed keys/history.
Validate SCD2 half-open windows and show the historical C001 segment lookup.
Run quality/quarantine, CDC crash/retry/duplicate/delete, and watermark tests on disposable copies.
Separate orchestration task success from data certification.
Compile the governed revenue metric and catch a deliberate 665 USD semantic drift.
Chapter 30 begins from the accepted AtlasMart state produced by
the earlier chapters:
10 paid order lines, 8 orders, 12 units, 820 USD gross
revenue, 495 USD cost, and 325 USD gross profit. The accepted source waterline remains 208.
The atomic sales grain remains
one paid order line identified by (order_id, line_no); gross_revenue_usd.v1 remains the governed
paid-line revenue metric in USD with UTC time semantics. CDC,
defect, backfill, performance, security, and disaster-recovery
exercises run against disposable copies and must reconcile back
to this baseline.
The capstone evidence was executed with Python 3.13.5 and SQLite 3.46.1 on synthetic data in UTC. SQLite supplies relational constraints, transactions, indexes, query-plan evidence, and file backup/restore. It does not reproduce distributed shuffle, cloud IAM, managed row policies, autoscaling, object-store catalogs, multi-region DR, or provider billing. Those concepts are represented only as explicit policy fixtures, scheduler/cost calculations, or architecture decisions. No production credentials, personal data, paid services, or network dependencies are required.
1. Problem frame: a green DAG can still publish the wrong 665 USD
AtlasMart’s capstone pipeline can complete every technical task
and still be wrong. If a semantic query silently excludes the
sales channel, the executable fixture returns
665 USD instead of governed
820 USD. The pipeline is therefore designed
around state transitions and certification evidence, not task
color.
2. Implement the atomic star with enforceable local constraints
The local SQLite harness uses three core structures:
fact_sales, dim_product, and
dim_customer_scd. The implementation is
deliberately small enough to inspect, yet it preserves
production concerns: composite fact grain, dimension
relationships, accepted values, SCD windows, and additive
measures.
PRAGMA foreign_keys = ON;CREATE TABLE dim_product( product_id TEXT PRIMARY KEY, product_name TEXT NOT NULL, category TEXT NOT NULL);CREATE TABLE fact_sales( order_id TEXT NOT NULL, line_no INTEGER NOT NULL, order_ts TEXT NOT NULL, customer_id TEXT NOT NULL, product_id TEXT NOT NULL REFERENCES dim_product(product_id), channel TEXT NOT NULL CHECK(channel IN ('web','store','sales')), status TEXT NOT NULL CHECK(status='paid'), qty INTEGER NOT NULL CHECK(qty > 0), revenue_usd INTEGER NOT NULL CHECK(revenue_usd >= 0), cost_usd INTEGER NOT NULL CHECK(cost_usd >= 0), PRIMARY KEY(order_id, line_no));
SQLite constraints are executable local evidence, but managed analytical engines vary in which constraints are enforced, informational, or optimizer-only. Production deployment must verify the actual engine.
3. SCD2 history: half-open windows prevent double matches
Customer C001 has three synthetic versions: SMB before
2026-09-18 08:30 UTC, Growth until 2026-09-19 00:00 UTC, then
Mid-Market. The validity rule is
effective_from <= event_time < effective_to,
with NULL end for the current row. Half-open boundaries let
adjacent versions meet without overlap.
| customer_sk | customer_id | segment | effective_from | effective_to | current |
|---|---|---|---|---|---|
| 1 | C001 | SMB | 1900-01-01 00:00Z | 2026-09-18 08:30Z | 0 |
| 2 | C001 | Growth | 2026-09-18 08:30Z | 2026-09-19 00:00Z | 0 |
| 3 | C001 | Mid-Market | 2026-09-19 00:00Z | NULL | 1 |
SELECT segmentFROM dim_customer_scdWHERE customer_id = 'C001' AND effective_from <= '2026-09-18T12:00:00Z' AND (effective_to IS NULL OR '2026-09-18T12:00:00Z' < effective_to);-- Growth
The executed test proves adjacent C001 windows are non-overlapping and exactly one row is current. It does not prove all future source change patterns are valid; production still needs overlap/gap tests for every business key.
4. Quality gate: schema constraints are necessary but not sufficient
The capstone stages 13 rows: the ten canonical facts plus an
exact duplicate, an orphan P404 product, and a
negative quantity. The controlled gate accepts ten and
quarantines three. The canonical fact remains unchanged.
| Defect | Observed reason | Action |
|---|---|---|
| duplicate O1006/1 | duplicate_grain | quarantine; do not double count |
| unknown product P404 | orphan_product | quarantine or resolve an explicit unknown-member policy |
| qty = -1 | nonpositive_qty | quarantine; do not let NOT NULL masquerade as quality |
quality: 13 staged -> 10 accepted / 3 quarantinedaccepted fact controls: (10, 8, 12, 820, 495, 325)production controls changed by quarantine? no
Quarantine preserves evidence and permits bounded repair/reprocessing. Silently dropping rows would make source-to-target reconciliation impossible.
5. CDC is a transactional state machine, not “insert newer rows”
The disposable CDC target begins at waterline 208. Event e209 inserts a new synthetic line; an injected crash occurs before commit, so both the row and watermark roll back. Retry applies once. Duplicate redelivery is ignored. Event e210 deletes the same synthetic row, returning controls to the baseline; its duplicate is ignored.
rolled_backappliedduplicate_ignoredappliedduplicate_ignoredbaseline (10, 8, 12, 820, 495, 325)after e209 insert (11, 9, 13, 870, 525, 345)after e210 delete (10, 8, 12, 820, 495, 325)
def apply_event(conn, event_id, seq, op, payload): conn.execute("BEGIN") watermark = read_watermark(conn) if seen(event_id) or seq <= watermark: conn.rollback() return "duplicate_ignored" mutate_target(conn, op, payload) remember_event(conn, event_id, seq) update_watermark(conn, seq) conn.commit() # target + event identity + watermark become one state transition
The accepted production waterline remains 208 because the capstone CDC exercise uses a disposable copy. That distinction is essential: a recovery test must not mutate course continuity just to demonstrate restart behavior.
6. Snapshots remain separate when measures are semi-additive
Sales events and inventory state must not be forced into one universal fact. The course established periodic inventory snapshots separately because on-hand units are additive across products/locations but generally not additive across snapshot dates. The capstone bus matrix preserves this separation. If an inventory snapshot is backfilled, its exact date/location/product grain and source evidence must be reconciled independently from sales.
7. Orchestration expresses dependencies; certification proves data state
The capstone DAG uses the logical sequence below. A scheduler can implement it in many products; the semantics are vendor-neutral.
source_ready -> ingest_immutable -> contract_and_quality_gate -> resolve_dimensions_history -> merge_atomic_fact_and_watermark -> metric_and_reconciliation_tests -> publish_materializations_marts -> certify_and_notify
A task can return exit code 0 while the source is incomplete or a metric drifts. Certification therefore inspects manifest completeness, atomic controls, history validity, metric golden results, access policy, and freshness. Scheduling alone never establishes correctness.
8. Governed semantic metric: same fact, one versioned meaning
gross_revenue_usd.v1 sums paid line revenue in USD
under the governed contract. The fixture’s deliberately wrong
query excludes channel='sales', reproducing the
earlier semantic incident pattern.
-- governed metric v1SELECT SUM(revenue_usd)FROM fact_salesWHERE status='paid';-- 820-- deliberate dashboard driftSELECT SUM(revenue_usd)FROM fact_salesWHERE status='paid' AND channel <> 'sales';-- 665
The 155 USD difference is not a performance issue or a source delay. It is semantic drift. Fixing it requires restoring the governed metric contract or explicitly creating a new metric/version—not relabeling 665 as “revenue.”
9. Controlled failure: unbounded append on retry
A failed batch is rerun with
INSERT INTO fact_sales SELECT ... and the
operator manually advances the watermark afterward. If the
first attempt partially wrote rows, the retry duplicates
facts; if the watermark advances first, rows can be skipped.
Repair: stable grain keys plus stable event identity, atomic transaction boundaries, bounded source scope, duplicate/delete semantics, and a watermark committed with the target. Then inject crash/retry/duplicate tests before relying on the mechanism.
10. Local lab: execute the capstone state transitions
# Expected deterministic checkpoints from the executed harnessassert controls == (10, 8, 12, 820, 495, 325)assert c001_segment_at_noon_sep18 == "Growth"assert quality == {"staged": 13, "accepted": 10, "quarantined": 3}assert cdc_states == [ "rolled_back", "applied", "duplicate_ignored", "applied", "duplicate_ignored"]assert governed_revenue == 820assert drifted_revenue == 665assert final_disposable_cdc_controls == controls
11. Verification checklist
- Atomic fact uniqueness is enforced at the declared line grain.
- SCD2 windows do not overlap and exactly one current row exists per current member.
- Defects are quarantined with reasons; control totals do not silently change.
- Crash before commit leaves row state and watermark unchanged.
- Retry and duplicate redelivery converge; explicit delete restores the baseline.
- Certification checks data/metric state, not only task completion.
- Metric v1 reproduces 820 USD; the 665 USD negative case is detected.
12. Production judgment and bridge to Lesson 3
The capstone now has reproducible logical state and restart semantics. Only after correctness is stable is it meaningful to choose indexes, materializations, marts, workload isolation, or cloud/lakehouse serving boundaries. Lesson 3 makes those choices from measured workload evidence while preserving the exact results established here.
Knowledge check
Why does the CDC crash test happen before commit?
Show answer
It verifies the transaction boundary: target mutation, seen-event identity, and watermark must either all commit or all roll back.
Why is 665 USD a semantic incident rather than a data-quality row defect?
Show answer
The rows are valid; the consumer applied an ungoverned filter that changed the metric definition.
Why keep inventory in a separate fact?
Show answer
Its periodic state has a different grain and semi-additive time behavior from atomic sales events.
Does a successful orchestration DAG prove the warehouse is correct?
Show answer
No. Readiness, completeness, reconciliation, history, metric, freshness, and policy checks must certify the data state.
Authoritative references
- SQLite — Transactions for the local atomic state-transition model.
- SQLite — Foreign Key Support for local relationship enforcement.
- Kimball Group — DW/BI resources for SCD, snapshot, grain, and fact/dimension techniques.
- Big Data Academy — Data Warehousing and Dimensional Modeling curriculum.