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.

Intermediate → Advanced150–190 minutesPipeline telemetry labDuration + volume + errors + lagLast reviewed: September 2026

Learning outcomes

01

Define pipeline duration, records/bytes, rejects, lag, retries, and resource use precisely enough to be operationally useful.

02

Separate task telemetry from data correctness and consumer availability.

03

Emit stable run IDs and event timestamps that let operators reconstruct a run.

04

Avoid high-cardinality labels and dashboards that monitor only CPU.

05

Connect each metric to a decision, owner, and failure mode.

Continuity: observability watches the governed system; it does not redefine it

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.

Executed local reliability fixture

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

Local metric event
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

Observed fixture metrics
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

  1. Confirm the healthy fixture reports 7.0 minutes and 10/10 rows.
  2. Confirm source readiness and certification timestamps are retained separately.
  3. Confirm rejects and retries are recorded independently from task status.
  4. Confirm no customer/order identifier is used as a metric label.
  5. Confirm pipeline success is not treated as a proof of data freshness or metric correctness.

Knowledge check

Acceptance questions

  1. Why is records_in == records_out not a universal quality test?
  2. Which timestamp best represents consumer availability?
  3. Why can low CPU coincide with a severe data incident?
  4. 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

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.

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.