Data Engineering For Ml

Batch vs Streaming Data for ML: When Fresh Beats Cheap

0 of 24 complete

0%

Contents

Back|Data Engineering For MlBatch vs Streaming Data for ML: When Fresh Beats Cheap
1/24
87 min left
  1. Home
  2. AI Engineering: Data, RAG and Agents
  3. Data Engineering for ML
  4. Batch vs Streaming Data for ML: When Fresh Beats Cheap
Related Topics
Feature Stores: Killing the Train and Serve Skew BugCore ConceptsStale Pieces in a Pipeline: A Cached Input, an Old Scaler, and What Each One HidesThe ML & AI LifecycleFine-Tuning vs RAG vs Prompting: Choosing Your ApproachLLM and GenAI OpsParameter-Efficient Fine-Tuning: LoRA and QLoRALLM and GenAI OpsEvaluating LLMs in Production: Grading Answers That Have No Right AnswerLLM and GenAI Ops
Next lesson
1 of 24
Data Ingestion Pipelines: Getting Data Into the ML Platform

System Design

  • Foundation
  • Intermediate
  • Advanced
  • Capstone

AI Engineering

  • Foundation
  • Data, RAG and Agents
  • Evaluation, LLM Ops and Security

systemdesign.academy

  • Home
  • Glossary
  • Interview prep
  • Reviews
  • About
  • Privacy
  • Terms

Where this shows up in interviews

This idea carries a full system design question on its own. Each walks through the full answer.

  • →Design a Fraud Detection System

The Paper and the Radio

Imagine two ways to find out about the roads before you drive to work. The first is the morning paper. It is printed once, in the night, from everything the reporters learned during the day before. It is cheap, it is complete, and it is right about yesterday. The second is the traffic report on the radio. Someone reads it out every few minutes, from calls that come in while you listen. It needs people on duty all day and all night, it is never finished, and it can tell you that a bridge closed ten minutes ago.

For some questions the paper is all you need: which roads are usually busy on a Monday, or how long some roadworks will last. For others only the radio will do: is the bridge open right now? Paying a radio station to tell you about next month's roadworks is waste. Reading last night's paper to find out whether the bridge is open right now is no use at all.

An illustration of a man at a desk with two screens of charts, beside text. Headed the paper and the radio, titled a count from last night, or a count from this hour? Beside him: a program guesses how many bikes will be rented this hour, and the number it leans on most is the rides in the last hour. Beneath: counted fresh every hour, its guesses were off by 30.5 rides on average; counted once, the night before, off by 56.4. A slower number, the average over the last week, was only 2.7 rides an hour off from fresh across the next day. Last: the question is not batch or stream for the whole program, but how fresh each number must be.

A program that learns from examples, which the rest of this course calls a model, has exactly this choice for every number it reads. Take a bank's check on card payments. Here is a made-up example: a card is used in Lagos, and ninety seconds later the same card is used in Toronto.

A program can stop the second payment if it looks at one number: how many times was this card used in the last five minutes? But that number has to be current, this minute. Say it was worked out by a job that ran at 2 AM. Then it still says whatever it said at 2 AM. The payment goes through, and the money is gone. Nothing in the program was wrong. The number it read was old.

Now take a program that guesses which customers will cancel next month, from what they did over the last 90 days. Does it matter whether its numbers are two seconds old or eight hours old? Not at all. Working them out once a night, in one big run, is far cheaper than updating them every time a customer does anything.

Same company, opposite needs. This lesson measures that choice on real data, and the answer turns out to depend less on the program than on each number it reads.

Where This Lesson Starts

This is the first lesson of the chapter on data engineering for ML. Data engineering is the work of getting data from where it is made to the model, in the right shape and at the right time. The first question that work has to answer is when each number is computed.

Two earlier lessons touch this from other sides. In the chapter on why production breaks, train/serve skew sent a model inputs computed a slightly different way at serving time. It used the same hourly bike rentals I use here. In the lifecycle chapter, stale pieces in a pipeline cached one input once a day and measured what the old value cost. This lesson asks the design question that comes before both: which inputs need to be computed fresh at all, and which can be computed once a night?

A flowchart headed the same rides, two ways to count them, titled where a number comes from decides how old it is. A box, rides happen, all day and all night, leads down two paths. On the left, nightly job at 00:00: reads the whole day, then stops, leading to a cylinder, a table: the same two counts all next day. On the right, stream: updates the two counts every hour, leading to a cylinder, the fresh counts. Both lead to the model guesses this hour's rides. Beneath: in the lab, one model reads the table, one reads the fresh counts, one reads a mix, and one is trained on the fresh counts but served the table.

The lesson has two halves. The first is a small lab on real data. One model guesses bike rentals. I serve its two inputs from a nightly job, from a stream, from a mix of the two, and from a job run every few hours. The second half covers the ideas every streaming system is built on, and I explain each one when we reach it. One of them is what a crash does to a count. Where I could, I measured those on the same data. Where I could not, I say so.

Nine Words for This Lesson

A hand-drawn list headed nine words for this lesson, titled how data arrives, and how old it is. Bounded: data with a first row and a last row, yesterday's rides. Unbounded: data that never ends, the rides still happening now. Batch job: a program that runs on a schedule, reads a bounded chunk, writes a result and stops. Stream: a program that never stops and handles each new record as it arrives. Input: a number the model reads to make a guess, like the rides in the last hour. Fresh: computed from everything up to now. Stale: computed at an older moment, and not updated since. MAE: mean absolute error, how many rides a guess is off by, on average. Skew: the model learned from one kind of number and is served another. Beneath: a stale input raises no error; the model guesses either way.

Data is bounded when it has a first row and a last row, like all of yesterday's rides. It is unbounded when it never ends, like the rides still happening now. A batch job is a program that runs on a schedule, reads a bounded chunk, writes a result and stops. A stream, or stream job, is a program that never stops and handles each new record, one row of data, as it arrives.

An input is a number the model reads to make a guess. Many teams call an input a feature; the two words mean the same thing here. An input is fresh when it was computed from everything up to now, and stale when it was computed at an older moment and not updated since. A stale input is not an error. It is a real number of the right type, in the right place. It is just old.

To say how good the guesses are, I use the mean absolute error, MAE. For every hour, take how far the guess was from the real number of rides, then average that over all the hours. An MAE of 30 means the guesses were off by 30 rides an hour on average. Lower is better.

Last, skew, short for train/serve skew: the model learned from one kind of number and is then served another. You will see it cost more than anything else in this lesson.

Bounded vs Unbounded: The Real Difference

Three tools come up again and again in this part of the course. Apache Spark runs batch jobs over many machines at once. Apache Flink runs stream jobs. Apache keeps records in order, in named streams called topics, so that other programs can read them. Strip away the tools, and there is one idea underneath everything. Tyler Akidau's essay "Streaming 101" (O'Reilly, August 2015) suggests two words for it that many engineers now use: bounded and unbounded data. His description of unbounded data is "a type of ever-growing, essentially infinite data set".

A batch job works on bounded data: a chunk with a clear start and a clear end, like all of yesterday's rides sitting in a folder of files. The job reads the chunk, computes something, writes the result and stops. It knows when it is finished, because the data runs out.

A stream works on unbounded data: records that keep arriving, one at a time, with no last row. The job never finishes. It handles each record as it lands, and then waits for the next one.

A hand-drawn sketch headed sketched from the real hours: yesterday and today, titled one day has a last hour; the stream never does. At the top, bounded: 7 August, all 24 hours, as the nightly job reads them, two rows of twelve boxes holding the real rides for each hour: 47, 18, 13, 6, 9, 36, 179, 502, 705, 327, 250, 214, then 283, 253, 261, 306, 445, 868, 814, 610, 448, 317, 224, 138, with the first and the last boxes marked. Under them: first hour, 00:00, 47 rides; last hour, 23:00, 138. Below, unbounded: 8 August, arriving one hour at a time: 58, 23, 6, 7, 7, 43, 173, 482, 737, then dots and an arrow running on, and the words 00:00 to 08:00 so far, there is no last row, only the next one. Beneath: the nightly job read the whole bounded day and stopped, and its last number, 138, became the next day's rides in the last hour for all 24 hours.

These are two real days from the lab's data. The nightly job reads the whole of 7 August 2012, 24 hours from 47 rides in the first hour to 138 in the last, and stops. The next day is still arriving. Remember that last number, 138. It is what the nightly job hands the model as "rides in the last hour" for every hour of 8 August.

Notice what falls out of the split. Because batch has a fixed chunk, it can be scheduled to run at 00:00. It can work through millions of rows at once. And it costs almost nothing between runs, because nothing is running.

