Instrument AtlasMart warehouse runs with operational and data-facing metrics so duration, volume, lag, retries, rejects, and resource signals can be interpreted in the context of correctness and consumer impact.
Pipeline Metrics: Duration, Records/Bytes, Error/Reject Counts, Lag, Retries, and Resource Use
Build an identity and access model for AtlasMart that separates humans from services, eliminates shared credentials, enforces environment boundaries, and proves least privilege with allow/deny evidence.
Learning outcomes
Define pipeline duration, records/bytes, rejects, lag, retries, and resource use precisely enough to be operationally useful.
Separate task telemetry from data correctness and consumer availability.
Emit stable run IDs and event timestamps that let operators reconstruct a run.
Avoid high-cardinality labels and dashboards that monitor only CPU.
Connect each metric to a decision, owner, and failure mode.
Chapter 25 begins from the accepted AtlasMart state established through Chapters 01–24: 10 current paid lines, 8 orders, 12 units, 820 USD gross revenue, 495 USD cost, and 325 USD gross profit. The fact grain remains one current accepted paid order line, source progress remains committed through sequence 208, Chapter 20 metric contracts remain authoritative, Chapter 21 certified marts remain dependent on conformed assets, Chapter 22 security/privacy controls remain in force, Chapter 23 lineage/ownership metadata supplies blast-radius context, and Chapter 24 tests remain the correctness gates. Observability adds continuous evidence and incident handling; it must not silently reinterpret business rules merely to make a dashboard look healthy.
Runtime: Python 3.13.5 + SQLite 3.46.1. Storage: local in-memory structures plus SQLite semantics where transactions matter. Clock: UTC with explicit timestamps. Security: synthetic identifiers only. Cost: free/local. Healthy run: 10 records in/out, 4,820 bytes, 7.0 minutes, zero rejects/retries, 860 ms fixture CPU, 34 MB fixture peak memory, and certification 12 minutes after source readiness. Important limitation: these CPU/memory numbers are fixture measurements used to teach signal relationships; they are not performance recommendations for any warehouse engine or cloud service.
1. The realistic problem: all tasks are green, but finance sees yesterday’s data
AtlasMart’s daily warehouse workflow reports “success.” CPU is low, SQL returned no errors, and every task exited zero. Finance nevertheless sees a stale certified dataset because the source arrived late. The business decision is whether the 08:30 UTC finance report can be trusted; the source state is a readiness manifest plus paid order-line events; the fact grain is one accepted paid line; the consumer is a certified finance mart. This is why pipeline telemetry must include data-facing time and volume semantics rather than infrastructure health alone.
2. Define the clocks before computing lag
| Timestamp | Meaning | Why it matters |
|---|---|---|
| event_time | When the business event occurred | Required for historical semantics; can be earlier than ingestion |
| source_ready_at | When the producer declares the partition/extract complete | Separates upstream lateness from downstream runtime |
| run_started_at | When orchestration begins after readiness | Scheduling/queue evidence |
| run_finished_at | When processing tasks complete | Task duration, not yet consumer readiness |
| certified_at | When quality/reconciliation/security checks make the product consumable | End-to-end availability to consumers |
Lag is meaningless without endpoints. “Lag = 20 minutes” could mean event-to-ingest, source-ready-to-start, source-ready-to-certified, or now-minus-latest-event. Record which one you are measuring.
3. The minimum useful run record
run = { "run_id": "run-healthy", "source_ready_at": "2026-09-21T08:04:00Z", "started_at": "2026-09-21T08:05:00Z", "finished_at": "2026-09-21T08:12:00Z", "certified_at": "2026-09-21T08:16:00Z", "records_in": 10, "records_out": 10, "bytes_in": 4820, "rejects": 0, "retries": 0, "cpu_ms": 860, "peak_mem_mb": 34, "status": "success"}
In production, metrics normally become time series and logs/traces carry richer context. The mechanism is vendor-neutral: stable run identity, timestamps with clear semantics, counts/bytes, failures/retries, and enough ownership/partition metadata to investigate.
4. Executed healthy-run evidence
duration_min 7.0source_to_certified_min 12.0records_in 10records_out 10bytes_in 4820rejects 0retries 0cpu_ms 860peak_mem_mb 34status success
This proves the local healthy fixture processed the expected small batch efficiently enough for its fixture objective. It does not prove freshness unless source readiness and certification time are compared to an SLO, and it does not prove semantic correctness unless Chapter 24-style tests/reconciliation also pass.
5. Counts are controls only when populations match
records_in == records_out may be a useful
completeness signal when every source row must become one
accepted target row. It is wrong for transformations that
legitimately filter pending orders, explode arrays, aggregate
rows, quarantine defects, or build SCD versions. AtlasMart
therefore records source manifest counts, accepted counts,
reject counts, and target-grain expectations separately. “Same
count” is evidence only after the population and grain are
declared.
6. Retries are not automatically bad—and zero retries is not automatically healthy
A transient I/O failure followed by one safe idempotent retry can be operationally acceptable. Repeated retries can indicate upstream instability, lock contention, rate limiting, or a deterministic bug being retried uselessly. Track retry reason and task identity. Never retry a non-idempotent write simply because an orchestrator can retry it; Chapter 15’s stable identity/transaction rules remain prerequisites.
7. Resource metrics need workload context
| Signal | Useful interpretation | Misleading interpretation |
|---|---|---|
| CPU time | Did compute demand change for comparable data/query shape? | “Low CPU means data is correct” |
| Peak memory | Did a comparable run spill/approach limits? | “More memory always makes the pipeline reliable” |
| Bytes read/written | Did scan/transfer work change materially? | Comparing different projections/populations as if equivalent |
| Duration | Did end-to-end or stage runtime regress? | Comparing cold/warm cache or different concurrency without disclosure |
Performance/cost conclusions require equivalent datasets, query semantics, cache/concurrency context, and engine/version disclosure. This lesson records small fixture numbers only so the incident signatures are observable.
8. Controlled failure: monitor only CPU
Wrong: alert only when CPU exceeds 80%. The late-source scenario has normal CPU because there is nothing to process until 09:05; the dashboard remains stale for almost an hour with no CPU alarm. Repair: instrument source readiness, latest accepted event/partition, run timing, certification time, completeness, rejects, and semantic reconciliation. CPU then explains resource symptoms instead of acting as a proxy for data reliability.
9. Cardinality is part of telemetry design
Labels such as environment, pipeline, task, partition date, and
bounded status are usually manageable. A metric label containing
customer_id, raw order_id, or
arbitrary error text can create unbounded time-series
cardinality and may leak sensitive data. Keep entity-specific
evidence in governed logs/quarantine tables where access and
retention are controlled.
10. Production judgment and bridge
Correctness/non-guarantee: metrics reveal behavior; they do not prove truth by themselves. Freshness/history: timestamp semantics must be explicit. Idempotency: retries require stable batch/event identity. Quality: record rejects and reconciliation outcomes, not only row throughput. Security: telemetry can leak identifiers or SQL/error payloads. Observability: every signal needs an owner and decision. Cost: more telemetry is not free; retain the signals that reduce uncertainty during incidents. Lesson 2 combines these run metrics with data-level freshness, volume, distribution, schema, lineage, and quality evidence.
11. Verification checklist
- Confirm the healthy fixture reports 7.0 minutes and 10/10 rows.
- Confirm source readiness and certification timestamps are retained separately.
- Confirm rejects and retries are recorded independently from task status.
- Confirm no customer/order identifier is used as a metric label.
- Confirm pipeline success is not treated as a proof of data freshness or metric correctness.
Knowledge check
Acceptance questions
-
Why is
records_in == records_outnot a universal quality test? - Which timestamp best represents consumer availability?
- Why can low CPU coincide with a severe data incident?
- What makes a retry safe?
Review the answers
1. Legitimate filtering, aggregation, expansion, quarantine, or history versioning can change row counts.
2. certified_at, because it
represents the point the governed product is consumable
after checks.
3. A late/missing source can leave nothing to compute while downstream data becomes stale.
4. Stable identity, idempotent semantics, and commit/restart state that cannot double-apply effects.
Authoritative references
- OpenTelemetry — Metrics specificationAuthoritative metric concepts and separation of instrumentation from collection/export behavior.
- OpenTelemetry — Metrics data modelUseful for reasoning about time-series measurements, aggregation, temporality, and preserving metric semantics.
- Google SRE Book — Service Level ObjectivesDefines SLIs/SLOs and emphasizes selecting indicators from user-relevant service behavior rather than everything that is easy to measure.
- Google SRE Book — Monitoring Distributed SystemsBackground on actionable monitoring signals and avoiding monitoring that produces noise without operator decisions.
- SQLite — TransactionsSupports the local demonstration of atomic repair/backfill state and safe commit boundaries.
- Python — statisticsUsed for small deterministic baseline examples; production anomaly detection needs workload-specific statistical validation.
- Python — hashlibUsed to fingerprint deterministic incident evidence and repaired outputs.
12. Lab cleanup/reset
The mandatory fixture is local and synthetic. Close the process
or rerun ch25_lab_internal.py from a clean process
to reproduce the metric record; no telemetry backend is
required.