ETL vs streaming
ETL vs Streaming
Understand how batch-oriented ETL pipelines and continuous streaming pipelines move, transform, validate, and deliver data, and learn how to choose the right processing model for a system.
Prerequisites
Before studying ETL and streaming, you should have a basic understanding of the following concepts:
Recommended Knowledge
- Relational databases, tables, rows, and SQL queries
- Files, object storage, and data warehouses
- Queues, event logs, producers, and consumers
- JSON and structured data formats
- Data partitioning and distributed processing
- Retries, idempotency, and failure handling
- Basic data aggregation and transformation
- Search indexes and ranking signals
What Is a Data Pipeline?
A data pipeline is an automated sequence of steps that moves data from one or more source systems to one or more destinations. During this journey, the pipeline may validate, clean, filter, enrich, join, aggregate, or restructure the data.
Sources may include operational databases, application events, APIs, log files, sensors, third-party systems, and object storage. Consumers may include dashboards, search indexes, recommendation engines, machine-learning models, reports, APIs, and alerting systems.
What Is ETL?
ETL stands for Extract, Transform, and Load. It is a data-integration process in which data is extracted from source systems, transformed into the required format, and then loaded into a destination.
Extract
Read data from one or more source systems.
Extraction may read database tables, API responses, CSV files, application logs, object-storage files, or exported business records. The pipeline should record which input was included so the run can be reproduced or audited.
Transform
Convert raw input into a clean and useful representation.
Transformations may remove invalid records, standardize units, normalize text, mask sensitive values, join datasets, derive columns, calculate aggregates, or apply business rules.
Load
Write the transformed output into a destination.
The destination may be a data warehouse, data lake, relational database, search index, reporting table, or machine-learning feature store.
Simple Analogy
ETL is like collecting products from several warehouses, inspecting and repackaging them at a processing center, and delivering the standardized products to a store.
Batch Processing in ETL
ETL commonly operates as a batch process. Instead of processing every record immediately, the system collects records during a defined interval and processes them together.
A batch may run every hour, every day, at the end of a business period, after files arrive, or on demand. Each run normally has a clear beginning, a bounded input, and a completion state.
Here, \(B_k\) is the batch for interval \(k\), \(e_i\) is an input event or record, and \(t_i\) is the time associated with that record. The batch contains the records that fall within the selected processing boundary.
Common ETL Use Cases
Where ETL Works Well
- Daily or hourly business reports
- Financial reconciliation
- Historical analytics
- Periodic search-index rebuilding
- Machine-learning training datasets
- Customer and product master-data synchronization
- Data migration between systems
- Compliance and audit exports
- Large-scale backfills
- Batch recomputation of ranking features
What Is Streaming?
Streaming is a processing model in which data is continuously consumed and processed as events become available. Instead of waiting for a large scheduled batch, the system continuously reacts to new information.
Streaming does not necessarily mean that the system processes exactly one event at a time. A streaming engine may group events into very small internal batches for efficiency while still presenting a continuous processing model.
Simple Analogy
Batch ETL is like filling a bucket and carrying it to its destination. Streaming is like using a pipe through which water continuously flows and is processed while moving.
Common Streaming Use Cases
Where Streaming Works Well
- Real-time fraud and abuse detection
- Live operational dashboards
- Application monitoring and alerting
- Clickstream and user-activity processing
- Real-time recommendations
- Search-index updates after content changes
- Inventory and availability updates
- Change data capture
- Location and sensor-event processing
- Continuous ranking-signal updates
ETL vs Streaming: Core Comparison
| Dimension | ETL or Batch Processing | Streaming Processing |
|---|---|---|
| Processing model | Processes a bounded collection of records | Continuously processes incoming events |
| Trigger | Schedule, file arrival, or manual execution | Event arrival or continuous log consumption |
| Latency | Usually determined by schedule and job duration | Designed for continuously updated results |
| Input boundary | Usually explicit and finite | Potentially unbounded |
| State management | Often scoped to one job or dataset | Maintained continuously across events and windows |
| Failure recovery | Retry a job, stage, or partition | Restart from a checkpoint or log offset |
| Late data | Can be handled through a later batch or backfill | Requires event-time and lateness policies |
| Operational complexity | Generally easier to schedule and reason about | Requires continuous monitoring and state management |
| Capacity usage | Resources may be allocated only during execution | Processing infrastructure normally remains active |
| Typical output | Periodic snapshot or aggregate | Continuously updated view or action |
| Backfills | Natural part of the batch model | Usually requires event replay or a separate batch path |
| Best fit | When delayed results are acceptable | When the value of data decreases quickly with delay |
ETL vs ELT
ETL should not be confused with ELT. Both patterns extract data, but they differ in where transformation takes place.
ETL
- Extract data from source systems
- Transform before loading into the destination
- Control the schema at the destination boundary
- Avoid loading unnecessary or prohibited fields
ELT
- Extract data from source systems
- Load raw or lightly processed data first
- Transform inside the target analytical platform
- Retain raw data for new transformations and reprocessing
| Factor | ETL | ELT |
|---|---|---|
| Order | Extract, transform, and load | Extract, load, and transform |
| Transformation location | External processing system | Destination warehouse or lakehouse |
| Raw-data retention | May not retain every original field | Usually preserves a raw layer |
| Destination requirement | Destination can receive prepared data | Destination needs transformation capability |
| Reprocessing | May require re-extracting source data | Can transform retained raw data again |
Processing Time and Event Time
Streaming systems must distinguish between the time an event occurred and the time the processing system received or processed it.
Event Time
Event time is when an event occurred in the source domain. For example, it may represent when an order was placed or when a user clicked a result.
Ingestion Time
Ingestion time is when the event entered the messaging or processing platform.
Processing Time
Processing time is when a worker actually processed the event. Queueing, retries, network problems, and service interruption can create a difference between event time and processing time.
Measuring only processing duration can hide time spent waiting in a producer, network, queue, or event log. End-to-end latency gives a more complete view of data freshness.
Windows in Stream Processing
A stream is unbounded, but many calculations require a bounded set of events. Windows divide the stream into manageable groups for aggregation.
Tumbling Window
Fixed-size, non-overlapping windows. For example, every event belongs to exactly one five-minute interval.
Sliding Window
Fixed-size windows that start at regular intervals and may overlap. An event may contribute to multiple windows.
Session Window
Groups events separated by less than a configured inactivity gap. It is useful for user sessions and bursts of activity.
Global Window
Places events into one logical window and relies on a custom trigger or processing rule to produce results.
SELECT
product_id,
window_start,
window_end,
COUNT(*) AS click_count
FROM TABLE(
TUMBLE(
TABLE product_clicks,
DESCRIPTOR(event_time),
INTERVAL '5' MINUTE
)
)
GROUP BY
product_id,
window_start,
window_end;
This conceptual streaming SQL query counts product clicks in non-overlapping five-minute event-time windows. Exact syntax depends on the selected stream-processing platform.
Late and Out-of-Order Events
Events do not always arrive in the same order in which they occurred. Mobile disconnection, network delay, retries, queue partitions, and producer failures can cause older events to arrive after newer events.
A stream processor therefore needs a policy that decides how long to wait before treating an event-time window as sufficiently complete.
The watermark \(W\) estimates how far event-time processing has progressed. Events older than the watermark may be treated as late, depending on the pipeline policy.
Wait Longer
- Includes more delayed events in the initial result
- Improves event-time completeness
- Increases result latency
- Requires state to be retained longer
Publish Earlier
- Produces results sooner
- May exclude delayed events initially
- May require corrections or result updates
- Can expose temporarily incomplete aggregates
Delivery Semantics
Because consumers and networks can fail, a pipeline must define how event redelivery affects its output.
| Semantic | Meaning | Design Concern |
|---|---|---|
| At-most-once | An event may be processed once or not processed | Failures may cause data loss |
| At-least-once | An event is retried but may be processed repeatedly | Consumers must handle duplicates safely |
| Effectively-once | Retries may occur, but visible effects are deduplicated | Requires idempotency, transactions, or deterministic output |
Idempotent Stream Processing
An idempotent consumer can receive the same event multiple times without producing multiple unintended effects. This is essential when the processing platform uses retries.
BEGIN;
INSERT INTO processed_events (
event_id,
processed_at
)
VALUES (
'event-98341',
CURRENT_TIMESTAMP
)
ON CONFLICT (event_id) DO NOTHING;
UPDATE product_metrics
SET click_count = click_count + 1
WHERE product_id = 501
AND EXISTS (
SELECT 1
FROM processed_events
WHERE event_id = 'event-98341'
)
AND NOT EXISTS (
SELECT 1
FROM applied_events
WHERE event_id = 'event-98341'
);
INSERT INTO applied_events (
event_id,
applied_at
)
VALUES (
'event-98341',
CURRENT_TIMESTAMP
)
ON CONFLICT (event_id) DO NOTHING;
COMMIT;
This example illustrates the need to connect deduplication and the business update within one safe transactional boundary. The exact implementation depends on the database and consistency model.
Checkpoints and Offsets
Streaming consumers need to remember their progress. An offset identifies a position in an ordered log partition, while a checkpoint records enough processing state to resume after a failure.
{
"consumerGroup": "search-ranking-consumers",
"checkpointId": "checkpoint-2841",
"partitions": {
"0": 193450,
"1": 201782,
"2": 188945
},
"stateVersion": "ranking-state-v7",
"status": "COMPLETED"
}
Batch ETL Example
The following JavaScript example demonstrates the structure of an incremental ETL job. It extracts records after a stored checkpoint, transforms them, writes to staging, validates the result, and then updates the checkpoint.
async function runOrderEtl() {
const checkpoint = await getCheckpoint("orders-etl");
const runId = createRunId();
try {
const orders = await extractOrders({
afterId: checkpoint.lastOrderId,
limit: 10000
});
const transformedOrders = orders.map(function (order) {
return {
orderId: order.id,
customerId: order.customer_id,
totalAmount: Number(order.total_amount),
orderDate: order.created_at.substring(0, 10),
status: normalizeStatus(order.status)
};
});
await writeToStaging(runId, transformedOrders);
await validateStagingData(runId);
await mergeIntoWarehouse(runId);
const lastOrder = orders[orders.length - 1];
if (lastOrder) {
await saveCheckpoint("orders-etl", {
lastOrderId: lastOrder.id,
completedRunId: runId
});
}
await markRunComplete(runId);
} catch (error) {
await markRunFailed(runId, error.message);
throw error;
}
}
In a production implementation, extraction should use a stable boundary. Depending only on an increasing identifier may be insufficient when old records can be updated or deleted.
Streaming Consumer Example
async function processProductEvent(event) {
validateEventSchema(event);
const alreadyProcessed = await eventStore.exists(event.eventId);
if (alreadyProcessed) {
return;
}
const normalizedEvent = {
eventId: event.eventId,
productId: event.productId,
eventType: event.eventType,
eventTime: new Date(event.eventTime),
sourceVersion: event.schemaVersion
};
await database.transaction(async function (transaction) {
await updateProductProjection(
transaction,
normalizedEvent
);
await eventStore.markProcessed(
transaction,
normalizedEvent.eventId
);
});
}
The consumer validates the event, checks for prior processing, updates the derived product view, and records the processed event within a controlled transactional operation.
ETL and Streaming in a Search System
A search platform may use both processing models because different parts of the index have different freshness and correctness requirements.
Batch ETL Path
Periodically read the authoritative product catalog, clean and enrich records, calculate expensive features, and build a new complete search-index version.
Streaming Path
Consume product, price, inventory, and availability events and update affected search documents continuously.
Reconciliation Path
Compare the search index with the authoritative source and correct missing, duplicated, or inconsistent documents.
Product Events Stream Processor Search Index
ETL and Streaming in Ranking
Ranking systems combine signals with very different update frequencies. Some features can be recomputed periodically, while others lose value quickly if delayed.
| Ranking Signal | Suitable Processing Model | Reason |
|---|---|---|
| Long-term popularity | Batch ETL | Usually calculated over a large historical window |
| Product quality score | Batch ETL | May require expensive joins and aggregations |
| Current inventory | Streaming | Search results should react quickly to availability changes |
| Recent click activity | Streaming or micro-batch | Recent behavior may affect ranking freshness |
| Price changes | Streaming | Stale values can create a poor user experience |
| Historical conversion rate | Batch plus incremental updates | Historical processing is large, but recent changes also matter |
Here, \(B\) represents batch-generated historical features, \(S\) represents continuously updated streaming features, and \(R\) represents request-time signals such as query relevance or user context.
Micro-Batching
Micro-batching groups records over short intervals and processes each group as a small batch. It sits between traditional scheduled batches and event-by-event processing.
Benefits
- Amortizes database and network overhead
- Improves write efficiency
- Can reuse batch-oriented transformation logic
- May provide sufficient freshness for many applications
Limitations
- Adds waiting time before processing begins
- Does not provide immediate event handling
- Still needs checkpoint and duplicate management
- Large micro-batches may create latency spikes
Hybrid Architecture
Many production systems combine batch and streaming rather than selecting only one. The streaming path supplies fresh updates, while the batch path supports historical processing, correction, reconciliation, and rebuilding.
Streaming Processing Fresh Serving View
Batch Processing Historical and Corrected View
Replay and Backfill
A pipeline may need to process historical data again after a bug fix, schema change, new business rule, or destination recovery. Batch systems usually describe this as a backfill. Streaming systems often describe it as replay.
| Concern | Batch Backfill | Stream Replay |
|---|---|---|
| Input selection | Date range, partition, file set, or primary-key range | Offsets, event-time interval, or retained log segment |
| Output safety | Replace partitions or merge idempotently | Deduplicate or rebuild the affected state |
| Resource impact | Can consume significant batch capacity | Can compete with live event processing |
| Versioning | Record job and transformation versions | Record consumer, schema, and state versions |
| Main risk | Overwriting correct output with incomplete data | Duplicating side effects or increasing consumer lag |
Common Failure Scenarios
Incomplete Batch Input
An ETL job starts before all expected files or source partitions arrive. Use manifests, data-readiness checks, and explicit cutoff boundaries.
Duplicate Event Processing
A streaming consumer retries an event after its output was already written. Use stable event identifiers, idempotent updates, and transactional processing where appropriate.
Consumer Lag
Events arrive faster than consumers can process them. Monitor lag, increase safe parallelism, optimize slow operations, and apply backpressure.
Schema Change
A producer changes its event or table structure before consumers are compatible. Use schema contracts, compatibility checks, and versioned migrations.
Incorrect Checkpoint
Progress is advanced before output is safely committed. A failure may then skip events. Coordinate the checkpoint with state and output commits.
Poison Event
One invalid event repeatedly fails and blocks a partition. Apply bounded retries, capture error context, and route unprocessable records to a quarantine or dead-letter flow.
Late Event Changes a Final Result
A delayed event arrives after a window result was published. Define whether to discard it, update the result, emit a correction, or rebuild the affected period.
Monitoring and Observability
Batch and streaming pipelines require different operational metrics because their processing lifecycles are different.
| Category | ETL Metrics | Streaming Metrics |
|---|---|---|
| Progress | Completed stages and partitions | Offsets and checkpoint progress |
| Latency | Job duration and publication delay | Event-time and processing-time lag |
| Throughput | Rows, files, or bytes per run | Events or bytes processed per interval |
| Failures | Failed jobs, tasks, and partitions | Consumer errors, retries, and failed events |
| Quality | Rejected records and validation failures | Invalid, duplicate, late, and out-of-order events |
| Freshness | Time since the last successful load | Difference between current time and latest processed event |
| Capacity | Worker utilization during the job | Consumer utilization and backlog growth |
Throughput and Capacity
A stable streaming system must process events at least as quickly as they arrive over a sustained period. Otherwise, the backlog continues to grow.
Here, \(\lambda_{\text{arrival}}\) is the arrival rate and \(\mu_{\text{processing}}\) is the sustainable processing rate. Capacity planning must also consider bursts, retries, maintenance, failed workers, and partition skew.
For a batch pipeline, the primary condition is whether the job can complete within its processing window.
Data Quality and Security
Fast data is not useful when it is incorrect or unsafe. Both ETL and streaming pipelines should enforce data contracts, access controls, encryption, auditing, and validation.
Important Controls
- Validate required fields, types, and acceptable ranges
- Version schemas and transformation logic
- Reject or quarantine malformed records
- Encrypt data in transit and at rest
- Restrict access to raw and sensitive fields
- Mask or tokenize protected information when required
- Maintain lineage from source to destination
- Audit administrative and publication actions
- Define retention and deletion policies
- Test recovery without exposing restricted data
Cost and Complexity Trade-offs
ETL Cost Profile
- Compute may run only during scheduled jobs
- Large scans can consume substantial resources
- Repeated full processing may become expensive
- Operations are often simpler for moderate freshness needs
Streaming Cost Profile
- Infrastructure usually operates continuously
- State and checkpoints require durable storage
- Low-latency systems need continuous monitoring
- Replay capacity must coexist with live traffic
How to Choose ETL or Streaming
Begin with the business freshness requirement. Do not choose streaming only because it appears more modern. Continuous processing introduces state, ordering, replay, duplicate, checkpoint, and operational concerns.
Decision Questions
Ask Before Designing
- How fresh must the output be?
- What happens if data is delayed?
- Is the input bounded or continuously generated?
- Can the source data be replayed?
- How will duplicate records be handled?
- Can events arrive late or out of order?
- How will schema changes be introduced?
- How will historical backfills be performed?
- What state must the processor maintain?
- How will incorrect output be corrected?
- What is the recovery objective?
- Does the value of low latency justify the complexity?
Example System Design
Consider an e-commerce search platform that needs searchable products, accurate inventory, updated prices, and useful ranking signals.
Proposed Architecture
- Product, order, click, price, and inventory services create authoritative records and events.
- A durable event log stores events for independent consumers and controlled replay.
- A streaming processor validates events and updates price, availability, recent popularity, and operational views.
- Raw events are also retained in durable storage for auditing and historical processing.
- A scheduled ETL pipeline joins catalog information with historical interactions and computes expensive ranking features.
- The batch pipeline builds a validated version of the complete search index.
- The serving layer combines textual relevance, batch features, streaming features, and request-time context.
- A reconciliation job compares derived systems with authoritative records and repairs inconsistencies.
Raw Storage ETL Historical Features
System Design Interview Discussion
In a system design interview, do not select ETL or streaming without explaining the requirements and operational trade-offs.
| Interview Area | What to Explain |
|---|---|
| Freshness | The maximum acceptable delay for each output |
| Volume | Expected input rate, size, bursts, and growth |
| Ordering | Whether global or per-entity order is required |
| Delivery | How duplicates and retries affect the destination |
| State | How aggregations, joins, windows, and checkpoints are stored |
| Late data | How delayed events affect published results |
| Failure recovery | How the system resumes, replays, or reruns work |
| Backfill | How historical data is reprocessed safely |
| Schema evolution | How producers and consumers remain compatible |
| Observability | How lag, freshness, failures, and quality are monitored |
| Cost | Why continuous processing is or is not justified |
Common Design Mistakes
Weak Design
- Choosing streaming without a freshness requirement
- Assuming all events arrive once and in order
- Committing offsets before durable output
- Using processing time when event time is required
- Ignoring historical replay and backfills
- Allowing incompatible schema changes
- Monitoring only infrastructure health
- Treating a completed job as proof of correct data
Strong Design
- Starts with measurable freshness requirements
- Defines ordering and duplicate behavior
- Uses checkpoints and idempotent output
- Defines late-data and window policies
- Supports replay, backfill, and reconciliation
- Versions schemas and transformations
- Monitors end-to-end freshness and correctness
- Uses batch and streaming where each provides value
Production Readiness Checklist
Pipeline Checklist
- Define the authoritative source for every dataset
- Document freshness and latency requirements
- Use stable event or record identifiers
- Version schemas and transformation logic
- Validate data before publishing it
- Design idempotent writes and consumers
- Coordinate checkpoints with durable output
- Define event-time and late-data policies
- Monitor consumer lag and batch completion time
- Provide quarantine handling for invalid data
- Test retries, replay, backfill, and recovery
- Plan for hot and uneven partitions
- Protect sensitive information throughout the pipeline
- Maintain data lineage and audit information
- Run reconciliation against authoritative data
- Create operational runbooks and ownership boundaries
Knowledge Check
What is the primary difference between ETL and streaming?
ETL commonly processes a bounded dataset during a scheduled or triggered run, while streaming continuously processes incoming events.
Is ETL the same as batch processing?
No. ETL describes extracting, transforming, and loading data. Batch describes a processing model. ETL is commonly implemented in batches, but the concepts are not identical.
Why are checkpoints important in streaming?
They record processing progress and state so a consumer can recover without starting from the beginning or silently skipping uncommitted work.
Why can events arrive out of order?
Network delay, producer retries, offline devices, partitions, and service failures can cause processing order to differ from event occurrence order.
Why do systems combine batch and streaming?
Streaming supplies fresh updates, while batch processing supports historical computation, backfills, reconciliation, and correction.
Summary
ETL extracts data from source systems, transforms it into the required structure, and loads it into a destination. It is commonly executed as a bounded batch and works well for reporting, historical analysis, migrations, backfills, and periodic recomputation.
Streaming continuously processes events and is appropriate when systems must react to changes with low delay. It requires careful handling of partitions, ordering, duplicates, checkpoints, state, event time, windows, late data, replay, and continuous operations.
ETL versus ELT determines where transformation occurs, while batch versus streaming determines how data is processed over time. These decisions should be evaluated separately.
Many search and ranking platforms use a hybrid architecture. Streaming maintains fresh prices, inventory, clicks, and recent activity, while batch pipelines compute historical features, rebuild indexes, perform backfills, and repair drift.
Key Takeaway
Choose the simplest processing model that satisfies the required freshness. Use batch ETL for bounded and delay-tolerant workloads, use streaming for continuously changing time-sensitive data, and combine them when the system needs both freshness and reliable historical recomputation.