Because a stream has no end, it must run all the time, so it costs all the time. It handles one record at a time. And it has to cope with messy things, like records that arrive late or out of order. Bounded gives you cheap and a little old. Unbounded gives you fresh and more expensive. Every design later in this lesson, Lambda, Kappa, watermarks and delivery guarantees, is a consequence of this one split.

Latency vs Throughput vs Cost

Engineers love to say streaming is better because it is faster. That is not the trade. The real trade has three sides, and you rarely get to win all three. here means how old a result is by the time it is used. means how many rows a system handles for the work it does.

SideBatchStreaming
LatencyMinutes to hours, set by the scheduleSeconds or less
ThroughputVery high, in bulk: one run scans millions of rows at onceAlso high on a big enough cluster, but each record is handled on its own, with saved state, so each row takes more work
CostLow: machines run only on scheduleHigher: machines run all the time, plus saved state
Running itSimple: a job that starts and stopsHard: late records, records out of order, always on call
Doing it againEasy: run the job again on the old filesHarder: replay the log, or keep a second path

What the Lab Ran

I wrote the lab's design into the docstring of its script, batch_stream_demo.py, on 30 September 2026, before it ever ran. The script is small on purpose: it is the same file you can copy from this lesson and run yourself.

An editorial page in four labelled zones, headed what the lab ran: batch_stream_demo.py, designed before it ran, titled one dataset, two counts, five ways to compute them. The data: UCI Bike Sharing, hourly rentals of Capital Bikeshare bikes in Washington D.C., 2011 and 2012; 17,379 hours have a row, and 165 clock hours have none and count as zero rides. The split: trained on 13,742 hours from 2011-01-08, tested on the last 3,476, from 2012-08-07 to 2012-12-31. The inputs: the calendar, meaning hour, weekday, month, year, working day and holiday, plus two counts of past rides, the last hour and the last 7 days per hour; no weather. What changed: only when the two counts were computed, and which kind the model learned from; every model HistGradientBoostingRegressor, seed 0. Beneath: the time of day, the shuffles, one sketched day, the counts written twice and one crash came after the results.

The data is UCI Bike Sharing, the hourly count of bikes rented from Capital Bikeshare in Washington D.C. in 2011 and 2012 (Hadi Fanaee-T and Joao Gama, 2013). It is the same data the chapter on why production breaks used. It has 17,379 rows, one per hour, in time order. The file has the year, the month, the weekday and the hour, but no day of the month. So the script rebuilds the calendar: a new day starts whenever the weekday changes. The report checks the rebuilt date against every row's own year, month and weekday, and every row agrees.

Two full years have 17,544 hours, so 165 clock hours have no row at all. Most of them fall in the small hours of the night. The longest run is 36 hours, from 01:00 on 29 October 2012 to 12:00 on 30 October. That gap matches Capital Bikeshare's shutdown for Hurricane Sandy: a local news site, ARLnow, reported on 31 October that "the system did shut down for about 36 hours". A count of rides is zero when no rides happen. So the script counts an hour with no row as zero rides, and never scores a guess for it.

The task: at the start of each hour, guess how many bikes will be rented in that hour. The model learns from the first 80% of rows: 13,742 hours, starting on 8 January 2011. The first week has no full week behind it, so it is skipped. The model is tested on the last 3,476 hours, from 7 August to 31 December 2012. The real test hours averaged 248.8 rides.

. Every model sees the calendar: the hour, the weekday, the month, the year, and whether the day is a working day or a holiday. I left the weather out on purpose. At the start of an hour nobody knows that hour's weather for sure. And I wanted the only inputs that depend on the pipeline to be two counts of past rides. The is the rides in the last hour, which changes every hour. The is the rides per hour over the last 7 days, 168 hours, which changes slowly.

The Main Run

Here is what the demo stored in results/bv-demo.json.

A two-column table headed the main run, bv-demo.json: mean absolute error, titled fresh counts cut the error by almost half; skew tripled it. Left, design; right, MAE in rides an hour. Calendar only, no counts: 67.8. Nightly, both counts from 00:00: 56.4. Stream, both counts fresh: 30.5. Split, fast fresh, slow nightly: 30.6. Skew, trained fresh, served nightly: 173.9. Beneath: real rides averaged 248.8 an hour; 3,476 test hours, one run each.

With no counts at all, the calendar alone was off by 67.8 rides an hour on average. Adding both counts from the nightly job brought that down to 56.4. Adding them fresh every hour brought it down to 30.5, so fresh counts cut the error by almost half against the nightly ones.

The split design is the one to look at twice. It took the fast count from the stream and the slow count from the nightly job, and it scored 30.6: 0.09 rides an hour worse than streaming both. The 0.09 itself is too small to read anything into. What matters is the rest of the drop. Streaming only the last-hour count took the nightly design's 56.4 down to 30.6, almost all the way to the 30.5 of streaming both. On this data one count needed a stream and the other did not, so the decision belongs to each input, not to the model.

A bar chart headed bv-demo.json: error on the 3,476 test hours, titled the error of each design, in rides an hour. Five bars on a scale from 0 to 180: calendar about 68, nightly about 56, stream about 31, split about 31, and skew, drawn as its own series labelled trained fresh, served nightly, about 174. Beneath: 67.8, 56.4, 30.5, 30.6, 173.9; lower is better; the split design lost only 0.09 rides against the full stream; one run each.

And then skew: 173.9. That model was trained on fresh counts and then served the nightly ones, and it was far worse than knowing nothing but the calendar. The next slides take these numbers apart. Which count needed to be fresh? Where in the day did the difference sit? Why did skew do so much harm? And how often must a batch job run before it gets close to a stream?

Which Count Needed to Be Fresh

Why did one count need a stream and the other not? The simplest measure needs no model at all. On every test hour, compare the count the nightly job served with the fresh count, and average how far apart they are.

Two panels headed the nightly job against fresh counts, all test hours, titled the same job: one count nearly useless, one nearly exact. Rides in the last hour: 186.1 rides off, on average; fresh mean 248.8, nightly mean 111.3. Rides an hour over 7 days: 2.7 rides an hour off, on average, on a fresh mean of 250.3. Beneath: the nightly job gave every hour of the next day the rides from 23:00 to midnight.

The fast count was off by 186.1 rides on average. The fresh last-hour count averaged 248.8 over the test hours. The nightly one averaged 111.3, because it is always the rides from 23:00 to midnight, a quiet hour. So for every hour of the next day, the nightly job served the number of a late evening. The slow count was off by only 2.7 rides an hour, on a count that averaged 250.3. A week's average barely changes in a day, because most of the hours it covers are the same whether you compute it at midnight or at 18:00.

The nightly values were 11.5 hours old on average when they were used. The job runs at 00:00, so an hour at 23:00 reads a value 23 hours old. But age alone did not decide the damage. What decided it was how much the value changed in that time. So the same age cost the fast count almost everything and the slow count almost nothing.

A hand-drawn table headed sketched: a real day, Wednesday 2012-08-08, every third hour, titled the nightly count said the same thing all day. Columns for the hours 00, 03, 06, 09, 12, 15, 18 and 21. Real rides: 58, 7, 173, 341, 280, 278, 862, 381. Last hour, fresh: 138, 6, 43, 737, 239, 236, 858, 500. Last hour, nightly: 138 in all eight. Stream guess: 65, 6, 158, 340, 286, 266, 813, 352. Nightly guess: 64, 18, 157, 329, 268, 269, 725, 294. Skew guess: 65, 93, 356, 188, 202, 198, 183, 129. Beneath: picked by a rule fixed before looking at any day, the first whole working day of the test; error over the whole day: stream 20.1, nightly 39.5, skew 201.7.

Here is one real day, picked by a rule I fixed before looking at any day's numbers: the first whole working day in the test. At 00:00 the two versions of the fast count agree, 138, because the nightly job has only just run. From then on the nightly count stays at 138. The fresh one follows the day: 6 at 03:00, 737 at 09:00 just after the morning rush, and 858 at 18:00.

The stream's guesses followed the real rides closely. The nightly model's guesses were close in the quiet hours and fell behind in the busy ones: 725 against a real 862 at 18:00. Over the whole day the stream was off by 20.1 rides an hour and the nightly model by 39.5.

Trained Fresh, Served Stale

The skew design was trained on fresh counts and served the nightly ones, and it did worse than no counts at all: 173.9 against the calendar's 67.8. The nightly design was served exactly the same nightly counts and scored 56.4. The only difference between the two is what each one learned from.

Two panels headed served the same nightly counts; trained two ways, titled trained on fresh, served stale: worse than no counts. Trained nightly: 56.4 rides off an hour, trained on the same stale counts it is served. Trained fresh: 173.9 rides off an hour; the calendar alone scored 67.8. Beneath, after the results: with the fast count shuffled, the stream model's error went from 30.5 to 224.6; with the slow count shuffled, to 34.8.

