fan-out
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.
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.
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
- 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
Common Forms of Fan-Out
Request Fan-Out
One synchronous request calls several downstream services. An API aggregator may request profile, order, payment, and delivery information in parallel.
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.
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.
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.
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.
- The client submits a search query.
- The coordinator analyzes filters and routing information.
- The query is sent to every required shard in parallel.
- Each shard evaluates matching documents locally.
- Each shard returns its highest-scoring candidates.
- The coordinator merges and reranks the candidates.
- 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:
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.
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.
{
"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.
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.
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.
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.
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.
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.
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
What is fan-out?
Fan-out is the distribution of one request, event, or task into multiple downstream operations.
What is fan-in?
Fan-in collects and combines multiple partial outputs into a final result, aggregate, or workflow state.
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.
What is fan-out on write?
It distributes or materializes data for recipients when the content is written, making later reads less expensive.
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.