Table of Contents

    fan-out

    SYSTEM DESIGN • CHAPTER 12.6

    Fan-Out

    Learn how fan-out distributes one request, event, or task across multiple workers, services, consumers, or data partitions, and how to control the resulting latency, load, and failure risks.

    Learning objective: By the end of this article, you will understand fan-out, fan-in, scatter-gather, fan-out on write, fan-out on read, amplification, partial failure, and practical fan-out strategies for search, feeds, notifications, and data pipelines.

    Prerequisites

    Recommended Knowledge

    • Client-server communication and APIs
    • Queues, topics, producers, and consumers
    • Database and search-index sharding
    • Parallel and asynchronous processing
    • Timeouts, retries, and idempotency
    • Load balancing and horizontal scaling
    • Caching and eventual consistency
    • Basic ranking and pagination concepts

    What Is Fan-Out?

    Fan-out is a distributed-system pattern in which one incoming request, event, or unit of work produces multiple downstream operations. Those operations may run across several services, workers, queues, consumers, partitions, or search shards.

    FAN-OUT MODEL
    One Input Many Independent Operations

    For example, a search coordinator may send one query to every relevant index shard. Each shard searches its local documents and returns its best candidates. Similarly, one published event may be delivered to several independent consumers.

    Simple Analogy

    A coordinator divides one large assignment among several specialists. Each specialist completes one portion in parallel. The coordinator may later collect and combine their results.

    Fan-Out and Fan-In

    Fan-out distributes work, while fan-in collects the resulting outputs. The two patterns frequently appear together.

    FAN-OUT AND FAN-IN
    Request Many Workers Combined Result

    Fan-Out

    • One input creates multiple operations
    • Distributes work across components
    • Enables parallel execution
    • Increases downstream traffic

    Fan-In

    • Collects multiple partial results
    • Merges, aggregates, or ranks outputs
    • Coordinates completion and timeouts
    • Produces one final response or dataset
    Fan-out improves parallelism and reach, but it also multiplies work. A scalable design must control both distribution and aggregation.

    Common Forms of Fan-Out

    1

    Request Fan-Out

    One synchronous request calls several downstream services. An API aggregator may request profile, order, payment, and delivery information in parallel.

    2

    Query Fan-Out

    One query is distributed to multiple database or search shards. Partial results are returned to a coordinator and merged into a global result.

    3

    Event Fan-Out

    One event is delivered to multiple independent subscribers. Each subscriber performs a different action without requiring the producer to invoke each consumer directly.

    4

    Task Fan-Out

    A large job is divided into smaller tasks and distributed across workers. This pattern is common in indexing, image processing, ETL, and batch recomputation.

    5

    Data Fan-Out

    One data change is copied into several derived systems, such as a search index, analytics store, cache, audit system, and recommendation platform.

    Fan-Out in Distributed Search

    A large search index is commonly divided into shards. Each shard holds only part of the indexed document collection. When a query can match documents on multiple shards, a coordinator distributes the query to those shards.

    DISTRIBUTED SEARCH FLOW
    Query Coordinator Search Shards Merge Top Results
    1. The client submits a search query.
    2. The coordinator analyzes filters and routing information.
    3. The query is sent to every required shard in parallel.
    4. Each shard evaluates matching documents locally.
    5. Each shard returns its highest-scoring candidates.
    6. The coordinator merges and reranks the candidates.
    7. The final page of results is returned to the client.

    Global Top-K Selection

    Suppose there are \(S\) shards and each shard returns its local top \(K\) candidates. The coordinator may need to examine as many as:

    CANDIDATE COUNT
    \[ C = S \times K \]

    If 20 shards each return 100 candidates, the coordinator receives up to 2,000 candidates before selecting the final global result. This is an example for understanding the calculation, not a required production configuration.

    async function distributedSearch(query, shards, limit) {
        const requests = shards.map(function (shard) {
            return searchShard(shard, query, limit);
        });
    
        const shardResponses = await Promise.allSettled(requests);
    
        const candidates = shardResponses
            .filter(function (response) {
                return response.status === "fulfilled";
            })
            .flatMap(function (response) {
                return response.value.results;
            });
    
        return candidates
            .sort(function (first, second) {
                return second.score - first.score;
            })
            .slice(0, limit);
    }

    This simplified example sends requests concurrently, retains successful responses, combines their candidates, and selects the highest-scoring results. A production implementation must also define timeout, retry, deduplication, pagination, and partial result behavior.

    Scatter-Gather

    Scatter-gather is a common request pattern closely related to fan-out and fan-in. The coordinator scatters requests across multiple destinations and gathers their responses.

    Stage Responsibility Main Risk
    Route Identify required shards or services Contacting unnecessary destinations
    Scatter Send requests concurrently Connection and traffic amplification
    Execute Process work locally Slow or overloaded workers
    Gather Collect responses Waiting indefinitely for stragglers
    Merge Combine, deduplicate, and rank results Coordinator CPU and memory pressure
    Respond Return the final result Incomplete or stale output

    Tail Latency

    In a synchronous fan-out operation, the overall response often depends on the slowest required downstream request. Adding more destinations increases the possibility that at least one operation will be slow or unavailable.

    SIMPLIFIED RESPONSE TIME
    \[ T_{\text{response}} \approx \max(T_1,T_2,\ldots,T_n) + T_{\text{merge}} \]

    This simplified formula assumes parallel execution. It shows why average shard latency alone is insufficient. The coordinator should also monitor slowest-shard latency, timeout rates, and the distribution of end-to-end response times.

    Tail-Latency Controls

    • Use strict per-request and end-to-end deadlines
    • Route only to shards that may contain relevant data
    • Keep shard workload balanced
    • Limit expensive query behavior
    • Use cached results where acceptable
    • Return partial results when the product permits it
    • Apply carefully controlled hedged requests
    • Cancel downstream work after the request deadline

    Event Fan-Out with Publish and Subscribe

    In publish-subscribe systems, a producer publishes one event to a topic. Multiple subscriber groups can independently receive that event and perform separate business operations.

    EVENT FAN-OUT
    Product Updated Search Indexer Cache Updater Analytics Consumer Audit Consumer
    {
        "eventId": "evt-6729",
        "eventType": "ProductUpdated",
        "entityId": "product-501",
        "occurredAt": "2026-09-24T05:00:00Z",
        "schemaVersion": 3,
        "data": {
            "name": "Wireless Keyboard",
            "price": 2499.00,
            "available": true
        }
    }

    Each consumer should use the stable event identifier to make retries safe. Consumers should also tolerate compatible schema evolution and avoid depending on the processing speed of other subscribers.

    Fan-Out on Write

    Fan-out on write, also called the push model, distributes an item when it is created. In a social feed, a newly published post may be inserted into the precomputed feed of each follower.

    WRITE AMPLIFICATION
    \[ W_{\text{fan-out}} \approx F \times C \]

    Here, \(F\) is the number of recipients and \(C\) is the number of downstream writes required for each recipient. A publisher with a large audience may therefore create a significant write burst.

    Advantages

    • Fast reads from precomputed recipient feeds
    • Less aggregation during feed retrieval
    • Suitable for frequently read feeds
    • Ranking can be partially prepared in advance

    Limitations

    • One write can generate many downstream writes
    • High-follower publishers may create severe bursts
    • Inactive recipients still receive precomputed entries
    • Deletion and permission changes require propagation

    Fan-Out on Read

    Fan-out on read, also called the pull model, stores content once in the publisher's stream. When a user requests a feed, the system reads recent items from relevant publishers and merges them.

    Advantages

    • New content requires fewer immediate writes
    • Avoids precomputing feeds for inactive users
    • Permission changes can be applied at read time
    • Well suited to extremely high-follower publishers

    Limitations

    • Feed reads require multiple lookups
    • Merge and ranking occur on the read path
    • Read latency grows with the number of sources
    • Cache misses can create expensive requests
    Dimension Fan-Out on Write Fan-Out on Read
    Primary work Publishing time Reading time
    Write cost Potentially high Usually lower
    Read cost Usually lower Potentially high
    Storage May duplicate feed references Stores publisher content more centrally
    High-follower publishers Can cause major write amplification Avoids mass writes during publication
    Inactive users May receive unnecessary materialized entries Generate work only when they read

    Hybrid Fan-Out

    A hybrid design applies different strategies to different publishers, recipients, queries, or workloads. For example, ordinary publishers may use fan-out on write, while unusually high-follower publishers are merged into the feed during reads.

    HYBRID PRINCIPLE
    Precompute predictable and affordable work. Defer highly amplified or infrequently consumed work until it is needed.

    The routing decision may consider follower count, user activity, content type, permission model, write rate, feed-read rate, and current system capacity.

    Fan-Out Amplification

    Fan-out creates amplification because one external operation becomes several internal operations.

    INTERNAL REQUEST RATE
    \[ Q_{\text{internal}} = Q_{\text{external}} \times F \]

    Here, \(Q_{\text{external}}\) is the incoming request rate and \(F\) is the average number of downstream destinations. Retries, replicas, and secondary workflows may increase the actual load further.

    Important: Capacity planning must use internal amplified traffic, not only the number of requests received from clients.

    Partial Failure

    A fan-out operation may succeed on some destinations and fail or time out on others. The system must define whether partial output is acceptable.

    Strategy Behavior Suitable Situation
    Fail the complete request Return an error if any required operation fails All components are required for correctness
    Return partial results Use available responses and mark incompleteness Availability is more important than completeness
    Use cached data Substitute a recent cached response Stale data is acceptable
    Retry asynchronously Complete missing side effects later The client does not need immediate completion
    Use a fallback Replace an unavailable advanced result with a simpler one Graceful degradation is supported

    Retries and Retry Storms

    Retrying every failed branch immediately can multiply load during an outage. If one request fans out to many services and each layer retries independently, the dependency may receive far more traffic than normal.

    Unsafe Approach Retry every failed branch immediately and indefinitely without considering the request deadline or dependency health.
    Safer Approach Use bounded retries, exponential backoff, randomized jitter, retry budgets, idempotency, circuit breakers, and one clearly defined retry-owning layer.
    async function callWithRetry(operation, options) {
        let attempt = 0;
    
        while (attempt < options.maxAttempts) {
            try {
                return await operation();
            } catch (error) {
                attempt += 1;
    
                if (attempt >= options.maxAttempts) {
                    throw error;
                }
    
                const exponentialDelay =
                    options.baseDelayMs * Math.pow(2, attempt - 1);
    
                const jitter =
                    Math.floor(Math.random() * options.jitterMs);
    
                await sleep(exponentialDelay + jitter);
            }
        }
    }

    Backpressure and Concurrency Control

    Unbounded fan-out can exhaust connection pools, threads, memory, queue capacity, or downstream service limits. The coordinator should restrict how many branches are active simultaneously.

    async function mapWithConcurrency(items, limit, worker) {
        const results = new Array(items.length);
        let nextIndex = 0;
    
        async function runWorker() {
            while (true) {
                const currentIndex = nextIndex;
                nextIndex += 1;
    
                if (currentIndex >= items.length) {
                    return;
                }
    
                results[currentIndex] =
                    await worker(items[currentIndex]);
            }
        }
    
        const workers = [];
    
        for (let index = 0; index < limit; index += 1) {
            workers.push(runWorker());
        }
    
        await Promise.all(workers);
    
        return results;
    }

    Concurrency Controls

    • Limit concurrent downstream requests
    • Apply per-service and per-tenant quotas
    • Use bounded queues
    • Reject or defer low-priority work during overload
    • Propagate deadlines and cancellations
    • Separate interactive and batch capacity
    • Monitor queue depth and consumer lag
    • Protect downstream dependencies with circuit breakers

    Selective Fan-Out

    A system should avoid broadcasting every request to every destination when routing metadata can identify a smaller target set.

    Broadcast Fan-Out

    • Contacts every shard or subscriber
    • Simple routing logic
    • Higher resource consumption
    • Greater sensitivity to slow destinations

    Selective Fan-Out

    • Contacts only relevant destinations
    • Reduces internal request volume
    • Lowers aggregation cost
    • Requires accurate routing metadata

    A search query filtered by region, tenant, date, or category may be routed only to shards that can contain matching documents. Routing metadata must remain accurate enough to avoid missing valid results.

    Ranking and Pagination

    Distributed ranking is more complex than independently sorting results on each shard. A document's local position does not automatically determine its global position.

    Each shard commonly returns more candidates than the final page requires. The coordinator merges candidates, applies consistent scoring rules, removes duplicates, and produces the requested page.

    Deep pagination concern: Asking every shard for a large offset can significantly increase computation, network transfer, and coordinator memory. Cursor or search-after pagination is generally safer for deep result sets.

    Idempotency and Deduplication

    Asynchronous fan-out frequently uses at-least-once delivery. Therefore, the same event or task may reach a consumer more than once.

    BEGIN;
    
    INSERT INTO processed_events (
        consumer_name,
        event_id,
        processed_at
    )
    VALUES (
        'search-indexer',
        'evt-6729',
        CURRENT_TIMESTAMP
    )
    ON CONFLICT (consumer_name, event_id) DO NOTHING;
    
    UPDATE search_documents
    SET
        title = 'Wireless Keyboard',
        price = 2499.00,
        is_available = TRUE
    WHERE product_id = 501;
    
    COMMIT;

    In a real implementation, the deduplication record and business update must be coordinated so a crash cannot mark an event complete without applying its effect, or apply the effect without recording completion.

    Monitoring Fan-Out

    Important Metrics

    • Fan-out width per request or event
    • Total downstream request rate
    • Successful, failed, timed-out, and cancelled branches
    • Per-shard and per-service latency
    • Slowest branch latency
    • Coordinator merge duration
    • Partial-result rate
    • Retry rate and retry amplification
    • Queue depth and oldest-message age
    • Duplicate-event rate
    • Hot publishers, tenants, queries, and partitions
    • Connection-pool and worker utilization

    Common Design Mistakes

    Weak Design

    • Broadcasting every request to every shard
    • Launching unlimited parallel operations
    • Waiting forever for every destination
    • Retrying independently at multiple layers
    • Ignoring duplicate message delivery
    • Publishing partial feed updates without tracking
    • Using push fan-out for every high-follower account
    • Monitoring only the coordinator

    Strong Design

    • Routes requests only to relevant destinations
    • Limits concurrency and queue size
    • Uses explicit deadlines and cancellation
    • Defines a bounded retry policy
    • Makes consumers idempotent
    • Supports partial responses where appropriate
    • Uses hybrid push and pull strategies
    • Measures downstream amplification

    System Design Interview Discussion

    Question What Your Design Should Explain
    What triggers the fan-out? A query, write, event, scheduled job, or user action
    How wide can it become? Average, peak, and exceptional destination counts
    Is it synchronous? Whether the caller waits for downstream completion
    What happens after partial failure? Error, partial response, fallback, or asynchronous retry
    How is overload controlled? Concurrency limits, quotas, queues, and backpressure
    How are duplicates handled? Idempotency keys and deterministic updates
    How are results combined? Merge, aggregation, deduplication, and ranking rules
    How is latency controlled? Deadlines, routing, caching, cancellation, and degradation

    Production Readiness Checklist

    Fan-Out Checklist

    • Measure average and maximum fan-out width
    • Estimate amplified internal traffic
    • Route only to relevant destinations
    • Set end-to-end and per-branch deadlines
    • Limit active parallel operations
    • Use bounded queues and backpressure
    • Define partial-failure behavior
    • Make retries bounded and idempotent
    • Deduplicate repeated events and results
    • Protect high-follower or hot-key workloads
    • Retain enough information for replay and repair
    • Monitor individual shards and downstream services
    • Test slow, failed, and unreachable destinations
    • Document fallback and degradation behavior

    Knowledge Check

    1

    What is fan-out?

    Fan-out is the distribution of one request, event, or task into multiple downstream operations.

    2

    What is fan-in?

    Fan-in collects and combines multiple partial outputs into a final result, aggregate, or workflow state.

    3

    Why does fan-out affect tail latency?

    A synchronous coordinator may need to wait for the slowest required branch, making the complete response sensitive to stragglers.

    4

    What is fan-out on write?

    It distributes or materializes data for recipients when the content is written, making later reads less expensive.

    5

    Why use a hybrid feed strategy?

    It balances fast reads for ordinary workloads with controlled write amplification for unusually large publishers.

    Summary

    Fan-out distributes one unit of work across multiple downstream destinations. It supports parallel processing, distributed search, publish-subscribe delivery, feed generation, notifications, indexing, and large-scale data processing.

    Distributed search commonly combines fan-out with fan-in. A coordinator sends a query to relevant shards, gathers local candidates, deduplicates them, merges their scores, and returns the final global results.

    Fan-out on write performs additional work during publication to make reads faster. Fan-out on read stores content more centrally and performs aggregation when the user requests it. Hybrid designs use both strategies according to audience size, activity, cost, and freshness requirements.

    Reliable fan-out requires selective routing, concurrency limits, backpressure, deadlines, bounded retries, idempotency, deduplication, partial-failure policies, and detailed monitoring of downstream amplification and tail latency.

    Key Takeaway

    Fan-out converts one operation into many. Use it to gain parallelism and reach, but control its width, protect downstream systems, plan for partial failure, and measure the total internal workload it creates.