To see why, I shuffled one input at a time across the test hours, after the results, and measured the stream model again. Shuffling breaks the link between an input and the right answer.

Each value on its own is still a real count, but the pairs can be odd: a shuffled count can put 800 rides next to 03:00. So how much worse the model gets tells you roughly how much it leans on that input. With the fast count shuffled, its error went from 30.5 to 224.6. With the slow count shuffled, it went to 34.8. The stream model leaned heavily on the last hour's rides. Served 138 at 18:00 on 8 August, it guessed 183 for an hour that had 862 rides.

The nightly model learned from the same stale value it would later be served. Here is one possible reading, which I did not test on its own. A late-evening count told it little about 18:00. So it learned to lean on the calendar and the week's average instead, and a stale count could not mislead it much.

This is the rule that matters most in this lesson, and it is not batch or stream: train on the same numbers you will serve. A stale input with a model that learned from stale values was a moderate loss here. A stale input fed to a model that expects a fresh one was a large loss, and nothing raised an error.

The damage depends on which way the mismatch runs. After the results, I measured the reverse too: a model trained on the nightly counts and served fresh ones. It scored 61.3. That is far better than 173.9, but still worse than the 56.4 the same model scored when it was served the nightly counts it learned from. A fresher input did not help a model that had learned what stale ones mean. The chapter on why production breaks measured the same kind of failure for inputs computed a different way, such as temperatures sent in Fahrenheit. Here the value was right in every way except the moment it was computed.

How Often Must the Job Run?

Between a nightly job and a stream there is a whole range of choices: run the same batch job every few hours. Spark's Structured Streaming works close to this by default. Its documentation says queries "are processed using a micro-batch processing engine": small batch jobs run one after another. The demo measured the range, with both counts from a job run every 1, 2, 3, 6, 12 and 24 hours, each model trained and served alike.

A bar chart headed bv-demo.json: the batch job run more often, trained and served alike, titled every step away from fresh cost something. Six bars on a scale from 0 to 70, for a job run every hour, every 2, 3, 6, 12 and 24 hours: about 31, 34, 37, 42, 46 and 56, rising from left to right, all below a dashed line near 68 labelled no counts. Beneath: 30.5, 33.8, 37.3, 42.2, 45.9, 56.4; the dashed line is the calendar alone, 67.8; one run each.

Every step away from fresh cost something. A job every hour scored 30.5, every 2 hours 33.8, and every 3 hours 37.3. Every 6 hours it scored 42.2, every 12 hours 45.9, and every 24 hours 56.4.

None of them came near the calendar alone, 67.8, so even a count from last night carried some information. How far the fast count was off rose from 38.0 rides for a job every 2 hours to 147.9 for a job every 6 hours. The job every 12 hours runs at 00:00 and at 12:00. Its fast count was off by a little less than the 6-hour one, 134.2, although its model still scored worse. I did not study why.

An isometric drawing of six blocks, heights to scale, headed bv-demo.json: how many times the job ran over the test weeks, titled the cost side: fresher means more runs. From left to right: 3,516 for every hour, the tallest; 1,758 for every 2 hours; 1,172 for every 3 hours; 586 for every 6 hours; 293 for every 12 hours; 146 for every 24 hours, the shortest. Beneath: their errors, 30.5, 33.8, 37.3, 42.2, 45.9 and 56.4; one update per ride would be 864,671 updates.

The other side of the trade is how often the job runs. Over the test weeks, a job every hour ran 3,516 times and a job every 24 hours ran 146 times. A real stream that updates on every ride would have done it 864,671 times. That count of runs is the real part of the cost here. What each run costs in money depends on the system, and I did not measure it.

So the practical question is not "batch or stream". It is: for this input, how much error does each extra run buy? Here, going from a job every 24 hours to a job every 3 hours cut the error from 56.4 to 37.3, for 1,172 runs instead of 146. Going the rest of the way to every hour cut it to 30.5, for 3,516. Whether that is worth it is a business question, and now it is one you can answer with numbers.

The Lambda Architecture: Have Both

For years, teams wanted the freshness of a stream and the exactness of batch at the same time. The Lambda architecture was the classic answer. Nathan Marz described it in a post called "How to beat the " (13 October 2011), with a batch layer and a realtime layer. The CAP theorem is a well-known result about systems spread over many machines. Sometimes part of the network is cut off from the rest, which is called a network partition. The theorem says that then the system has to give up either consistency or availability.

His book Big Data, written with James Warren, "presents the Lambda Architecture" under that name, and calls the fast path the speed layer.

The idea: send every event down two paths at once.

An architecture diagram with product logos, headed Nathan Marz's design, 2011, drawn with today's tools, titled Lambda: two paths, one answer, the logic written twice. At the top, the Apache Kafka logo: every ride, as an event, one log feeds both paths. On the left, the batch layer: the Amazon S3 logo, all history, as files, the complete record, cheap to reread; then the Apache Spark logo, exact counts, nightly, recomputed from all history. On the right, the speed layer: the Apache Flink logo, the newest hours, counts what the batch has not reached yet. Both paths lead down to the Apache Cassandra logo: one serving store, batch values for old hours, speed values for new ones. Beneath: drawn with today's tools; his post used Hadoop and Storm. Every count lives in two programs, one for each path, and the two must agree to the last ride.

The batch layer rereads the full history and recomputes exact values. It is slow, running on a schedule, but its output is complete, so it is the source of truth. The speed layer follows the same events and covers the hours the batch layer has not reached yet.

A serving store merges the two. For older hours it serves the batch value. For the newest hours it serves the speed layer's value, until the batch layer catches up and replaces it. In Marz's own words: "everything the realtime layer computes is eventually overridden by the batch layer." Because its values only live until the exact batch value arrives, the speed layer can afford to be roughly right.

It works. But look at the diagram and you will see the tax: the same count is written twice, once for the batch engine and once for the stream engine. That is two programs that must agree exactly, or the model's inputs quietly change depending on which path produced them. Jay Kreps put it plainly in 2014. Keeping code that must give the same result in two systems, he wrote, "is exactly as painful as it seems like it would be."

The Kappa Architecture: Just Stream Everything

On 2 July 2014, Jay Kreps, one of the original creators of Apache , wrote "Questioning the Lambda Architecture" for O'Reilly Radar. His complaint was the tax on the last slide: why keep two copies of every computation? His answer was one path, the stream. He named it almost in passing: "Maybe we could call this the Kappa Architecture, though it may be too simple of an idea to merit a Greek letter."

A hand-sketched diagram headed Jay Kreps, Questioning the Lambda Architecture, July 2014, titled Kappa: one program; to reprocess, replay the log. A cylinder, the log, kept for 30 days, feeds a box, job, version 1, which writes a cylinder, table A, read until the switch. The log also feeds a second box, version 2, same code, read from the start, which writes table B. Table B leads to the model, which switches to B once B catches up. Beneath: his words, maybe we could call this the Kappa Architecture; the logic lives in one program, and reprocessing is a second copy of it reading the log from the start.

The trick is that the event log keeps enough history to serve as the source of truth. Reprocessing, the thing the batch layer was for, becomes a replay of the log. His recipe goes like this. First, keep as much history in Kafka as you may need to reprocess. In his words: "if you want to reprocess up to 30 days of data, set your retention in Kafka to 30 days."

When the logic changes, start a second copy of the same stream job. It reads from the beginning of the log and writes to a new table. When it has caught up, switch the application to the new table, then stop the old job and delete the old table. One program. Fix a bug once, and the replay applies the fix everywhere. The switch also gives you a clean way back: if the new version is wrong, you never switch, and the old job is still running.

The cost is real. A full recompute means replaying the whole log through a stream engine. A stream engine reads records one at a time and keeps a running state for each key, while a batch engine scans compressed files in bulk. For a pure historical total, batch is often the faster tool. And you must keep a lot of history in the log. Kafka can now move old parts of its log to cheap object storage. This tiered storage arrived as an early-access feature in Kafka 3.6 (October 2023). It was declared production-ready in Kafka 3.9 (November 2024).

Event Time, Processing Time, and Watermarks

The next idea decides whether a streaming input is right at all.

Event time is when a thing actually happened. Processing time is when your job saw it. In Akidau's words, event time is "the time at which events actually occurred" and processing time is the time "at which events are observed in the system". In a stream the two drift apart all the time, because networks lag and phones go through tunnels. A ride that really started at one moment might reach the job several seconds later.

It is tempting to say batch has no such problem, because all the data is already in the file. It has the same problem, hidden at the cut-off. A nightly job over "yesterday" reads whatever had arrived by 00:00. A record from 23:59 that arrives at 00:03 is simply missing from the result, and stays missing until someone runs the job again.

