Table of Contents

    batch recomputation

    SYSTEM DESIGN • CHAPTER 12.8

    Batch Recomputation

    Learn how large-scale systems periodically rebuild search indexes, rankings, recommendations, analytics, and other derived datasets from authoritative source data.

    Learning objective: By the end of this article, you will understand when batch recomputation is appropriate, how a recomputation pipeline works, and how to design it for correctness, scalability, reliability, freshness, and safe publication.

    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.

    CORE IDEA
    \[ \text{Derived State} = f(\text{Authoritative Source Data}) \]

    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
    A derived dataset should normally be reproducible from durable source data and versioned transformation logic.

    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.

    1

    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.

    2

    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.

    3

    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.

    4

    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.

    5

    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.

    6

    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.

    7

    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.

    PIPELINE FLOW
    Snapshot Validate Transform Verify Publish

    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.

    INCREMENTAL MODEL
    \[ O_{t} = O_{t-1} + \Delta(I_{t}) \]

    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.

    Choose Full Recomputation When The dataset can be rebuilt within the available processing window, correctness and simplicity are more important than minimum compute cost, the transformation logic has changed significantly, or the existing derived state may be corrupted.
    Choose Incremental Recomputation When The complete dataset is very large, only a small percentage changes between runs, freshness requirements are tighter, and the team can safely manage checkpoints, deletions, late data, and duplicate processing.
    Practical design: Many systems combine frequent incremental updates with a less frequent full recomputation. Incremental processing maintains freshness, while the complete rebuild repairs accumulated drift.

    Example: Recomputing Product Ranking

    Consider a product-search system that recalculates a ranking score from textual relevance, popularity, freshness, and product quality.

    EXAMPLE RANKING FUNCTION
    \[ Score = 0.45R + 0.25P + 0.20F + 0.10Q \]

    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.

    Unsafe Approach Add each recomputed value directly to the current result without recording which build or partition produced the update.
    Safer Approach Use deterministic output keys, build identifiers, partition identifiers, upserts, or complete partition replacement so that retrying the same work produces the same final state.
    IDEMPOTENCY PROPERTY
    \[ f(f(x)) = f(x) \]

    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.

    APPROXIMATE PARALLEL RUNTIME
    \[ T_{\text{job}} \approx \max(T_1, T_2, \ldots, T_n) + T_{\text{coordination}} \]

    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.

    DATA STALENESS
    \[ \text{Staleness} = \text{Current Time} - \text{Latest Reflected Source Time} \]

    Increasing recomputation frequency can improve freshness, but it also increases compute usage, operational load, and the chance that one run overlaps with the next.

    DESIGN RULE
    The output refresh interval should be driven by the business cost of stale data, not merely by the fastest schedule the platform can support.

    Common Failure Scenarios

    1

    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.

    2

    Worker Failure

    A worker crashes after processing part of the data. Use deterministic partitions, retries, idempotent writes, and checkpointed completion records.

    3

    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.

    4

    Overlapping Runs

    A new scheduled run begins before the previous run finishes. Use concurrency controls, build identifiers, cancellation policies, or isolated output versions.

    5

    Publishing Partial Output

    Consumers are redirected before all partitions are ready. Require completion and validation gates before changing the serving pointer.

    6

    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.

    ESTIMATED PROCESSING TIME
    \[ T \approx \frac{D}{N \times P \times E} \]

    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.

    HIGH-LEVEL ARCHITECTURE
    Database and Event Log Raw Storage Batch Compute Staging Index Search Alias

    Proposed Flow

    1. Business services write authoritative product and transaction data.
    2. Raw records and events are copied into durable storage.
    3. The orchestrator creates a build identifier and source cutoff.
    4. Distributed workers clean records and calculate ranking features.
    5. Indexing workers create a new version of the search index.
    6. Automated checks validate document counts, schema, permissions, and known search queries.
    7. The search alias is atomically moved to the new index.
    8. 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

    1

    Why should a recomputation job write to staging?

    Staging prevents consumers from reading incomplete output and allows validation before publication.

    2

    Why is idempotency important?

    Distributed jobs can retry work after failures. Idempotency ensures that repeated processing does not create duplicate or inconsistent results.

    3

    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.

    4

    What is the main advantage of incremental recomputation?

    It reduces processing work by applying only the changes since the previous valid state.

    5

    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.