Batch and stream processing
Data pipelines move data from where it is produced (application databases, event logs, third-party APIs) to where it is analysed or used. The two basic styles are batch and streaming.
ETL and ELT
- ETL (extract, transform, load) transforms data before loading it into the warehouse. This made sense when warehouse storage and compute were expensive.
- ELT (extract, load, transform) loads raw data first and transforms it inside the warehouse with SQL, often managed with a tool like dbt. Modern cloud warehouses make this cheap, and keeping raw data means you can rerun transformations when logic changes.
Batch processing
A batch job processes a bounded chunk of data on a schedule: “every night, compute yesterday’s revenue by country”.
- Simple to build, test and reason about.
- Easy to rerun: fix the bug, rerun the job for the affected days.
- Results are only as fresh as the schedule, typically hours old.
Most reporting, billing and machine-learning training data work well in batch.
Stream processing
A streaming system processes events continuously as they arrive, usually from a log such as Kafka, Kinesis or Pub/Sub: “flag a card within seconds of a suspicious purchase”.
- Low latency, from seconds to milliseconds.
- More moving parts: state, ordering, failures and replays all need handling.
- Harder to test and to correct after the fact.
Choose streaming when the value of the result decays quickly: fraud detection, live dashboards, alerting, real-time recommendations. Otherwise, batch is usually the better trade.
Event time versus processing time
In streaming, every event has two times:
- Event time: when it actually happened, for example when the user tapped “buy”.
- Processing time: when your pipeline saw it.
They differ because phones go offline, queues back up and services retry. If you aggregate by processing time, a delay makes Monday’s sales appear on Tuesday. Aggregate by event time.
Windows and late data
To compute “purchases per 5 minutes” over an endless stream, you group events into windows:
- Tumbling: fixed and non-overlapping (00:00–00:05, 00:05–00:10).
- Sliding: overlapping, for example a 5-minute window every minute.
- Session: groups activity separated by gaps of inactivity.
But when is a window complete? A watermark is the system’s estimate that “no events older than time T are still coming”. Events arriving after the watermark are late data: decide whether to drop them, update the earlier result, or route them to a correction job.
Delivery guarantees
- At most once: events may be lost, never duplicated.
- At least once: events are never lost, but may be processed twice after a retry. This is the most common guarantee.
- Exactly once: each event affects the result once. It is achievable within some systems, but end-to-end it usually depends on idempotent writes.
The practical rule: assume at least once and make every step idempotent. Upsert by a unique event id instead of blindly inserting, so processing an event twice has the same effect as processing it once.
Pipelines you can rerun
Whether batch or streaming, design pipelines so any run can be repeated safely:
- partition outputs by date and overwrite a partition instead of appending to it;
- keep raw input immutable, so it can be reprocessed;
- make transformations deterministic, with no “current time” hidden inside them;
- test with data quality checks: row counts, uniqueness of keys, null rates and freshness.
A pipeline that can’t be rerun turns every bug into a manual data-repair project.