batch recomputation
Batch Recomputation
Learn how large-scale systems periodically rebuild search indexes, rankings, recommendations, analytics, and other derived datasets from authoritative source data.
Prerequisites
Before studying batch recomputation, you should have a basic understanding of the following concepts:
Recommended Knowledge
- Databases, tables, indexes, joins, and aggregation
- ETL pipelines and data transformation
- Object storage and distributed file systems
- Queues, logs, and stream-processing fundamentals
- Search indexing and ranking signals
- Partitioning and parallel processing
- Idempotency, retries, and checkpointing
- Basic SQL syntax
What Is Batch Recomputation?
Batch recomputation is the process of rebuilding a derived dataset by reading a bounded collection of source data and recalculating some or all of the output. The computation usually runs on a schedule, after an important data change, or when an operator initiates a backfill or repair.
A derived dataset is information calculated from other data. It may include a search index, product ranking, recommendation list, daily report, machine-learning feature table, autocomplete dictionary, or materialized aggregate.
The function \(f\) represents the transformation logic. During a full recomputation, the system applies this logic to the complete relevant source dataset and produces a new version of the derived state.
Simple Analogy
Imagine calculating the final results of an examination. Instead of adjusting the published result after every mark changes, the school periodically reads all current marks, recalculates every total, and publishes a new result sheet. That complete calculation is similar to batch recomputation.
Authoritative Data and Derived Data
A well-designed recomputation pipeline clearly separates authoritative data from derived data.
Authoritative Data
- The accepted source of truth
- Preserves the original business facts or events
- Must be durable and recoverable
- May include database records, event logs, or raw files
Derived Data
- Generated from authoritative information
- Optimized for search, analytics, ranking, or serving
- Can usually be rebuilt if it becomes incorrect
- May be stale between recomputation cycles
Why Batch Recomputation Matters in Search and Ranking
Search and ranking systems frequently depend on information collected from many sources. Examples include document content, product metadata, click activity, popularity, freshness, quality, availability, and business rules.
Updating every output immediately after every source change may be too expensive or unnecessarily complex. A batch job can periodically combine the latest source data and build a consistent output snapshot.
| Use Case | Source Data | Recomputed Output |
|---|---|---|
| Search | Documents, metadata, permissions, and categories | Search index and document fields |
| Ranking | Clicks, engagement, quality, popularity, and freshness | Scores or ranked candidate lists |
| Autocomplete | Search queries, entities, frequency, and vocabulary | Prefix dictionary with suggestion scores |
| Recommendations | User activity, item metadata, and interaction history | Recommended items for users or segments |
| Analytics | Transactions, events, and operational records | Aggregated reports and dashboards |
How a Batch Recomputation Pipeline Works
A production batch-recomputation workflow usually contains multiple controlled stages. Each stage should be observable, restartable, and protected from incomplete output.
Capture a Stable Input Boundary
Decide exactly which source records belong to the run.
The job may process a database snapshot, all files available before a cutoff time, or all events up to a recorded log offset. A stable boundary prevents the job from reading an unpredictable mixture of old and new data.
Read and Validate Source Data
Load authoritative data and verify that required inputs exist.
Validation may check schemas, required partitions, file manifests, record counts, duplicate identifiers, null values, and expected value ranges.
Transform the Data
Clean, join, enrich, aggregate, tokenize, or score records.
Search pipelines may tokenize documents and build postings. Ranking pipelines may combine popularity, relevance, quality, and freshness signals into a final score.
Write to a New Output Version
Avoid modifying the currently served dataset in place.
The pipeline writes into a staging table, temporary index, versioned directory, or newly created object-storage prefix. Users continue reading the previous valid version while the new one is built.
Validate the New Version
Confirm correctness before making the output visible.
Validation can include document counts, query tests, ranking quality checks, duplicate detection, schema checks, checksum comparison, and comparison against the previous version.
Publish Atomically
Switch users from the previous output to the new output.
Publication may update a search alias, metadata pointer, database view, catalog entry, or routing configuration. The switch should appear atomic so consumers never receive a partially built dataset.
Retain and Clean Up
Keep enough history for rollback, then remove expired versions.
Retaining the previous successful output allows a fast rollback if the newly published version produces incorrect results.
Full Recomputation
Full recomputation processes the entire relevant source dataset and rebuilds the complete output. It is conceptually simple because the final result depends on the current source data rather than a complicated history of updates.
Advantages
- Simple correctness model
- Repairs historical processing errors
- Supports major transformation-logic changes
- Can rebuild state after corruption
- Reduces dependence on old intermediate state
Limitations
- Reads and processes unchanged data again
- May require substantial computing capacity
- Can take a long time for very large datasets
- May increase storage and network costs
- Produces stale output while the new version is building
Incremental Recomputation
Incremental recomputation processes only new or changed input and applies the resulting differences to existing derived state. It generally reduces work, but it introduces additional state and correctness challenges.
Here, \(O_t\) is the new output, \(O_{t-1}\) is the previous output, and \(\Delta(I_t)\) represents the changes calculated from new, modified, or deleted input records.
The formula uses addition conceptually. In practice, applying a delta may involve inserts, updates, deletions, score adjustments, or replacement of affected partitions.
| Factor | Full Recomputation | Incremental Recomputation |
|---|---|---|
| Input processed | Complete relevant dataset | New or changed data |
| Implementation | Usually simpler | Usually more complex |
| Processing cost | Higher as the dataset grows | Lower when the change set is small |
| Recovery from old errors | Naturally rebuilds the output | May preserve earlier incorrect state |
| Freshness | Depends on the full-run schedule | Can be refreshed more frequently |
| State management | Limited dependency on previous output | Requires checkpoints, offsets, or change tracking |
| Best fit | Moderate datasets or periodic correction | Large datasets with relatively small changes |
Choosing Between Full and Incremental Processing
The decision should be based on freshness requirements, dataset size, update frequency, processing cost, operational complexity, and the cost of an incorrect result.
Example: Recomputing Product Ranking
Consider a product-search system that recalculates a ranking score from textual relevance, popularity, freshness, and product quality.
In this example:
- \(R\) represents normalized textual relevance.
- \(P\) represents normalized popularity.
- \(F\) represents normalized freshness.
- \(Q\) represents normalized product quality.
The exact weights are business decisions and should be evaluated using relevance metrics, experiments, and user feedback rather than copied blindly.
CREATE TABLE product_ranking_build (
build_id VARCHAR(50) NOT NULL,
product_id BIGINT NOT NULL,
relevance_score DECIMAL(10, 6) NOT NULL,
popularity_score DECIMAL(10, 6) NOT NULL,
freshness_score DECIMAL(10, 6) NOT NULL,
quality_score DECIMAL(10, 6) NOT NULL,
final_score DECIMAL(10, 6) NOT NULL,
computed_at TIMESTAMP NOT NULL,
PRIMARY KEY (build_id, product_id)
);
INSERT INTO product_ranking_build (
build_id,
product_id,
relevance_score,
popularity_score,
freshness_score,
quality_score,
final_score,
computed_at
)
SELECT
'ranking-2026-09-24',
p.product_id,
p.relevance_score,
p.popularity_score,
p.freshness_score,
p.quality_score,
(0.45 * p.relevance_score) +
(0.25 * p.popularity_score) +
(0.20 * p.freshness_score) +
(0.10 * p.quality_score),
CURRENT_TIMESTAMP
FROM product_features p
WHERE p.is_active = TRUE;
The build identifier keeps one recomputation separate from another. The new result can be validated before the serving layer starts using it.
Example Pipeline Pseudocode
async function runBatchRecomputation(cutoffTime) {
const buildId = createBuildId(cutoffTime);
await registerBuild({
buildId: buildId,
cutoffTime: cutoffTime,
status: "RUNNING"
});
try {
const input = await readSourceSnapshot(cutoffTime);
await validateSourceData(input);
const partitions = partitionInput(input);
await processInParallel(partitions, async function (partition) {
const output = recomputeDerivedRecords(partition);
await writeToStaging(buildId, output);
});
const validation = await validateStagingOutput(buildId);
if (!validation.isValid) {
throw new Error("Output validation failed");
}
await publishBuildAtomically(buildId);
await markBuildComplete(buildId);
} catch (error) {
await markBuildFailed(buildId, error.message);
throw error;
}
}
This workflow writes every result into staging storage and publishes only after validation succeeds. A retry can reuse or safely replace output associated with the same build identifier.
Idempotency and Duplicate Protection
Batch jobs must be safe to retry. A worker may finish writing its output but fail before reporting success. Without idempotency, a retry could create duplicate records or apply the same aggregate more than once.
In practical terms, processing the same input for the same build multiple times should not produce additional or inconsistent output.
Partitioning the Work
Large datasets are divided into partitions so multiple workers can process them in parallel. Common partitioning strategies include:
Partitioning Strategies
- Time partitions, such as day, hour, or month
- Hash partitions based on an entity identifier
- Range partitions based on ordered keys
- Geographical or tenant-based partitions
- Source-file or object-storage partitions
- Search-index shards
Balanced partitions improve parallelism. Poor partitioning can create stragglers, where most workers finish early while one worker continues processing a disproportionately large or expensive partition.
The slowest partition often determines the completion time of the whole stage. Partition-size distribution is therefore as important as the average amount of work.
Safe Publication with Versioned Outputs
Consumers should not read the output while it is only partially built. A safer design creates an immutable version and changes a small pointer only after validation succeeds.
{
"activeBuild": "ranking-2026-09-24",
"previousBuild": "ranking-2026-09-23",
"sourceCutoff": "2026-09-24T00:00:00Z",
"status": "PUBLISHED"
}
The active-build pointer may be stored in a database, metadata service, search alias, configuration store, or catalog. The publication operation should be atomic from the consumer's perspective.
In-Place Replacement
- Readers may see incomplete output
- Rollback can be difficult
- Concurrent updates can create mixed versions
- Validation happens too late
Versioned Publication
- The current version remains available during the build
- The new version can be validated independently
- Publication requires only a pointer or alias switch
- The previous version supports fast rollback
Checkpoints, Manifests, and Resumability
A large recomputation should not necessarily restart from the beginning after a single worker fails. The pipeline can record progress at stable boundaries.
A manifest describes the expected input, output partitions, code version, source cutoff, and completion state of a particular build.
{
"buildId": "search-index-2026-09-24",
"sourceCutoff": "2026-09-24T00:00:00Z",
"transformationVersion": "ranking-v4",
"expectedPartitions": 4,
"completedPartitions": [
"partition-00",
"partition-01",
"partition-03"
],
"pendingPartitions": [
"partition-02"
],
"status": "RUNNING"
}
After a failure, the orchestrator can retry the missing partition instead of processing every successful partition again. The resumed job must still verify that the input snapshot and transformation version have not changed.
Freshness and Recomputation Frequency
Batch output represents source data only up to a particular cutoff. The difference between the current time and the newest source data reflected in the output is commonly described as data staleness.
Increasing recomputation frequency can improve freshness, but it also increases compute usage, operational load, and the chance that one run overlaps with the next.
Common Failure Scenarios
Incomplete Source Snapshot
The job starts before all expected source files or partitions arrive. The output appears valid but silently excludes data. Use manifests, expected-partition checks, and explicit cutoff conditions.
Worker Failure
A worker crashes after processing part of the data. Use deterministic partitions, retries, idempotent writes, and checkpointed completion records.
Bad Transformation Logic
A code deployment produces incorrect rankings or aggregates. Record the transformation version, validate output, use a controlled rollout, and preserve the previous output.
Overlapping Runs
A new scheduled run begins before the previous run finishes. Use concurrency controls, build identifiers, cancellation policies, or isolated output versions.
Publishing Partial Output
Consumers are redirected before all partitions are ready. Require completion and validation gates before changing the serving pointer.
Schema Incompatibility
The output schema changes before consumers are ready. Validate contracts and use compatible schema evolution or coordinated deployment.
Validation Before Publication
A job completing successfully does not prove that its output is correct. Technical and business-level validations should run before publication.
| Validation Type | Example Check |
|---|---|
| Completeness | All expected partitions were processed |
| Uniqueness | No duplicate document or product identifiers exist |
| Schema | Required fields and compatible data types are present |
| Range | Ranking scores remain within acceptable boundaries |
| Distribution | Major changes from the previous version are investigated |
| Search quality | Known queries continue returning expected results |
| Permission safety | Restricted documents are not exposed to unauthorized users |
| Serving readiness | The new output can be opened and queried successfully |
Monitoring and Observability
Monitoring should explain whether the job started, what it is processing, where time is being spent, whether output is valid, and whether the new version was published.
Important Metrics
- Job start time, completion time, and total duration
- Records, files, or bytes processed
- Processing throughput per worker or partition
- Successful, failed, pending, and retried partitions
- Slowest partition and partition-size distribution
- Input and output record counts
- Validation failures and rejected records
- Time since the last successful publication
- Current source-to-output freshness lag
- Compute, storage, and network consumption
Capacity and Cost Estimation
A basic estimate can help determine whether a full recomputation can finish within the available processing window.
In this simplified model:
- \(D\) is the total amount of input data.
- \(N\) is the number of workers.
- \(P\) is the effective throughput of one worker.
- \(E\) is the parallel-efficiency factor, where \(0 < E \leq 1\).
Real systems also include startup time, data shuffling, skew, throttling, retries, output commits, and validation. Therefore, this formula should be treated as an initial estimate rather than a guaranteed runtime.
Optimization Techniques
Ways to Improve the Pipeline
- Partition data so workers can process it independently
- Use columnar and compressed formats for analytical input
- Filter unnecessary records and columns early
- Place computation close to the data
- Reuse stable intermediate outputs where correctness permits
- Separate exceptionally expensive records from normal work
- Use checkpointing to avoid restarting completed partitions
- Apply incremental processing when the change set is small
- Autoscale workers within controlled capacity limits
- Retain previous versions for fast rollback
Example: Idempotent Partition Replacement
One safe pattern is to compute a complete partition in staging and replace the corresponding destination partition only after the staged data passes validation.
BEGIN;
DELETE FROM daily_product_metrics
WHERE metric_date = '2026-09-23';
INSERT INTO daily_product_metrics (
metric_date,
product_id,
views,
clicks,
purchases
)
SELECT
metric_date,
product_id,
views,
clicks,
purchases
FROM daily_product_metrics_staging
WHERE build_id = 'metrics-2026-09-24'
AND metric_date = '2026-09-23';
COMMIT;
Replacing the complete destination partition avoids adding the same aggregation twice. For very large partitions, the platform may offer more efficient partition-exchange or versioned-table mechanisms.
Example System Design
Consider a product-search platform that receives product records, catalog updates, inventory changes, click events, and purchase events. It needs to periodically refresh both its search index and ranking features.
Proposed Flow
- Business services write authoritative product and transaction data.
- Raw records and events are copied into durable storage.
- The orchestrator creates a build identifier and source cutoff.
- Distributed workers clean records and calculate ranking features.
- Indexing workers create a new version of the search index.
- Automated checks validate document counts, schema, permissions, and known search queries.
- The search alias is atomically moved to the new index.
- The previous index remains available for rollback.
System Design Interview Discussion
During a system design interview, avoid saying only that a cron job will rebuild the index. Explain the correctness and operational decisions behind the job.
| Question | What Your Design Should Explain |
|---|---|
| What data is processed? | The source of truth and the stable input boundary |
| How is the work distributed? | The partition key, number of workers, and skew handling |
| What happens after a failure? | Retries, checkpoints, idempotency, and resumability |
| How is partial output prevented? | Staging, validation, and atomic publication |
| How does rollback work? | Version retention and pointer reversal |
| How fresh is the output? | Cutoff time, schedule, runtime, and freshness objective |
| How does the design scale? | Partitioning, parallelism, storage layout, and capacity |
| How is correctness verified? | Technical checks, business checks, and comparison tests |
Best-Practice Checklist
Production Readiness
- Identify the authoritative source of truth
- Record a stable source cutoff for every run
- Assign a unique build identifier
- Version transformation code and configuration
- Make partition processing idempotent
- Checkpoint completed work
- Write into isolated staging storage
- Validate output before publication
- Publish through an atomic pointer or alias change
- Retain the previous successful version
- Monitor runtime, failures, skew, cost, and freshness
- Test recovery, rollback, and historical backfills
- Prevent unsafe overlap between scheduled runs
- Define ownership and an operational runbook
Knowledge Check
Why should a recomputation job write to staging?
Staging prevents consumers from reading incomplete output and allows validation before publication.
Why is idempotency important?
Distributed jobs can retry work after failures. Idempotency ensures that repeated processing does not create duplicate or inconsistent results.
When is full recomputation useful?
It is useful when the dataset can be rebuilt economically, transformation logic changes significantly, or existing derived state must be repaired.
What is the main advantage of incremental recomputation?
It reduces processing work by applying only the changes since the previous valid state.
What makes versioned publication safer?
The existing version remains available until the new version is complete and validated, and rollback can restore the previous pointer.
Summary
Batch recomputation rebuilds derived state from a bounded set of authoritative data. It is commonly used for search indexes, ranking features, recommendations, analytics, autocomplete, and materialized aggregates.
Full recomputation offers a simpler correctness model and can repair accumulated errors, but it becomes expensive as data grows. Incremental recomputation reduces repeated work and can improve freshness, but it requires careful handling of state, deletions, retries, late data, and duplicates.
A reliable design uses stable input boundaries, versioned transformation logic, balanced partitions, idempotent workers, checkpoints, staging output, comprehensive validation, atomic publication, monitoring, and rollback support.
Key Takeaway
Do not rebuild data directly in front of users. Build a new version from a defined source snapshot, validate it, publish it atomically, and preserve the previous version so the system can recover quickly.