Akidau makes the same point about batch jobs that cut data into fixed windows. He asks: "what if some of your events are delayed en route to the logs due to a network partition?" A network partition is the break in the network described on the Lambda slide. The difference is that a stream has to decide, every minute, how long to wait.

Before the timing problem, there is a shape problem. You cannot count or average over an unbounded stream until you cut it into finite pieces, called windows. And the rule you pick for the cut changes the number. Apache Flink's documentation describes the three common choices.

A hand-drawn sketch headed three ways to cut a stream into pieces, titled tumbling, sliding, session: the cut changes the number. Three lines carry the same eleven dots, one dot per record. Tumbling: four boxes of equal length side by side, touching, never overlapping. Sliding: three bars of equal length, each starting a quarter of the way along from the one before, so they overlap. Session: four boxes of different lengths, each around one burst of dots, with empty gaps between them. Beneath: tumbling, fixed pieces that never overlap; sliding, a fixed length that moves on a shorter step, so pieces overlap; session, a piece ends after a gap with no records.

The window type is a modelling choice, not a detail. Tumbling windows "have a fixed size and do not overlap". They suit a number like rides per clock hour, the way the lab's data is cut. Sliding windows have a fixed size and move forward on a shorter step, so they overlap. They are how "transactions in the last 5 minutes" stays current every 30 seconds.

Exactly-Once vs At-Least-Once, and How This Feeds ML

One more streaming decision has direct consequences for your data: the delivery guarantee. Stream jobs crash and restart, and what the restart does to your counts is the delivery guarantee.

At-least-once means every record is processed, but after a crash some may be processed twice. Flink's documentation says it plainly: "Nothing is lost, but you may experience duplicated results (at least once)". For a count, a duplicate is an extra.

Exactly-once means every record affects the result one time, even across crashes. The name describes the effect on the result, not on the wire. A record may still be read again after a crash, but it lands in the result once. Flink does this by saving its state together with its position in the input, which is called a checkpoint.

The promise holds end to end only if the place the job writes to takes part. In Flink's words, "your sources must be replayable, and your sinks must be transactional (or idempotent)". A sink is where the job writes its results. Idempotent means that writing the same thing twice leaves the same result as writing it once. A producer is the program that writes records into , and a transaction is a group of writes that either all become visible or none do. Kafka added idempotent producers and transactions in version 0.11, in 2017, so that the log itself can take part.

After the results, I built one crash on purpose, on the real data, to see what a duplicate does to a count. The stream job keeps one counter per clock hour in a database, and adds one to it for every ride. The fast count reads the last hour's counter, and the slow count sums the last 168.

The job records how far it got every 6 hours. It crashes once, at the start of 09:00 on Monday 13 August 2012, the first Monday of the test. It has counted 06:00, 07:00 and 08:00, but has not yet recorded that it did. On restart it goes back to 06:00 and counts those three hours again.

A sequence diagram with four columns, the rides, stream job, counters and the model, headed one crash, built on purpose: 2012-08-13, after the results, titled at-least-once: three hours counted twice. Step 1, at 06:00 the stream job saves its place. Step 2, the rides from 06:00 to 09:00 arrive. Step 3, the job adds 159, 436 and 673 to the counters. Step 4, at 09:00 it crashes before saving its place. Step 5, it restarts at 06:00 and adds 159, 436 and 673 again. Step 6, the model asks the counters for the last hour. Step 7, the counters send 1,346, really 673. Beneath: the 09:00 guess was 470 instead of 322, where the real count was 305; the 7-day count ran 7.55 rides an hour high for 166 hours, then 6.60 and 4.01; exactly-once saves the counters and the place together.

How Uber, Netflix, and DoorDash Actually Split It

Three teams that wrote about their systems each run both batch jobs and streams. I checked every quote below against the company's own post on 30 September 2026.

Uber's 2017 post about its Michelangelo platform, by Jeremy Hermann and Mike Del Balso, describes two ways its models' inputs are computed. Three names in it need a word each. HDFS, the Hadoop Distributed File System, stores very large files across many machines. Apache Cassandra is a database that spreads its data over many machines. Samza is a stream processing framework, like Flink. The first is "bulk precomputing and loading historical features from HDFS into Cassandra on a regular basis": batch. The second is to "publish relevant metrics to Kafka and then run Samza-based streaming compute jobs to generate aggregate features at low ": a stream.

Its example is the same pair as this lesson's lab. UberEATS computes "a restaurant's average meal preparation time over the last seven days" in batch, and the same average "over the last one hour" in the stream. The streaming values are also logged back to HDFS for later training. The post says the near-real-time path "ensures that the same data is used for training and serving": this lesson's skew rule, built into the platform.

DoorDash wrote in 2020, in a post by Arbaz Khan and Zohaib Sibte Hassan, that "almost all of the features get updated every day". Its example of a real-time input, "average delivery time for orders from a store in the past 20 minutes", is instead "updated uniformly throughout the day". In 2021, Allen Wang and Kunal Shah described Riviera, a framework built on Apache Flink that computes such inputs. One example is the "total orders confirmed by a store in the last 30 minutes", over a rolling window that refreshes every minute.

Riviera serves them from a built on Redis: a system that keeps a model's input values ready to be read quickly.

Netflix wrote in 2013, in a post by Xavier Amatriain and Justin Basilico, that its recommendations are computed in three ways. Offline computation "allows for more choices in algorithmic approach such as complex algorithms and less limitations on the amount of data that is used". Online computation "can respond quickly to events and use the most recent data". Nearline computation sits between the two, and stores its results instead of serving them at once. Netflix's Keystone platform runs its streams on and Flink (2018), and a 2022 paper describes Netflix search adapting to members' "interactions from the current session".

Try It Yourself

This script is the lab. It downloads the Bike Sharing data and rebuilds the clock. It computes the two counts as a stream and as a nightly job would, and trains the five designs. Then it runs the batch job every 1, 2, 3, 6, 12 and 24 hours. It prints each design's error, how far the nightly counts are from fresh ones, and how many times each job ran. It does not need a GPU; when I ran it, it finished in under ten seconds.

A real screenshot of VS Code with batch_stream_demo.py open at the top of the file: the docstring, which says what the script needs and how to run it, and holds the design written on 2026-09-30 before the first run, with the data, the fast and slow counts, the five designs and what is reported; then the imports; the code that downloads Bike Sharing and rebuilds each row's clock hour from its weekday and hour is further down. Beneath: copy it from the box on the slide.

Before you run this lab. You need Python 3 and two libraries: pip install scikit-learn pandas. scikit-learn holds the model and the download, and brings NumPy with it; pandas holds the table. The first run downloads the Bike Sharing data from OpenML (under 1 MB), so it needs an internet connection once. After that, scikit-learn keeps a copy in a folder in your home directory (scikit_learn_data).

I ran it with scikit-learn 1.9.1 on a Mac. scikit-learn runs the same way on Windows and Linux, but I have not checked the numbers there. Another version of scikit-learn may give different decimals, which is why the first line printed is the version. Give it a file name, python batch_stream_demo.py out.json, and it also saves every number at full precision; that is how results/bv-demo.json was made.

"""Batch or streaming? The same two inputs, computed once a night and every hour.

Lesson 1 of 'Data Engineering for ML', made small. It needs Python 3 with
scikit-learn and pandas (pip install scikit-learn pandas). The first run
downloads UCI Bike Sharing from OpenML (under 1 MB) and keeps a copy.
    python batch_stream_demo.py            # print the results
    python batch_stream_demo.py out.json   # and save every number

Design, written 2026-09-30 before the first run:
  Data: 17,379 hours of bike rides in Washington D.C., 2011 and 2012, in
  time order. At the start of each hour a model predicts that hour's rides.
  It learns from the first 80% of rows and is tested on the last 20%. The
  weather is left out on purpose, so the only inputs a pipeline computes
  are two counts of past rides:
    fast  rides in the last hour               (changes every hour)
    slow  rides per hour over the last 7 days  (changes slowly)
  An hour with no row in the file counts as zero rides and is not scored.
  Training skips the first 7 days, which have no full week behind them.
  A stream updates both counts every hour. A nightly batch job computes
  both at 00:00, and every hour of the next day reads those values.
  Five designs. Each model is trained and tested on the same kind of input,
  except the last one:
    calendar  no ride counts at all: hour, weekday, month and so on
    nightly   both counts from the nightly job
    stream    both counts fresh every hour
    split     the fast count from the stream, the slow one nightly
    skew      trained on the stream's counts, served the nightly ones
  Then the batch job run every 1, 2, 3, 6, 12 and 24 hours. Every 1 hour
  must equal stream, and every 24 hours must equal nightly.
  Reported: mean absolute error (MAE, rides an hour) on the test hours,
  how far each nightly count is from the fresh one, and how many times each
  design runs its job. HistGradientBoostingRegressor with seed 0, one run,
  no significance test.

Author: Roni Das
Created: 2026-09-30
"""
import json
import sys

