Fan-out and fan-in with Durable Functions
Processing 5,000 records one after another takes too long, and running them all at once in one function times out. Durable Functions runs them in parallel and collects the results.
A common integration job: get a list of 5,000 accounts, call an external API for each, and produce a summary at the end. In a single function, doing them one at a time is too slow, and starting them all at once with Task.WhenAll risks hitting the function's timeout and loses all progress if the instance restarts.
Durable Functions solves this with an orchestration: a function that coordinates other functions, survives restarts and can wait for many parallel tasks.
The pattern
- Fan out: start one activity per item.
- Fan in: wait for all of them and combine the results.
public class AccountSync(CrmClient crm, SyncService syncService)
{
[Function(nameof(SyncAllAccounts))]
public async Task<SyncSummary> SyncAllAccounts(
[OrchestrationTrigger] TaskOrchestrationContext context)
{
var accountIds = await context.CallActivityAsync<List<string>>(nameof(GetAccountIds));
var tasks = accountIds
.Select(id => context.CallActivityAsync<bool>(nameof(SyncAccount), id))
.ToList();
var results = await Task.WhenAll(tasks);
return new SyncSummary(Total: results.Length, Failed: results.Count(ok => !ok));
}
[Function(nameof(GetAccountIds))]
public Task<List<string>> GetAccountIds([ActivityTrigger] object? input, FunctionContext ctx)
=> crm.GetActiveAccountIdsAsync();
[Function(nameof(SyncAccount))]
public async Task<bool> SyncAccount([ActivityTrigger] string accountId, FunctionContext ctx)
{
// call the external API, write the result; return false on a handled failure
return await syncService.SyncAsync(accountId);
}
}
(The example uses the isolated worker model, with dependencies injected through the class's constructor.)
Each activity runs as its own function execution, spread across instances. The orchestrator checkpoints its progress to storage, so if an instance restarts halfway through, completed activities aren't run again.
Start the orchestration from any trigger:
string instanceId = await durableClient.ScheduleNewOrchestrationInstanceAsync(nameof(SyncAllAccounts));
The orchestrator rules
An orchestrator function is replayed from its history every time it wakes up. To replay correctly, its code must be deterministic:
- No I/O in the orchestrator. Database calls and HTTP requests go in activities.
- No
DateTime.UtcNow. Usecontext.CurrentUtcDateTime. - No random values or new GUIDs generated directly; generate them in an activity, or use
context.NewGuid().
Breaking these rules causes errors or strange behavior on replay, often only under load.
Don't overwhelm the other side
Starting 5,000 activities at once can easily exceed the external API's rate limit. Two ways to control it:
- Limit concurrency per instance with
maxConcurrentActivityFunctionsinhost.json. - Fan out in batches: process 100 items, wait, then the next 100.
Large lists
Each orchestration's history grows with every activity it starts. For very large jobs, split the work into sub-orchestrations (for example, one per batch of 500) with CallSubOrchestratorAsync, so no single history gets huge.
Takeaway
When a job needs to process many items in parallel and must survive restarts, use a Durable Functions orchestration: fan out an activity per item, fan in with Task.WhenAll, keep the orchestrator deterministic, and limit concurrency so you don't flood the systems you're calling.