Roadmap To Be A Data Engineer / Lesson 03

Lesson 03 Batch vs streaming Fundamentals §3 Ingestion & Transformation About 5 min read

Ghost Listings

How fresh data really needs to be, and what each minute of freshness costs.

1The situation

Tuesday, 09:14. A pre-owned Céline shoulder bag sells for €680. It is the only one — every item in the catalogue has quantity 1. The order is written to the shop database in about 40 ms and the item is marked sold.

The Google Shopping feed for that item was generated at 03:00 that morning by a nightly job that exports the full catalogue. It says available, €680. It will keep saying that until the job runs again — at 03:00 tomorrow.

For the next 17 h 46 min the company pays to advertise a bag it does not have. Three people click the ad. One writes to support. No alert fires; the pipeline is green all day.

Now multiply: ~2,000 items sell per day, each advertised until the next nightly refresh — on average for half a day.

2The mechanism

Average staleness ≈ half the refresh interval + processing time. Worst case is the full interval.

“Yesterday’s data” sounds like a small uniform inaccuracy. Twelve hours of average error is the true sentence — and it is why the fix for a stale dashboard is almost never “schedule the job earlier”.

Latency is a budget, spent in five places:

source commit        40 ms   the sale is written to the shop DB
   ↓ capture         ?       when do we notice? poll schedule or log stream
   ↓ transport       ?       moving the row to the warehouse
   ↓ transform       ?       rebuilding the model on top of it
   ↓ serve           ?       cache / feed / dashboard refresh
total = end-to-end latency

A 60-second stream feeding a dashboard cached for 30 minutes is a 30-minute pipeline. The budget is the sum, dominated by the slowest hop — the most common self-inflicted wound in data engineering.

3What each interval costs (real output of the exercise below)

Refresh Avg stale Ghost h/day Wasted ads Pipeline Total/day
nightly batch 720.1 min 24,002 €144.00 €0.90 €144.90
4× per day 180.4 min 6,014 €36.10 €3.60 €39.70
hourly 30.5 min 1,016 €6.10 €21.60 €27.70 ← cheapest
every 15 min 8.0 min 267 €1.60 €86.40 €88.00
streaming (CDC) 1.0 min 33 €0.20 €45.00 €45.20

“Ghost hours” = item-hours of the catalogue being wrong. Unit costs illustrative: €0.006 per ghost item-hour, €0.90 per full feed rebuild, €45/day for an always-on CDC connector.

Cost falls, then rises, then falls again — a valley, not a slope. And read the bottom two rows together: every 15 minutes costs €88/day and is eight times staler than streaming, which costs €45. Cranking the cron frequency up made things worse on both axes at once.

The reason is structural: the 15-minute job rebuilds the entire catalogue every run to find the ~20 items that changed; streaming reads the DB change log and moves only those 20 rows.

If you want a big freshness jump, change the architecture, not the schedule. Full-refresh batch is cheap when rare and absurd when frequent. Incremental/streaming is expensive when rare and cheap when constant.

4Three architectures compared

Batch Micro-batch Streaming
Cadence hours to a day 1–15 minutes continuous, event by event
Unit of work “all rows since X” a small slice of rows one record
Typical tools Airflow + SQL, dbt, Fivetran/Airbyte Spark Structured Streaming, dbt incremental, scheduled CDC Kafka/Kinesis/Pub-Sub + Flink, Debezium CDC
Re-run a bad day trivial — re-run the job usually fine hard — replay from the log, if you kept it
Failure feels like red DAG at 04:00 a lagging queue a pager at 02:00, lag graph climbing
Real cost driver compute per run compute per run × many runs always-on infra + the people operating it
Use when decision is daily/weekly (reporting, finance, ML training) ops dashboards, near-real-time inventory a machine acts in seconds: fraud, pricing, availability

What actually gets harder in streaming

  • Event time vs processing time — sale at 09:14, arrived 09:17. Which one does “sales per hour” use?
  • Late / out-of-order events — an offline app flushes yesterday’s events today, into a window already reported.
  • Windowing — “last hour” stops being a WHERE clause and becomes state the system must hold and release.
  • Delivery semantics — at-least-once means a record can be processed twice: fine for a log, fatal for a revenue counter. (This is why idempotency is Lesson 04.)