import numpy as np
import sklearn
from sklearn.datasets import fetch_openml
from sklearn.ensemble import HistGradientBoostingRegressor

data = fetch_openml("Bike_Sharing_Demand", version=2, as_frame=True, parser="auto").frame
rides = data["count"].to_numpy()

# The file has the year, month, weekday and hour, but no day of the month.
# A new day starts when the weekday changes; a jump of two skips a day.
weekday, hour = data["weekday"].to_numpy(), data["hour"].to_numpy()
day = np.r_[0, np.cumsum((weekday[1:] - weekday[:-1]) % 7)]
clock = day * 24 + hour  # hours since 2011-01-01 00:00

# Rides on a full hourly clock. total[k] is every ride before hour k.
on_clock = np.zeros(clock[-1] + 1, dtype=int)
on_clock[clock] = rides
total = np.r_[0, np.cumsum(on_clock)]


def counts(asof):
    """The two inputs as a job that ran at the start of hour `asof` saw them.

    The first week's rows get meaningless values here; none is used."""
    fast = total[asof] - total[asof - 1]
    slow = (total[asof] - total[asof - 168]) / 168
    return fast, slow


calendar = np.column_stack([
    hour, weekday, data["month"], data["year"],
    data["workingday"].astype(str) == "True", data["holiday"].astype(str) == "True",
]).astype(float)


