0%
Data-Intensive Applications
Foundations of Data Systems
Distributed Data
Encoding and Evolution
Stream Processing
Data Quality and Governance
Operational Patterns
Batch Processing Patterns
Batch processing begins with moving data from where it is created to where it can be analyzed. Unlike stream processing (which handles events one at a time as they arrive), batch processing accumulates data over a time window (hourly, daily, weekly) and processes it all at once. This makes batch processing inherently higher latency (you wait for the batch window to close) but also simpler and more efficient: you amortize the overhead of job startup, resource allocation, and I/O across millions of records instead of paying it per-event.
ETL (Extract, Transform, Load) is the traditional pattern: pull data from source systems, reshape it into an analytical schema, and write it to a data warehouse. ELT (Extract, Load, Transform) flips the last two steps: dump raw data into the warehouse first, then transform it using the warehouse's own compute engine. The distinction matters because it determines where your compute budget goes and how flexible your pipeline is when requirements change.
ETL vs ELT
In ETL, transformation happens in a dedicated processing layer (Spark, custom scripts, Informatica) before data reaches the warehouse. This works well when storage is expensive: you only store the cleaned, modeled data. The transformation stage typically converts raw, normalized source schemas (third normal form from OLTP databases) into denormalized analytical schemas (star or snowflake schemas optimized for aggregation queries). This includes joining dimension tables to fact tables, computing derived columns, deduplicating records, and applying business logic (currency conversion, timezone normalization, data type casting).
The downside is rigidity. If an analyst needs a column you dropped during transformation, you must modify the pipeline, reprocess the data, and wait for the next batch run. In organizations with hundreds of pipelines, this change-request cycle can take weeks: the request goes to the data engineering team, gets prioritized against other work, requires code review and testing, and finally runs in the next scheduled batch window.
ELT became dominant as cloud warehouses (BigQuery, Snowflake, Redshift) made storage cheap and compute elastic. You load raw data first, preserving every column and every row. Transformations run inside the warehouse using SQL (often orchestrated by dbt). When requirements change, you write a new SQL model against the raw data that already exists. No pipeline modification, no reprocessing from source.
The practical difference shows up during debugging. In ETL, when an analyst reports wrong numbers, you must trace backward through the transformation code to find where the logic went wrong, then reprocess from source. In ELT, the raw data is already in the warehouse. You can write ad hoc SQL queries against it to diagnose the issue immediately, then fix the transformation model and rebuild the downstream table in place. The debugging cycle shrinks from hours (wait for a reprocessing run) to minutes (run a corrected SQL query).
In interviews, frame the ETL vs ELT choice as a flexibility-cost tradeoff. ETL saves storage by discarding raw data early. ELT saves engineering time by deferring transformation decisions. Modern data stacks overwhelmingly favor ELT because storage is cheap, but ETL remains necessary when data must be cleaned for compliance reasons before it touches the warehouse (PII redaction, GDPR filtering).
Orchestration and DAG Scheduling
A production batch pipeline is not a single script. It is a directed acyclic graph (DAG) of dependent tasks: extract from source A, extract from source B, join A and B, compute aggregates, load into warehouse, update dashboards. Tools like Apache Airflow, Dagster, and Prefect define these dependencies and execute tasks in the correct order.
The key property of a well-designed DAG is that each task is independently retriable. If the "join A and B" task fails, you rerun just that task, not the entire pipeline from the beginning. This requires that each task reads from and writes to well-defined locations (staging tables, object storage paths) rather than passing data through in-memory variables. The intermediate outputs serve as both checkpoints for retry and audit trails for debugging.
DAG scheduling also handles backfills. When you deploy a new pipeline or fix a bug in an existing one, you need to reprocess historical data. A good orchestrator lets you trigger the DAG for a specific date range (say, the last 90 days) and processes each date partition independently. Without this capability, backfills become manual, error-prone scripts that someone runs from their laptop.
Failure handling in DAGs requires careful thought about retry semantics. If the "load into warehouse" task fails, retrying it is only safe if the load is idempotent (overwrites the partition rather than appending). If the "extract from source A" task fails, retrying depends on whether the source supports re-reading the same data range (it does if you query by timestamp; it does not if you consume from a queue that auto-acknowledges). Each task in the DAG should document its retry behavior: idempotent and safe to retry, or non-idempotent and requires manual intervention.
Data Quality Gates
Production pipelines need automated quality checks between extraction and loading. A data quality gate validates each batch before it progresses to the next stage. Common checks include row count validation (today's extract should be within 20% of yesterday's), null rate monitoring (critical columns should have less than 0.1% nulls), schema validation (column types and names match expectations), and freshness checks (the most recent timestamp in the batch should be within the expected window).
When a gate fails, the pipeline halts and alerts the on-call engineer. The previous valid data remains in the warehouse. This prevents a common failure mode where a source system bug sends empty or malformed data, the pipeline loads it without checking, and dashboards suddenly show zero revenue or missing users.
A simple quality gate in SQL:
The cost of a quality gate is a few seconds of compute per batch. The cost of loading bad data is hours of investigation and potential business decisions made on wrong numbers.
Change Data Capture (CDC)
The extraction phase has its own design decision: do you query the source database with SELECT statements on a schedule, or do you capture changes as they happen? CDC-based ingestion reads the database's transaction log (MySQL binlog, PostgreSQL WAL, MongoDB oplog) to capture every INSERT, UPDATE, and DELETE in real time.
CDC has three advantages over scheduled queries.
First, it captures deletes. A scheduled SELECT only sees rows that currently exist; it misses rows deleted between runs. If a user deletes their account, the SELECT-based pipeline never notices: the row simply vanishes from future extracts. CDC captures the DELETE event explicitly, allowing the downstream system to mark the record as deleted.
Second, it captures intermediate states. If a row is updated three times between batch runs, a SELECT sees only the final state, while CDC captures all three changes. For audit trails, compliance logging, and analytics that depend on state transitions (how many times did the order status change before delivery?), intermediate states are essential data that point-in-time queries discard.
Third, it reduces load on the source database. Reading the transaction log is a sequential scan of an append-only file, far cheaper than running analytical queries against production tables. A scheduled SELECT that scans a 100-million-row table puts significant load on the database during peak hours. CDC reads the log stream with negligible impact because the database already writes the log for its own durability guarantees.
The tradeoff is operational complexity. CDC requires managing log positions (offsets), handling schema changes in the source database, and running connectors (Debezium, Maxwell) that must stay in sync with the source. If the connector falls behind and the database purges old log segments, you lose events and must fall back to a full snapshot.
Schema evolution in CDC is particularly tricky. When the source database adds a column, renames a field, or changes a data type, the CDC connector must handle the transition without crashing or losing data. Debezium, for example, captures schema changes as separate events in the stream. Downstream consumers must process these schema change events and adapt their deserialization logic accordingly. Teams that adopt CDC without planning for schema evolution discover the problem the first time a developer runs an ALTER TABLE on a monitored table and the pipeline breaks.
Another CDC operational concern is initial bootstrapping. When you first set up CDC for an existing table with 500 million rows, you cannot simply start reading from the current log position because you would miss all the existing data. The initial load requires a consistent snapshot: read the entire table at a point in time, then start streaming changes from that exact point forward. Debezium handles this automatically (the "initial snapshot" mode), but the snapshot itself can take hours for large tables and places significant read load on the source database. Planning the initial snapshot during a low-traffic window is essential.
Incremental Processing vs Full Reprocess
Every batch pipeline faces a fundamental scheduling question: should each run process all data from scratch (full reprocess) or only the data that changed since the last run (incremental)?
Full reprocessing is simple and correct by construction. There is no state to maintain, no edge cases around late-arriving data, and no risk of accumulated drift between incremental results and ground truth. The cost is compute: processing 1 TB of data every hour when only 10 GB changed is wasteful. But "wasteful" is relative. If your batch window is 4 hours and full reprocessing finishes in 2, the extra compute cost may be worth the operational simplicity. You never debug checkpoint issues because there are no checkpoints.
Incremental processing is efficient but complex. You maintain a checkpoint (a timestamp, an offset, a watermark) that marks where the last run ended. The next run processes only records newer than the checkpoint. This works well for append-only data (event logs, sensor readings). It breaks down when source records are updated in place: if you only process new records, you miss updates to old ones.
Checkpoints themselves introduce failure modes. If the checkpoint is stored in a different system than the output (e.g., checkpoint in ZooKeeper, output in S3), a crash between writing the output and updating the checkpoint causes the next run to reprocess data that was already written, creating duplicates. If the crash happens after updating the checkpoint but before writing the output, the next run skips that data entirely, creating gaps. The solution is to make the checkpoint and the output atomic: store the checkpoint alongside the output data (e.g., as a metadata file in the same S3 prefix) so both are committed or neither is.
The hybrid approach is common in practice: run incremental processing for daily batches and a full reprocess weekly or monthly to correct any drift. The full reprocess serves as a reconciliation step that catches anything the incremental runs missed.
A refinement of the hybrid approach is "tiered reprocessing." Daily runs process only new data (fast, cheap). Weekly runs reprocess the last 7 days to catch late-arriving data and corrections. Monthly runs reprocess the full dataset to eliminate any accumulated drift. Each tier has a different cost-accuracy tradeoff: daily is cheap but potentially incomplete, weekly catches most issues, and monthly guarantees correctness. This three-tier pattern is common at companies like Netflix and Airbnb, where the daily tier feeds real-time dashboards, the weekly tier feeds business reports, and the monthly tier feeds financial reconciliation.
The tiered approach also maps to different SLAs. The daily tier can tolerate a few percent of missing data because dashboards show trends, not exact numbers. The weekly tier needs accuracy within 0.5% because business reports drive staffing and inventory decisions. The monthly tier must be exact because financial reconciliation feeds audited financial statements where even a $1 discrepancy must be explained. Each tier's processing budget is proportional to its accuracy requirement.
Idempotent Processing
Network failures, pod restarts, and cloud spot instance preemptions mean your batch job will fail partway through at some point. The question is not whether it will fail, but whether you can safely rerun it.
An idempotent pipeline produces the same output regardless of how many times you run it with the same input. Formally, a pipeline is idempotent when
for every input . This sounds simple, but most pipelines violate idempotency in subtle ways.
The batch framing here is deliberately narrow: a whole job rerun over a whole input. The record-level version, where individual messages carry idempotency keys checked against a durable store, is a different mechanism with different failure modes and is developed in building reliable data pipelines. Batch idempotency is usually cheaper when you can get it, because overwriting a partition needs no key store, no TTL and no lookup on the hot path.
A pipeline that appends rows to a table without checking for duplicates doubles the data on rerun. A pipeline that assigns sequential IDs generates different IDs on rerun. A pipeline that calls an external API with side effects (sending emails, charging credit cards) cannot be idempotent without additional safeguards. Even pipelines that appear idempotent may not be: a pipeline that reads from a source table and writes to a target table is idempotent only if the source data has not changed between the original run and the rerun. If the source is a live table receiving new writes, the rerun processes a different input and produces a different output.
The standard technique is "write to staging, then atomically swap." Each pipeline run writes its output to a temporary location (a staging table, a temporary directory). Once the run completes successfully, you atomically replace the final output with the staging data. If the run fails, the staging data is discarded. If you rerun, the new staging data overwrites the old staging data before the swap. The final output is always the result of exactly one complete run.
In Spark, partition-based overwrite is the idiomatic implementation:
This overwrites only the partitions present in the current DataFrame. If today's run produces data for date=2026-05-31, only that partition is replaced. All other date partitions remain untouched. Rerunning for the same date overwrites the same partition with identical data.
In practice, partition-based overwrite is the most common implementation. If your output table is partitioned by date, each pipeline run overwrites exactly one partition (today's date). Rerunning the pipeline for the same date replaces the same partition. Other partitions are untouched. This gives you idempotency at the partition level without touching historical data.
Another technique is deterministic output paths. Instead of writing to output/results.parquet, write to output/results_2026-05-31_run-abc123.parquet where the path includes the logical date and a deterministic run identifier. A rerun produces the same path and overwrites the same file. This makes every run's output traceable and rollback trivial: point consumers back to the previous run's file.
The worst pattern is "append to a shared table and rely on deduplication later." This pushes the correctness burden downstream. Every consumer must know how to deduplicate, and if they forget (or their deduplication logic has a different key than yours), they get wrong results. Idempotency should be solved once, at the source, not repeatedly at every consumer.
Timestamps and non-deterministic functions are another idempotency trap. A pipeline that sets processed_at = NOW() in its output produces different results on every run. A pipeline that generates UUIDs for new records creates different IDs on rerun. Replace NOW() with the batch's logical timestamp (the scheduled run time, not the wall clock time). Replace random UUIDs with deterministic identifiers derived from the input data (a hash of the natural key).
External dependencies also break idempotency in subtle ways. A pipeline that calls a currency exchange rate API gets different rates on rerun because exchange rates change. A pipeline that reads a configuration file gets different behavior if the config changed between runs. The fix is to snapshot all external inputs at the start of each run: fetch the exchange rates once, store them alongside the batch data, and use the snapshot for the entire run. On rerun, use the same snapshot. This makes the pipeline a pure function of its inputs, which is the definition of idempotency.
Mid-level engineers build pipelines that work. Senior engineers build pipelines that fail gracefully and can be rerun without manual intervention. Staff engineers design the framework that makes idempotent-by-default the easy path for every pipeline in the organization, often through partition-based overwrite semantics and deterministic output paths.