Overview
Batch closes a time slice after it ends. Stream processes events as they land. The brief's SLA and lateness pick the clock.
On this page7 sections
The decision
Batch processing closes a slice of business time, then runs a job on that slice. Stream processing reads events from a queue as they land and uses watermarks to close windows.
Orchestration covers schedulers, idempotent loads, and backfills. Streaming covers topics, partitions, keys, and at-least-once delivery. Neither engine runs in this tab. You choose which practiced pattern the brief needs.
A batch system asks: what happened yesterday? A stream system asks: what is happening now? The word 'batch' means you collect a group of events, close the group, then process them together. The word 'stream' means each event is processed as it arrives (or within seconds).
What is at stake
The wrong clock wastes infrastructure and misses SLAs. A nightly file with a dashboard-due-at-08:00 requirement fits batch. A live counter with second-level freshness fits stream.
Many systems need both clocks for different consumers. Naming that split early prevents one path from silently serving two incompatible requirements.
Imagine your food delivery startup has two requests. Finance wants yesterday's gross revenue by 08:00 for bank reconciliation. The ops team wants a map of active deliveries updating every five seconds. A single pipeline cannot serve both: batch is too slow for the live map, streaming is over-engineered for a nightly report. You need two paths with a shared bronze layer.
Option A vs Option B
Batch closes a slice, then works. Stream works as events land, then closes a window with a watermark.
Circle one column per row. Mixed circles mean two paths, not one magic box.
| If this is true | Lean batch | Lean stream |
|---|---|---|
| Source is a nightly or hourly file | Yes: object prefix + DAG | Only if you must convert files to a topic |
| SLA is dashboard at 08:00 | Night job with a buffer | Overkill unless a second consumer needs seconds |
| SLA is click on a live counter | Will miss | Queue + watermark |
| Late data is next file restates | Backfill + idempotent load | Possible, heavier state |
| Late data is event hours after event_time | Next batch picks it up | Watermark + allowed lateness |
| You need per-key order | Sort in the job | Partition by key |
Scheduled shift versus always-on conveyor
A scheduler opens the dock at 02:00, processes yesterday's data, and clocks out. A stream consumer keeps moving while events arrive, and a watermark marks when you stop waiting for late events.
Orchestration why-a-scheduler named clock, retry, dependency, and record. idempotent-load makes a second run a skip instead of a second insert. backfills-and-catchup reruns a date after a fix.
Extract, land bronze, transform silver, test, then gold. Cloud managed-orchestration hosts this graph in production.
Streaming why-a-queue explains why services do not all write the warehouse directly. topics-partitions-keys sizes parallelism. at-least-once is the delivery promise you should assume: duplicates will happen.
Producers, topic, consumer group, sink. Watermarks close windows in Spark streaming.
Trace the clock
Read the brief. Predict batch versus stream before you watch two consumers get two tables.
A worked comparison
Finance needs yesterday's GMV by 08:00. The ops team needs active deliveries on a map every five seconds. The brief from the previous lesson forces two paths, not one magic box.
Circle one column per row. Mixed circles mean two paths.
| Consumer | SLA | Lean batch? | Lean stream? |
|---|---|---|---|
| Finance close | Yesterday by 08:00 | Yes | No |
| Live delivery map | 5 seconds | No | Yes |
| Vendor hourly file | Next batch run | Yes | Only if converting to topic |
# Conceptual batch pipeline skeleton
def run_batch_pipeline(execution_date: str):
"""Process one slice of time, then stop."""
raw_files = list_prefix(f"gs://bronze/events/dt={execution_date}/")
df = read_parquet(raw_files)
silver = transform(df, deduplicate_on="event_id")
silver.write_merge(target="silver.events", key="event_id")
gold = aggregate(silver, group_by=["dt", "country"], metric="sum(gmv)")
gold.write_overwrite(target="gold.fct_daily_gmv", partition="dt")
# Scheduler calls this at 02:00 with execution_date = yesterday# Conceptual stream consumer skeleton
def stream_consumer():
"""Continuously read from the queue and emit results."""
for message in consumer.poll(topic="orders"):
event = parse(message.value)
window = get_window(event.event_time, size="5min")
window.add(event)
if watermark_passed(window):
result = aggregate(window, metric="count")
emit_to_dashboard(result)
# Consumer runs forever; watermark decides when a window is "done"Hybrids land files to bronze on a schedule and also publish purchases to a topic for a live counter. Two consumers, two clocks, one bronze of record if you can afford it. Finance reads the night job; the growth widget reads the stream; both keys reconcile to order_id.
- Batch default: file or dump, SLA in hours, Orchestration graph.
- Stream default: live producers, SLA in seconds or minutes, queue vocabulary.
- Hybrid: two serving paths, one grain, two SLAs written down.
Stream is not faster batch
A stream that micro-batches every hour is a night job with more moving parts. Add a queue only when a named consumer pays for the freshness.
Do not add a queue because it sounds modern
If every source is a file dump and the SLA is hours, adding a queue means converting files to messages for no consumer benefit. Let the source shape guide the architecture.
PySpark structured streaming
You practiced watermarks and windows in the PySpark track. The stream column of this decision table is where those skills deploy in production.
Common beginner questions
How do I know if I need streaming or batch?
Look at your brief. If the freshness SLA is in hours and the source is a file, batch wins. If the SLA is in seconds and the source is a live service, stream wins. If you have both, you likely need a hybrid with two paths.
Can I start with batch and add streaming later?
Yes. Many teams launch with a nightly batch, then add a stream path for the one consumer that truly needs seconds. Starting batch-only is cheaper and simpler. Add a queue when a named consumer pays for the freshness.
What happens if the system gets 10x more data?
A well-partitioned batch job scales by adding Spark workers. A stream scales by adding partitions and consumer instances. Both scale, but you need the partition key design from the start.
What comes next
The next lesson sizes the ingest layer for 500 million events per day. You will use the batch-or-stream decision you made here to pick the right front door: object prefix or message queue.
Practice
Given a brief with SLA 'gold.fct_purchases by 08:00' and source 'hourly vendor gzip files', circle batch or stream for each row in the decision table. Write one sentence defending your choice for the finance consumer.
No editor on this lesson. The diagrams are the work. Mark it read when you can explain the design out loud.