def table(fast_every=None, slow_every=None):
    """The calendar, plus each count from a job run every N hours."""
    cols = [calendar]
    if fast_every:
        cols.append(counts(clock // fast_every * fast_every)[0])
    if slow_every:
        cols.append(counts(clock // slow_every * slow_every)[1])
    return np.column_stack(cols)


cut = int(0.8 * len(rides))  # learn from the first 80% of rows
train = np.flatnonzero((np.arange(len(rides)) < cut) & (clock >= 168))
test = np.arange(cut, len(rides))


def mae(learn_from, serve_with):
    model = HistGradientBoostingRegressor(random_state=0)
    model.fit(learn_from[train], rides[train])
    error = model.predict(serve_with[test]) - rides[test]
    return float(np.abs(error).mean())


out = {"scikit_learn": sklearn.__version__, "rows": len(rides),
       "train_rows": len(train), "test_rows": len(test),
       "hours_without_a_row": int(len(on_clock) - len(rides)),
       "test_mean_rides": float(rides[test].mean())}
print(f"scikit-learn {sklearn.__version__}, {len(test):,} test hours")
print(f"real rides an hour in the test: {rides[test].mean():.1f}")

X = {"calendar": table(), "nightly": table(24, 24),
     "stream": table(1, 1), "split": table(1, 24)}
out["mae"] = {name: mae(X[name], X[name]) for name in X}
out["mae"]["skew"] = mae(X["stream"], X["nightly"])
print("error (MAE, rides an hour), one run each:")
for name, value in out["mae"].items():
    print(f"  {name:<9} {value:6.1f}")

# How far the nightly job's counts are from fresh ones, on the test hours.
fresh, night = counts(clock), counts(clock // 24 * 24)
out["nightly_off_by"] = {k: float(np.abs(night[i] - fresh[i])[test].mean())
                         for i, k in enumerate(["fast", "slow"])}
print("nightly count off from the fresh one, on average:")
print(f"  fast  {out['nightly_off_by']['fast']:6.1f} rides")
print(f"  slow  {out['nightly_off_by']['slow']:6.1f} rides an hour")

# The same batch job, run more often: both counts, trained and served alike.
out["every"] = {}
print("batch job every N hours: runs in the test, MAE")
for n in [1, 2, 3, 6, 12, 24]:
    runs = int(clock[-1] // n - (clock[cut] - 1) // n)  # starts in the test
    out["every"][n] = {"runs": runs, "mae": mae(table(n, n), table(n, n))}
    print(f"  every {n:>2}h {runs:>6,} runs  {out['every'][n]['mae']:6.1f}")
out["test_rides"] = int(rides[test].sum())
print(f"one update per ride instead: {out['test_rides']:,}")

if len(sys.argv) > 1:  # a file name was given: save every number too
    json.dump(out, open(sys.argv[1], "w"), indent=1)

The Lab Report

A real terminal recording headed python batch_stream_report.py, titled every table in this lesson, from the stored run and the data. It opens: Bike Sharing, 17,379 hours, tested on 3,476, 2012-08-07 to 2012-12-31, 30 checks: all agree. Then seven numbered sections: 1, the headline, with the five designs' MAE and the job every 1 to 24 hours with its runs, MAE and how far each count is off; 2, the error in eight 3-hour blocks for calendar, nightly, stream and skew, with how far the nightly fast count is off; 3, the stream model's error with the fast count shuffled, 224.6, and the slow count shuffled, 34.8; 4, one real day, Wednesday 2012-08-08, every third hour; 5, the same counts written twice, the fast count different on 5 test hours and the slow count on 690; 6, one crash and a replay on 2012-08-13, with 1268 extra rides; 7, after a review, the reverse of skew, trained on the nightly counts and served fresh ones, at 61.3, and the replay's slow count 7.55 too high on 166 hours, then 6.60 and 4.01. Beneath: the lab's own report; it recomputes the counts another way, fits every model again, and stops unless every stored number comes back.

The report lives in scripts/labs/dataeng/batch_stream_report.py. It reads the demo's stored files, results/bv-demo.json and the printed run, and the Bike Sharing data from scikit-learn's local copy. It does not trust the demo's arithmetic. It rebuilds the calendar with real dates, and checks it against every row's year, month and weekday. It computes the two counts a different way, with pandas on an hourly clock instead of the demo's running total. It fits every model again. And it stops unless every stored error, count and run comes back exactly. Every check agreed. It changes nothing in the demo's files.

Its json mode writes every number to results/bv-report.json, which the figures read. The demo mode checks the demo's printed run line by line, and the box mode writes the playground below and checks it against both files.

Here is what came before the run, in the demo's docstring: the data, the split, the two counts and the five designs. The job every 1 to 24 hours, and what to report, were there too. Then there is what came after I saw the results, written into the report's docstring before the report first ran. It covers the error by time of day, the shuffles, and the sketched day with the rule that picked it. It also covers the two versions of the counts, and the crash and replay. Every figure that shows one of those says "after the results".

Count It Yourself, No Model

This box has no model in it. It holds the real rides for every clock hour, from eight days before the test to the end of 2012. It also holds the guesses of three designs, stream, nightly and skew, for every test hour, to a hundredth of a ride. Each number is stored as two or three characters, in base 64, so that the box stays small. Base 64 just means counting with 64 different digits instead of 10. It runs in your browser.

The counts in the box are exact. So how far a job's counts are from fresh ones matches the lab to the last digit. The report checked that for a job every 1, 2, 3, 6, 12 and 24 hours. The guesses are stored to a hundredth of a ride. The report checked that every MAE the box prints, for the whole test and for each 3-hour block, reads the same as the lab's.

As it is, the box prints the number of test hours, their real mean, and the MAE of the three designs. Then it prints how far two jobs are off from fresh counts. A job run every 24 hours is off by 186.1 and 2.7, and a job run every hour by 0.0.

Then try off(3) and off(6) to see how far a job every 3 or 6 hours is off. Compare them with the chart of errors on the slide about how often the job must run. Try by_block('skew') next to by_block('nightly') to see where in the day the skewed model went wrong, and mae('stream', range(15, 18)) for one block on its own. Ask yourself which of these you could have measured before building anything.

The Code, Part by Part

Loading. fetch_openml("Bike_Sharing_Demand", version=2) downloads the data once and reads the local copy after that. rides is the number of rides in each row's hour.

The clock. The file has no day of the month, so day counts days. Each step between two rows adds the change in weekday, taken modulo 7 (the remainder after dividing by 7). That change is 0 within a day, 1 at midnight, and 2 when a whole day has no row. clock is then the number of hours since 00:00 on 1 January 2011. on_clock puts every row's rides at its clock hour, with zero for the 165 hours that have no row. total is a running total: total[k] is every ride before hour k.

counts(asof). The two inputs as a job that ran at the start of hour asof would have seen them. The fast count is total[asof] - total[asof - 1], the rides in the hour before. The slow count is the rides in the 168 hours before, divided by 168. Because both use the running total, each count is one subtraction.

table. The calendar columns, plus each count from a job run every N hours. clock // N * N is the start of the hour when that job last ran, so is the stream, the nightly job and the split.

How to Choose, One Input at a Time

A hand-sketched column of six boxes joined by arrows, headed choosing, one input at a time, titled measure, then choose; train on what you serve. 1, list every input and how often its value changes. 2, measure how far last night's value is from a fresh one. 3, shuffle it and see how much the model leans on it. 4, try the batch job every few hours before a stream. 5, train on the same numbers you will serve. 6, store each value with the time it was computed. Beneath: here, steps 2 and 3 alone said stream the fast count and keep the slow one nightly; that split scored 30.6 against 30.5.

List every input and how often its value changes. A customer's country barely changes. A week's average changes a little every hour. A count over the last hour changes completely every hour. Write it down for each input before you choose anything.

Measure how far last night's value is from a fresh one. On past data, compute each input the way a nightly job would and the way a stream would, and average the difference. It needs no model. Here it was 186.1 rides for the fast count and 2.7 rides an hour for the slow one, and those two numbers alone pointed at the right design.

See how much the model leans on it. Shuffle the input across rows and measure the model again. If the error barely moves, freshness cannot matter much for that input.

Try a batch job every few hours before a stream. It is the same code as the nightly job, run more often. Here a job every 3 hours scored 37.3 against 56.4 for a nightly one.

Train on the same numbers you will serve. Build the training table the same way the serving path computes its inputs. Here, training on fresh counts and serving nightly ones cost more than any other choice: 173.9. The reverse, which I measured after the results, cost less, 61.3 against 56.4, but still more than training and serving alike.

Store each value with the time it was computed. Then how old an input is becomes a number you can check, and a stale input can raise an alarm instead of passing silently.

When to Stream, When to Batch

A two-column table headed grounded in this lesson's numbers, titled batch it, or stream it? Left, batch it when: its value barely moves between runs, the nightly 7-day count was off by 2.7 rides an hour across the next day; the model does not lean on it; a job every few hours is close enough; you are building a training table from history. Right, stream it when: its value moves within hours, the nightly last-hour count was off by 186.1 rides; the model leans on it, shuffled, the error went from 30.5 to 224.6; even a job every 2 hours cost 3.31 rides an hour here; the guess is needed now, not tomorrow. Beneath: batch is the default; stream only what the model needs fresh.

Batch it when its value barely moves between runs. A 7-day average, a customer's country, last month's total. Here the nightly 7-day count was off from the fresh value by 2.7 rides an hour on average, across the next day. Keeping it nightly cost 0.09 rides an hour.

Batch it when the model does not lean on it. Even a fast-moving input is not worth a stream if shuffling it barely changes the error.

Batch it when a job every few hours is close enough. It keeps the simplicity of batch: a job that starts and stops, and runs again on the same files after a failure.

Batch it for training data. Months of history, computed cheaply and exactly, is what a batch job is for, as long as serving computes the same thing.

Stream it when its value moves within hours and the model leans on it. Here the nightly last-hour count was off from the fresh value by 186.1 rides on average, across the next day. Shuffling it took the error from 30.5 to 224.6. That input earned its stream.

Stream it when the guess is needed now. A card payment, a delivery estimate, a bike dock about to run empty. If a guess from last night's numbers is no use by the morning, the input that drives it must be fresh.

Do not stream everything because it sounds modern. A stream runs all the time, needs a watermark, a delivery guarantee and someone on call. Pay for it only for the inputs that need it.

What These Numbers Can and Cannot Tell You

A two-column page headed read before you trust these numbers, titled what these runs are, and what they are not. They are: one dataset, bike rides, Aug to Dec 2012; one model type, one run each; hourly data, a stream of hours; the design came first; blocks, shuffles and a day, after; the crash I built on purpose; two versions of two counts. They are not: not a rate for other data; not every kind of model; not a per-ride stream; not changed after the run; chosen knowing the results; not a failure found in the wild; not every way code can differ.

One dataset, one model type, one run each. Everything here is one city's bike rides in the second half of 2012, and one kind of model. With other data, other inputs and other models, the numbers would move. No significance test was declared, so I do not call any gap significant. The 0.09 between split and stream is too small to read anything into; what the split shows is that one count carried nearly all of the gain.

A stream of hours, not of rides. The data only has hourly totals, so my "stream" updates once an hour. A real stream updates on every ride, which would make the fresh count fresher still. I did not measure event time, watermarks or late records for the same reason.

No money. I counted how many times each job ran. What a run costs depends on the system, and I did not measure it.

The design came first. The five designs and the job every 1 to 24 hours were written down before the first run and not changed. The time-of-day blocks, the shuffles, the sketched day, the two versions of the counts and the crash came later. I designed them after I saw the results, and wrote them down before the report first ran. The reason I give for the nightly model's resistance to stale counts is my reading, labelled as such.

Faults I built on purpose. The crash, the replay and the second version of the counts are faults I put in to measure them. They are the kind of fault that happens in real systems, but they are not failures I found in someone's system.

What to Do Next

A hand-drawn list headed before you build a stream, titled five questions for every input. How fast?: how far is last night's value from a fresh one, on average? How much?: how much worse is the model with this input shuffled? How often?: would a job every few hours be close enough? Same both?: is the model trained on the same numbers it will be served? When?: is the time each value was computed stored with it? Beneath: here, a count from last night was 186.1 rides off; a week's average, 2.7.

Take a model your team runs and list every input it reads. For each one, ask the five questions on the card. The first two need only past data and an afternoon. Compute each input the nightly way and the fresh way, and average the difference. Then shuffle it to see how much the model leans on it. Most inputs will turn out to be like the 7-day count, and they can stay in batch. A few will be like the last-hour count, and those are the ones worth a stream.

Then check the fourth question for every input: is the model trained on the same kind of number it is served? The skew design in this lesson is the common way to get it wrong. The training table is built from fresh values for every past hour, because history makes that easy. Serving, meanwhile, reads a table the nightly job refreshes once a day.

Here that cost more than anything else measured: 173.9. The reverse costs less but is still a mistake. After the results I trained on the nightly counts and served fresh ones. It scored 61.3, worse than the 56.4 of serving the nightly counts it learned from.

A closing card headed to keep, titled choose freshness per input; train on what you serve. In large type: 30.5 or 56.4. Beneath: rides off per hour, with both counts fresh or both from last night. Then: keeping only the week's average nightly cost 0.09 rides an hour more; keep it in batch. Then: trained fresh and served stale, 173.9, worse than no counts at all. Last: one dataset, one run each: a way to measure, not a law.

The card keeps three numbers from the lab. With both counts fresh the guesses were off by 30.5 rides an hour, and with both from last night by 56.4. Keeping only the week's average nightly cost 0.09 more. Training on fresh counts and serving stale ones cost the most: 173.9.

The next lesson in this chapter is about data ingestion pipelines. It covers how records get from where they are made into the log or the lake that both kinds of job read from.

Knowledge Check

Knowledge Check

4 questions - Score 80% to pass

Q1

The split design took the last-hour count from the stream and the 7-day count from the nightly job, and scored 30.6 against 30.5 for streaming both. What does that show?

Q2

The skew design was trained on fresh counts and served the nightly ones. Why did it score 173.9, worse than the calendar alone?

Q3

Two versions of the 7-day count, one working on clock hours and one on the rows a job receives, disagreed on 690 test hours. What made them differ?

Q4

What does exactly-once do that at-least-once does not?

It helps to see the two shapes as real pipelines. Here is a common way to build the nightly job. The day's records land as files in object storage. A scheduler starts a job at 00:00. A batch engine reads the whole day and computes the numbers, and a table of results waits for the model.

Four cards joined by down arrows, headed one common way to build the nightly job, titled batch: read the whole day, compute, stop. Amazon S3, with its logo: the day's records land here as files, all day. Apache Airflow, with its logo: a scheduler, starts the job at 00:00. Apache Spark, with its logo: reads the whole day, computes both counts, then shuts down. A table of counts: the model reads it at every hour of the next day. Beneath: these are real tools that teams use for this, not the lab's; the lab does the same job in a few lines of NumPy, and over the test weeks it ran 146 times.

The economics of batch live in the gaps. You pay for the machines only while the job runs. The price of that deal is that every number is exactly as old as the last run. Nothing is computed again between runs, so at 18:00 the table still says what was true at 00:00.

Now the same records, wired differently. In a stream, each record goes onto a durable log the moment it happens. A durable log is a list of records kept on disk in the order they arrived, which readers can read again from any point. A job that never stops updates the numbers as each record arrives. The results land in a fast store, which the model reads when it needs them. , a common choice for that store, keeps its data in memory. Its documentation describes "sub-millisecond reads from memory", and the network between the model and the store adds its own time on top.

Four cards joined by down arrows, headed one common way to build the stream, titled streaming: never stop, update on every record. Apache Kafka, with its logo: a log, every ride as an event, kept in order for as long as it is set to keep them. Apache Flink, with its logo: always running, updates both counts as each event arrives. Redis, with its logo: holds the latest counts in memory, so a read takes under a millisecond. The model: reads the fresh counts at the start of each hour. Beneath: again real tools, not the lab's; the lab's stream updates once an hour, 3,516 updates over the test weeks, and one update per ride would be 864,671.

Here freshness has a running meter. The log, the job and its saved state all stay up around the clock, and you pay for every hour whether or not anyone asks for a guess. What you buy is a number that is never more than one record old.

Here is the made-up card example from the first slide, built as a stream. Step through it and notice two choices that later slides come back to. The count lives in Flink's own saved state. And the model may score a payment before that payment is in the count, which is fine here, because the swipe ninety seconds earlier is.

Read it as a set of levers. Streaming buys you latency, and you pay for it in cost and effort. Batch buys you cheap bulk work, and you pay for it in freshness. There is no free option. The skill is knowing which lever the input in front of you needs pulled, and refusing to pay for the others.

The old version of this lesson drew this trade as a chart and a set of cost bars. No numbers stood behind their shapes, so I have taken them out. In the lab below I measured the freshness side on real data and counted how many times each design runs its job. I did not measure money: the cost of a run on my laptop says nothing about a cloud bill.

The choice is not per model. It is per input, and it comes down to a short set of questions. Start from freshness, because freshness is the only thing a stream buys that a batch job cannot.

A flowchart headed per input, not per model, titled three questions before you stream an input. The first question, far from fresh the next day? No: compute it nightly. Yes: does the model lean on it? No: compute it nightly. Yes: is every few hours enough? Yes: run the batch job more often. No: stream it. Beneath: whichever you pick, train on the same numbers you will serve. Here, across the next day, the nightly slow count was off from fresh by 2.7 rides an hour and the fast count by 186.1 rides; shuffling the fast count took the stream's error from 30.5 to 224.6; a job every 2 hours scored 33.8.

Each question is something you can measure before you build anything, and the lab below measures all three. The last line under the chart is the rule that matters most, and it is not about batch or stream at all. Whichever you pick, the model must learn from the same kind of number it will be served.

If you remember one summary of the two shapes, make it this one. Read down each column and the character of each system is clear. Batch is the cheap, simple workhorse that is always a little old. A stream is the fresh, expensive specialist that asks more of the team.

A two-column table headed the two shapes, side by side, titled what each one is good at, and what it costs. Left, batch job: reads a bounded chunk, a first row and a last row; as old as its last run, here 11.5 hours old on average; runs on a schedule and costs nothing between runs; after a failure, run it again on the same files; a late record is missed until someone reruns the job; feeds training tables and inputs that change slowly. Right, stream: reads an unbounded stream, there is no last row; as old as the last record, here at most one hour; runs all the time, so it costs all the time; after a failure, replay from a saved position, and a record can count twice; a watermark decides how long to wait for a late record; feeds inputs that change within minutes or hours. Beneath: cheap and simple, and always a little old; fresh, and it asks more of the team.

Neither is better. They answer different questions. Batch answers "what is true over a long history, cheaply", and a stream answers "what is true right now". A mature platform runs both, and decides input by input which one each number comes from.

The inputs
fast count
slow count

A table of five rows headed the five designs, fixed before the run, titled the same model, trained and served five ways. Calendar: no ride counts at all, hour, weekday, month and so on. Nightly: both counts from the job at 00:00, for training and for serving. Stream: both counts fresh every hour, for training and for serving. Split: the fast count from the stream, the slow count from the nightly job. Skew: trained on the stream's counts, served the nightly ones. Beneath: then the batch job run every 1, 2, 3, 6, 12 and 24 hours, trained and served alike; every hour must equal stream, and every 24 hours must equal nightly.

The designs. A stream updates both counts at the start of every hour. A nightly batch job computes both at 00:00, and every hour of the next day reads those same values. Each model is trained and tested on the same kind of input, except the last one. That one, skew, is trained on fresh counts and then served the nightly ones. It is what happens when a team builds its training table one way and its serving path another.

Then the same batch job is run every 1, 2, 3, 6, 12 and 24 hours. On hourly data, a job run every hour is the stream. A job run every 24 hours is the nightly design. So those two had to come back exactly equal to those designs, and they did.

Every model is scikit-learn's HistGradientBoostingRegressor, many small trees of yes-or-no questions built one after another, with seed 0. Each design ran once. The design declared no significance test, which is a calculation of how likely a difference is to come from chance alone. So I describe the numbers and do not call any gap significant.

A sequence diagram with four columns, new hour, the model, the table and nightly job, headed one served hour, nightly design: Wednesday 2012-08-08, 18:00, titled the table answers with a number from last night. Step 1, the nightly job at 00:00 counts rides so far. Step 2, it sends the table last hour: 138. Step 3, the new hour asks the model to guess 18:00. Step 4, the model asks the table for the counts. Step 5, the table sends back last hour: 138. Step 6, the model answers with a guess of 725, where the real count was 862. Beneath: streaming, step 5 would send 858, the rides from 17:00 to 18:00, and the stream's guess was 813; nothing in this call says the value is 18 hours old.

Here is one served hour under the nightly design, from a real day in the test. At 18:00 on Wednesday 8 August 2012 the model asked the table for the counts. It got 138 for the last hour: the rides from 23:00 to midnight the night before. The real number for 17:00 to 18:00 was 858. The nightly model guessed 725 rides for 18:00, and the real count was 862. The stream model, which was sent 858, guessed 813. Nothing in that exchange says how old the value is.

After the results, I measured where in the day the fresh counts paid off, in blocks of three hours.

A bar chart headed error by time of day, 3-hour blocks, after the results, titled fresh counts paid off most from midday to evening. Eight pairs of bars, nightly and stream, on a scale from 0 to 120, for blocks starting at 00:00, 03:00, 06:00, 09:00, 12:00, 15:00, 18:00 and 21:00. Nightly: 15.7, 10.4, 57.6, 52.4, 84.7, 108.1, 82.2, 38.6. Stream: 12.0, 5.3, 39.4, 33.5, 36.1, 54.7, 40.2, 22.4. Beneath: real rides averaged 47 an hour from 00:00 and 458 from 15:00.

In the night the two designs were close: from 00:00 the nightly model was off by 15.7 and the stream by 12.0. The gap grew with the traffic. From 12:00 it was 84.7 against 36.1, and from 15:00, the busiest block with 458 rides an hour on average, 108.1 against 54.7. A stale count cost most exactly when there was the most to get wrong. For a real system that is worth knowing: the hours when a stale input hurts are often the hours that matter most to the business.

After the results, I measured how easily two careful versions of the same count disagree. I wrote the two counts a second time, the way a job that only sees the rows it receives would write them. The fast count became the rides in the previous row. The slow count became the average of the previous 168 rows. The demo's version works on the clock, and counts an hour with no row as zero. Both sound like the same definition, and neither is a bug you would spot in a code review.

Two panels headed the same two counts written twice, after the results, titled two correct-looking versions disagreed on 690 hours. Rides in the last hour: 5 hours of 3,476 differed, by at most 22 rides. Rides an hour over 7 days: 690 hours differed, by at most 82.3 rides an hour. Beneath: one counts clock hours, the other counts the rows it received; served the row version, the stream model scored 30.4 instead of 30.5, and 236 guesses moved; neither raised an error.

The fast count differed on only 5 of the 3,476 test hours, by at most 22 rides. The slow count differed on 690 hours, by up to 82.3 rides an hour. The worst was at 19:00 on 3 November 2012, a few days after the 36-hour Sandy shutdown. The clock version counted those 36 hours as zero rides and said 168.8. The row version skipped them, reached back further than a week, and said 251.1. Both versions sum whole rides, so neither difference comes from rounding.

Served the row version, the stream model scored 30.4 instead of 30.5, and 236 of its guesses moved. Here the disagreement happened to cost nothing, because this model leans little on the slow count. It raised no error either way. On an input the model does lean on, the same disagreement would not be free.

Lambda
Kappa
Paths to keep runningTwo (batch and speed)One (stream)
Logic writtenTwiceOnce
ReprocessingRun the batch layer againReplay the log
Speed of a full recomputeHigh (a tuned batch job)Lower (a full stream replay)
Main riskTwo programs drift apartLong retention, heavy replays

Kreps also looked at a middle way. You write the logic once, in a framework such as Summingbird, and it turns the logic into both a batch job and a stream job. His verdict was careful: "This definitely makes things a little better, but I don't think it solves the problem." In practice many teams sit somewhere between the two shapes. The question the lab keeps asking still applies: whichever design you run, do training and serving see the same numbers?

Session windows close when no record has arrived for a set time, which is how you measure one visit to a website or one shopping trip. Pick the wrong window and the input answers a question you did not ask. It shows up as a quietly worse model, not as an obvious crash.

Now the timing problem. Say you count transactions in the last 5 minutes. If you put each record in a window by when it arrived, using processing time, a late record lands in the wrong window and the count is wrong. So streaming engines put records in windows by event time, and they use a watermark to decide when a window is done.

A hand-sketched timeline headed when it happened against when the job saw it, titled a late record belongs to the window it happened in. Two arrows run left to right: when it happened, above, and when the job saw it, below. Above the top line, a box marks one window, followed by a shorter box marked grace. Three records, A, B and C, happen inside the window, and a line runs down from each to the moment the job saw it. A and B are seen soon after; the line from C slants to the right, so it is seen after the window has ended, at a point that falls within the grace period. A vertical line at the end of the window is labelled watermark: probably all seen up to here. Beneath: grouped by when it was seen, C lands in the wrong window; grouped by when it happened, it counts where it belongs, if it arrives within the grace period.

A watermark is the engine's moving promise: I have probably now seen every record up to this time. When the watermark passes the end of a window, the engine closes the window and sends out the result. To catch stragglers you set a grace period, which Flink calls allowed lateness. Then a record that arrives a little after the watermark still lands in the right window.

Flink's default allowed lateness is zero: "elements that arrive behind the watermark will be dropped", unless you send them to a separate output for late records. Set it too tight and you drop real records. Set it too loose and you hold windows open for a long time, spending memory and delaying results. Batch hides this cost at its cut-off. A stream has to pay it in the open.

My lab does not measure event time. Its data is already cut into clock hours, one row per hour, so every ride is in its hour before the model ever sees it. That is a limit of the data, and this slide gives definitions and sources, not measurements.

The three hours had 159, 436 and 673 rides, so 1,268 rides were counted twice. The fast count at 09:00 read 1,346 instead of 673. So the stream model guessed 470 rides for an hour that had 305; with the right count it guessed 322. The slow count ran 7.55 rides an hour too high for the next 166 hours. Then the doubled hours left its 7-day window one at a time. So it was 6.60 too high for one more hour, 4.01 for the hour after, and then right again.

Over those 168 hours the error went from 26.9 to 27.7 rides an hour. Over the whole test, the one crash moved the MAE from 30.52 to 30.56. So one crash made one guess badly wrong and a week of guesses slightly wrong, and nothing raised an error. I did not measure a job that crashes often.

Now connect it back to machine learning. These pipelines feed a model in two places, and both have to agree. The batch or historical path usually builds the training data: inputs computed over months of history, cheap and exact, used offline to fit the model. The stream path produces the inputs used at serving time. It computes the same definitions live and writes them to a fast store, which is read when a guess is needed.

Say the training data computes a 7-day count one way, and the serving path computes it a slightly different way, or double-counts after a crash. Then the model sees different inputs in training and at serving. The lab measured both of those for the 7-day count, and here they cost little. The second version of the count moved the stream model from 30.5 to 30.4. The one crash moved it from 30.52 to 30.56. The large cost, 173.9 against 56.4, came from the last-hour count, served stale to a model trained on fresh values. How much a mismatch hurts depends on how much the model leans on that input.

Here is a sketch of the settings that decide whether a streaming count is right. It is not a file any tool reads; the real setting names are in the comments.

# A sketch of the settings that decide a streaming count's correctness.
# Not a real config file: each comment names the real setting.
delivery:
  checkpointing_mode: exactly-once # Flink CheckpointingMode.EXACTLY_ONCE
  kafka_producer:
    enable.idempotence: true       # already Kafka's default if nothing conflicts
  kafka_consumer:
    isolation.level: read_committed  # read only committed transactions
  sink:                            # where the job writes its counts
    type: kafka                    # Flink's KafkaSink
    delivery_guarantee: exactly-once  # DeliveryGuarantee.EXACTLY_ONCE: writes in a
                                   # Kafka transaction, committed on each checkpoint
    transactional_id_prefix: ride-counts  # setTransactionalIdPrefix, required here
    # transaction.timeout.ms must cover a checkpoint plus a restart

time:
  # Event time has been Flink's default since version 1.12.
  watermark:
    strategy: bounded-out-of-orderness  # WatermarkStrategy.forBoundedOutOfOrderness
    max_out_of_orderness: 10s      # the watermark trails the newest event time by 10s
  window:
    type: sliding
    size: 5m                       # "transactions in the last 5 minutes"
    slide: 30s
    allowed_lateness: 1m           # a window is kept 1m past the watermark, so a record
                                   # up to about 70s behind the newest still counts;
                                   # later records are dropped

state:
  checkpointing:
    interval: 30s                  # how often state and position are saved together
    storage: s3                    # somewhere durable, so state survives a crash

Two details are easy to miss. Kafka's producer already has idempotence on: its documentation says "Idempotence is enabled by default if no conflicting configurations are set". And exactly-once in Flink is not complete without the sink. Flink's documentation says a KafkaSink set to exactly-once "will write all messages in a Kafka transaction that will be committed to Kafka on a checkpoint". It also needs a transactional id prefix. Without such a sink, a crash can still write the same counts twice.

Every one of those lines is a correctness decision, not a speed setting. Choose event time or you count late records in the wrong window. Choose exactly-once, with a sink that takes part, or a crash counts records twice. The watermark and the grace period add up: here a record about 70 seconds behind the newest one still counts, and a later one is dropped. Set them too tight and you drop stragglers. That many correctness settings is why a stream is the tool you reach for on purpose, not the default you turn on everywhere.

A two-column table headed from each company's own posts, checked 2026-09-30, titled three teams, each running both. Left, the slow part, batch; right, the fast part, stream. Uber, 2017: a restaurant's average meal preparation time over the last seven days, bulk-loaded from HDFS into Cassandra; against the same average over the last one hour, from Kafka through Samza streaming jobs into Cassandra. DoorDash, 2020: almost all of the features get updated every day; against DoorDash, 2021: orders confirmed by a store in the last 30 minutes, on Flink, stored in Redis. Netflix, 2013: offline computation, with fewer limits on the algorithm and the amount of data; against online computation, which responds quickly to events and uses the most recent data. Beneath: long windows and heavy work, batch; the last minutes or the last hour, a stream.

The pattern across all three is the one the lab measured. Long windows and heavy computation go to batch. The last minutes or the last hour go to a stream. Uber's post also says, in so many words, that training and serving must see the same data. None of the three posts describes streaming everything, and none describes serving a last-hour number from a nightly job. That discipline, applied input by input, is the entire lesson.

This is a real run in VS Code's terminal (python batch_stream_demo.py).

A real screenshot of VS Code's terminal after running python batch_stream_demo.py. It prints scikit-learn 1.9.1, 3,476 test hours; real rides an hour in the test: 248.8; then the error, MAE in rides an hour, one run each: calendar 67.8, nightly 56.4, stream 30.5, split 30.6, skew 173.9; the nightly count off from the fresh one, on average: fast 186.1 rides, slow 2.7 rides an hour; the batch job every 1, 2, 3, 6, 12 and 24 hours, with 3,516, 1,758, 1,172, 586, 293 and 146 runs and MAE 30.5, 33.8, 37.3, 42.2, 45.9 and 56.4; and one update per ride instead: 864,671.

When I ran it, it printed scikit-learn 1.9.1, 3,476 test hours and a real mean of 248.8 rides. The five designs came out at 67.8, 56.4, 30.5, 30.6 and 173.9, and the nightly counts were off by 186.1 and 2.7. The job every 1 to 24 hours scored 30.5, 33.8, 37.3, 42.2, 45.9 and 56.4. All of it matches the stored bv-demo.json, and the longest printed line was 49 characters. The report's demo mode checks every one of those lines against the file.

To try something I have not run, change both 168s in counts to 24, which turns the slow count into a one-day average. I cannot tell you what it prints, because I have not run it. The question to ask is whether the split design still costs almost nothing once the slow count moves faster.

Four brand cards headed the tools, with their logos, titled what ran where. scikit-learn: the models and the data download. pandas: the report's calendar and its rolling counts. NumPy: the demo's running total of rides. Python: the demo, the report and the box.

The split between the tools is deliberate. scikit-learn fitted every model and fetched the data. The demo computes its counts with a NumPy running total. The report computes them again with pandas on a real hourly calendar. So the check that the numbers come back uses a different method from the one that made them. Plain Python holds the playground, because it has to run in a browser.

table(1, 1)
table(24, 24)
table(1, 24)

mae. Fits a HistGradientBoostingRegressor with seed 0 on the training rows of one table, and scores it on the test rows of another. Passing the same table twice trains and serves alike; passing the stream's table and then the nightly one is the skew design.

The rest. How far the nightly counts are from fresh ones, the loop over 1, 2, 3, 6, 12 and 24 hours, and the count of rides in the test. With a file name on the command line, json.dump saves every number at full precision.