The honest counter-question to “can we make it real-time?” is: which decision changes if the data is 60 seconds old instead of 12 hours old? A real answer (an ad stops, a price moves, a payment is blocked) pays for streaming immediately. “It would feel better” means hourly is right.

5Three questions to ask in a standup

  1. “What is our end-to-end latency — event to someone seeing it?” Not the job’s schedule: the sum of all five hops. Teams quote the 15-minute sync and forget the 30-minute cache behind it.
  2. “What does one hour of staleness cost us here?” No answer → the freshness debate is aesthetic and hourly wins by default. An answer → you have a number to weigh the pipeline bill against.
  3. “Is this job incremental, or does it rebuild everything each run?” The best predictor of what happens to the bill if you double the frequency.

620-minute hands-on (DuckDB)

pip install duckdb. Every statement below was run as written; the output is the real one.

Show the full code (32 lines)
-- STEP 1 — one day of sales. Every item is unique: stock = 1.
CREATE OR REPLACE TABLE sales AS
SELECT i AS item_id,
       TIMESTAMP '2026-08-30 00:00:00' + INTERVAL (((i * 137) % 1440)) MINUTE AS sold_at
FROM range(1, 2001) t(i);

-- STEP 2 — the freshness law as a function: how long until the next refresh notices?
CREATE OR REPLACE MACRO staleness_minutes(sold_at, every_min) AS
  (CAST(every_min AS BIGINT)
   - (date_diff('minute', TIMESTAMP '2026-08-30 00:00:00', sold_at) % CAST(every_min AS BIGINT)));

-- STEP 3 — the options, with their cost model.
CREATE OR REPLACE TABLE options AS
  SELECT a AS every_min, b AS label, c AS run_cost, d AS flat_cost
  FROM (VALUES
    (1440,'nightly batch',   0.90, 0.0),
    (360, '4x per day',      0.90, 0.0),
    (60,  'hourly',          0.90, 0.0),
    (15,  'every 15 min',    0.90, 0.0),
    (1,   'streaming (CDC)', 0.00, 45.0)) v(a,b,c,d);

-- STEP 4 — the whole lesson in one query.
SELECT label, every_min,
  round(avg(staleness_minutes(sold_at, every_min)),1)                     AS avg_stale_min,
  round(sum(staleness_minutes(sold_at, every_min))/60.0, 0)               AS ghost_hours,
  round(sum(staleness_minutes(sold_at, every_min))/60.0 * 0.006, 1)       AS wasted_ad_eur,
  round(any_value(flat_cost) + any_value(run_cost)*(1440.0/every_min), 1) AS pipeline_eur,
  round(sum(staleness_minutes(sold_at, every_min))/60.0 * 0.006
        + any_value(flat_cost) + any_value(run_cost)*(1440.0/every_min), 1) AS total_eur
FROM sales, options
GROUP BY label, every_min
ORDER BY every_min DESC;
│ nightly batch   │ 1440 │ 720.1 │ 24002.0 │ 144.0 │  0.9 │ 144.9 │
│ 4x per day      │  360 │ 180.4 │  6014.0 │  36.1 │  3.6 │  39.7 │
│ hourly          │   60 │  30.5 │  1016.0 │   6.1 │ 21.6 │  27.7 │
│ every 15 min    │   15 │   8.0 │   267.0 │   1.6 │ 86.4 │  88.0 │
│ streaming (CDC) │    1 │   1.0 │    33.0 │   0.2 │ 45.0 │  45.2 │

Then change one number. Set run_cost to 0.05 — what happens when the job becomes incremental instead of rebuilding the whole feed — and re-run. Every-15-minutes drops from €88.00 to €6.40 and becomes the cheapest row.

That is the payload: the sweet spot is not a property of your business, it is a property of your pipeline’s design.

7Takeaway

“Real-time” is not a speed. It is a budget with two sides: what staleness costs you, and what freshness costs to run. Find the valley between them — then make the pipeline incremental and look again.

Vocabulary

batchmicro-batchstreamingend-to-end latencyfreshness SLAdata stalenessCDC (change data capture)full refresh vs incrementalevent time vs processing timelate-arriving datawindowingat-least-once / exactly-onceconsumer lagreplay

Next: Lesson 04 — idempotency & backfills.

Back to top