A batch job processes a bounded set of data, produces a result, and stops. Nightly revenue reports, daily feature tables for machine learning, monthly invoices and weekly search-index rebuilds are all batch jobs. Batch is the oldest data-processing model and still carries most of the world's analytical and financial workloads, because it is simple to reason about: the input is known, the output is complete, and a run can be repeated.
That last property is not automatic. Most batch pain, such as double-counted revenue, half-written tables and backfills that corrupt current data, comes from jobs that cannot be rerun safely. This article covers the architecture that makes reruns safe, independent of engine: how to define the unit of work, how to write and publish output, how jobs signal completeness to each other, and how to backfill, retry and validate. Engine internals for Hadoop are covered in the MapReduce overview, and the orchestrator that runs pipelines in workflow orchestration.
What makes a job a batch job
Three properties define batch. The input is bounded: the job knows when it has read everything. The output is complete for its scope: a daily table for 2026-09-30 either exists in full or not at all. And the job is optimized for throughput, not latency: it can take an hour as long as it finishes before the deadline, and it uses that freedom to sort, join and aggregate large volumes efficiently.
The deadline is usually expressed as data freshness, for example "yesterday's revenue table is available by 06:00 UTC". That is the service level objective for a batch system, and everything else in the design, including parallelism, retries and alerting, exists to meet it.
The unit of work: a logical interval, not "now"
The most important design decision is to parameterize every run by the interval it processes, called the logical date, data interval or partition. The run for 2026-09-30 reads the input partitions for 2026-09-30 and writes the output partition for 2026-09-30, regardless of when it actually runs. A job that instead reads "everything since the last run" or uses the wall clock gives different results when it runs late, twice, or during a backfill.
With a logical interval, the job becomes close to a pure function: output[d] = f(inputs[d], reference_data_as_of[d], code_version). Rerunning it for the same d with the same code produces the same output. That property makes retries, backfills, audits and debugging tractable. Keep raw inputs immutable and partitioned by the same interval, so the inputs to f do not change underneath a rerun. If a source can only be read as a current snapshot, such as an operational database table, extract it once into a dated raw partition and have every downstream job read that copy.
Idempotent output: overwrite the partition
A run must produce the same end state whether it runs once or five times. Appending rows breaks this: a retried run appends the same rows again, and revenue doubles. The standard fix is partition overwrite: each run replaces the whole output partition for its interval. When the output is not naturally partitioned by interval, for example a table of customer states, use a keyed merge (upsert) on a deterministic key instead, so repeated application converges to the same rows. The same reasoning appears for message consumers in idempotency.
import os, shutil, uuid
from pathlib import Path
def run(day: str, root: Path):
src = root / "raw" / "orders" / f"dt={day}"
final = root / "marts" / "revenue_daily" / f"dt={day}"
staging = root / "marts" / "revenue_daily" / f"_staging_{day}_{uuid.uuid4().hex[:8]}"
rows = aggregate(read_partition(src)) # pure transform of this day's input only
write_parquet(staging, rows) # never write straight to the final path
check_quality(staging, day) # raises, leaving final untouched
# Publish: swap the directory. A rename within one POSIX file system is atomic,
# so readers see the old partition or the new one, never a mix.
old = final.with_name(final.name + ".old")
if final.exists():
os.replace(final, old)
os.replace(staging, final)
shutil.rmtree(old, ignore_errors=True)
(final / "_SUCCESS").touch() # completeness marker for downstream jobsThe sketch has a small window between the two renames in which the partition is missing. Readers that tolerate a retry are fine; readers that are not need a pointer switch instead, such as a view or catalog entry that is updated to point at a new versioned directory in one operation.
Atomic publish on object stores
Object stores such as S3, Google Cloud Storage and Azure Blob Storage do not offer an atomic rename of a directory. A "rename" of a prefix is a copy and delete of every object, and a failure halfway leaves a mix of old and new files. Early Hadoop-era committers built on rename were slow and unsafe on object stores for exactly this reason.
The modern answer is to move atomicity into metadata. Table formats such as Apache Iceberg, Delta Lake and Apache Hudi write new data files to unique paths and then commit a new table snapshot by atomically updating a single pointer, a metadata file or a transaction log entry. Readers resolve the table through that pointer, so they see either the previous snapshot or the new one. An overwrite of a partition is a commit that removes the old files from the snapshot and adds the new ones. Old files are deleted later by a retention job, which also gives you time travel for recovery. If you use plain files on an object store, emulate this with a manifest: write data under a run-specific prefix, then write a small manifest listing the files, and have readers go only through manifests.
Completeness and dependencies
A downstream job must not start until its inputs are complete, and "the directory exists" is not the same as "complete". Use explicit signals. Hadoop's output committer writes an empty _SUCCESS file after a job commits, and many pipelines copy the convention. Table formats let you check for a committed snapshot or a partition entry. Orchestrators model the dependency directly, with a sensor that waits for the upstream marker or a task dependency within the same pipeline.
Late-arriving data complicates completeness. Events for 2026-09-30 may still arrive on October 2 from a mobile client that was offline. Decide on a policy per dataset: close the interval after a fixed lateness allowance (for example, run at 06:00 with a 6-hour allowance and accept that later events are dropped or counted in a correction), or rerun the last N intervals every day so late data is folded in. The rolling rerun is simple and robust precisely because runs are idempotent overwrites.
Parallelism inside a job
Within a run, engines such as Spark split input into tasks, shuffle data by key for joins and aggregations, and write one or more files per output task. Three sizing rules prevent most performance problems. Aim for output files in the range of about a hundred megabytes to around a gigabyte, because thousands of tiny files slow every later reader and the metadata service. Watch for skew: one key with a large share of the rows makes one task run for an hour while the others finish in minutes; salting the key or handling hot keys separately fixes it. And size the cluster for the deadline, not the average: a job that normally takes 40 minutes must still finish in time on the day after a holiday sale.
Backfills
A backfill runs the job for past intervals, because of a new job, a bug fix or a changed definition. With logical intervals and overwrite semantics, a backfill is just many runs with old dates, but it needs operational care:
- Limit concurrency. Running 365 days at once can overwhelm the cluster and any source systems the job reads. Cap parallel runs, and schedule large backfills away from the nightly window.
- Pin the code version. Record which code version produced each partition. A table half-built by old logic and half by new logic is worse than either.
- Use point-in-time reference data. Joining 2025 orders to today's product catalogue silently rewrites history. Keep dated snapshots or validity ranges for reference tables.
- Backfill into a shadow table first for large changes, compare it with the current table, and swap once the differences are understood.
Retries, bad records and failure handling
Failures happen at two levels. Task-level failures, such as a lost executor or a transient read error, are retried by the engine; keep the retry count small and make sure tasks write to attempt-specific paths so a failed attempt's partial files are never committed. Job-level failures are retried by the orchestrator with backoff; because the job is an idempotent overwrite, a retry needs no cleanup.
Bad records need a policy rather than a crash. A single malformed row in a billion should not fail the nightly run, but a schema change that breaks every row must. Route unparseable rows to a quarantine location with the error and source position, count them, and fail the job when the count exceeds a threshold such as 0.1 percent of input. This is the batch equivalent of a dead letter queue, and the quarantine should be reviewed, not just accumulate.
Quality gates before publish
The cheapest time to catch bad data is before it is published. Run checks against the staging output and fail the run, leaving the previous partition in place, when they fail. Useful checks are cheap and specific: row count within a band relative to the same weekday last week; no nulls in key columns; unique primary keys; totals that reconcile with an independent source, such as payment totals against the ledger; and no values outside known domains. A published partition that is a day late is an inconvenience. A published partition that is wrong spreads into dashboards, model training and invoices before anyone notices.
Worked example: a daily revenue mart
Suppose the orders source produces 60 million events a day, about 120 GB of compressed Parquet. The job filters to completed orders, joins to a currency table and a product snapshot for the day, and aggregates by country, product line and hour. The deadline is 06:00 UTC, and the upstream extract lands by 02:30.
Budget backwards. With 3.5 hours of window, leave an hour of margin for one full retry and a slow day, so the job should finish in under 90 minutes on a normal day and under 2.5 hours on a peak day at three times normal volume. If a test run on a small cluster processes about 2 GB per core-minute through the join and aggregation, 120 GB is about 60 core-minutes, a few minutes on a 32-core cluster. The peak day triples that. The real risk is not compute but skew: one country with half the orders makes a single task dominate, so pre-aggregate by hour inside each task before the shuffle. The output is small, a few hundred thousand rows, written as one or two files per partition. The quality gate compares total revenue with the payment ledger to within 0.5 percent and checks that every country seen yesterday appears today.
Batch or streaming?
| Concern | Batch | Streaming |
|---|---|---|
| Freshness | Hours to a day | Seconds to minutes |
| Correctness model | Complete interval, simple reruns | Watermarks, late data, state management |
| Cost per record | Lowest: large sequential reads, spot capacity | Higher: always-on resources |
| Reprocessing | Rerun the interval | Replay from a log, rebuild state |
| Good fit | Finance, reporting, training sets, bulk exports | Alerting, personalization, fraud scoring |
Many platforms run both: streaming for the fresh view, and a batch job that recomputes the authoritative numbers each night. Micro-batch sits between them, running the batch model on small intervals of a few minutes; it keeps rerun semantics while cutting latency, at the cost of many small output files that need compaction. Change streams feeding batch tables are covered in change data capture.
Operating batch pipelines
Monitor data, not only jobs. The primary signal is freshness per dataset: when did the latest complete partition land, relative to its deadline? Alert when a dataset is late, even if every job reports success, because an upstream job that silently skipped a day produces no errors. Track duration per run as a trend, since slowly growing inputs push jobs towards their deadlines months before they breach. Record lineage, meaning which input partitions and code version produced each output partition, so that after a bug you can list every affected partition and backfill exactly those.
Failure modes
| Failure | Cause | Fix |
|---|---|---|
| Doubled numbers after a retry | Append instead of overwrite | Partition overwrite or keyed merge |
| Readers see half a table | Writing directly to the final path | Stage, then atomic swap or snapshot commit |
| Downstream ran on partial input | Directory existence used as readiness | Success markers or committed snapshots, sensors |
| Backfill rewrote history wrongly | Current reference data joined to old facts | Dated reference snapshots, pinned code version |
| Nightly job misses its deadline | Skew, growth, small files | Salting, capacity for peak, compaction, duration trend alerts |
| Silent missing day | Upstream skipped, job succeeded on empty input | Row-count bands, freshness alerts per dataset |
What to do next
- Check that every batch job in your pipeline takes its logical interval as a parameter and never reads the wall clock for data selection.
- Find any job that appends output, and convert it to partition overwrite or a keyed merge.
- Make publishing atomic: staging plus rename on a POSIX file system, or a table format commit on object storage.
- Add a quality gate with row-count bands, key uniqueness and one reconciliation total to your most important dataset.
- Add a freshness alert per dataset based on its deadline, not on job status.
- Run a one-week backfill into a shadow table and compare it with production to prove your jobs are really idempotent.