The Azure Scale Controller is Wilding: Build Your Own Java Fan-Out Orchestrator
If you work with Azure Durable Functions, you likely use the Fan-Out / Fan-In pattern to process data in parallel. In practice, trusting default queue scaling in production often triggers HTTP 429 (Too Many Requests) storms, connection pool exhaustion, and downstream outages.
The Problem: Aggressive Autoscaling
Queues decouple workloads, but they don’t rate-limit them. When an orchestrator dumps thousands of task messages into an Azure Storage Queue within milliseconds, the Scale Controller sees a spike and panics—rapidly spinning up worker instances to clear the backlog. Those workers instantly bombard your downstream databases, legacy systems, or third-party APIs with hundreds of parallel requests, crushing infrastructure under its own weight.
Why Global Settings Fall Short
Setting static concurrency caps in host.json or setting instance limits in the Azure Portal creates major trade-offs:
-
host.json Limits Are Global: A static concurrency limit applies to all activity functions in the app. Lightweight logging and heavy database writes get choked equally.
-
Instance Caps Affect Every Trigger: Restricting scale-out caps HTTP endpoints or Service Bus triggers running in the same app.
-
The Data Engineer’s Dilemma: Data pipelines consume from different data sources with different capacities:
-
Legacy Relational DB: Fails above 5 concurrent connections.
-
SaaS REST API: Rate-limits above 20 requests/second.
-
Cloud Blob Storage: Handles 100+ parallel streams with ease.
-
The Solution: The Sliding Window Pattern
Instead of letting Azure’s autoscaler take over, control execution speed directly inside your Java Orchestrator. Using a Sliding Window pattern, you maintain a fixed worker pool size per batch run. As soon as one task finishes, the orchestrator immediately spawns the next item.
Code Implementation
-
HTTP Starter Function Triggers the orchestrator and returns standard status-check URLs.
package com.example; import com.microsoft.azure.functions.ExecutionContext; import com.microsoft.azure.functions.HttpMethod; import com.microsoft.azure.functions.HttpRequestMessage; import com.microsoft.azure.functions.HttpResponseMessage; import com.microsoft.azure.functions.annotation.AuthorizationLevel; import com.microsoft.azure.functions.annotation.FunctionName; import com.microsoft.azure.functions.annotation.HttpTrigger; import com.microsoft.durabletask.DurableTaskClient; import com.microsoft.durabletask.azurefunctions.DurableClientContext; import com.microsoft.durabletask.azurefunctions.DurableClientInput; import java.util.Optional; public class StartOrchestrationFunction { @FunctionName("start-fan-out") public HttpResponseMessage start( @HttpTrigger( name = "req", methods = {HttpMethod.POST}, authLevel = AuthorizationLevel.ANONYMOUS) HttpRequestMessage<Optional<String>> request, @DurableClientInput(name = "durableContext") DurableClientContext durableContext, final ExecutionContext context) { DurableTaskClient client = durableContext.getClient(); String instanceId = client.scheduleNewOrchestrationInstance("mini-fan-out-orchestrator"); return durableContext.createCheckStatusResponse(request, instanceId); } } -
Sliding Window Orchestrator Caps running tasks at maxConcurrency and feeds new items into the worker pool as running tasks finish.
package com.example; import com.microsoft.azure.functions.annotation.FunctionName; import com.microsoft.durabletask.Task; import com.microsoft.durabletask.TaskOrchestrationContext; import com.microsoft.durabletask.azurefunctions.DurableOrchestrationTrigger; import java.util.HashSet; import java.util.LinkedList; import java.util.List; import java.util.Queue; import java.util.Set; public class MiniOrchestrator { @FunctionName("mini-fan-out-orchestrator") public void run(@DurableOrchestrationTrigger(name = "ctx") TaskOrchestrationContext ctx) { int maxConcurrency = 2; List<String> items = List.of("Item1", "Item2", "Item3", "Item4", "Item5", "Item6", "Item7"); Queue<String> remaining = new LinkedList<>(items); Set<Task<?>> runningTasks = new HashSet<>(); while (runningTasks.size() < maxConcurrency && !remaining.isEmpty()) { runningTasks.add(ctx.callActivity("process-item", remaining.poll(), Void.class)); } while (!remaining.isEmpty()) { Task<?> completed = ctx.anyOf(runningTasks.toArray(new Task[0])).await(); runningTasks.remove(completed); runningTasks.add(ctx.callActivity("process-item", remaining.poll(), Void.class)); } while (!runningTasks.isEmpty()) { Task<?> completed = ctx.anyOf(runningTasks.toArray(new Task[0])).await(); runningTasks.remove(completed); } } } -
Activity Function (process-item) Executes the discrete unit of work against your target system.
package com.example; import com.microsoft.azure.functions.ExecutionContext; import com.microsoft.azure.functions.annotation.FunctionName; import com.microsoft.durabletask.azurefunctions.DurableActivityTrigger; public class ProcessItemActivity { @FunctionName("process-item") public void processItem( @DurableActivityTrigger(name = "item") String item, final ExecutionContext context) { context.getLogger().info("Processing batch item: " + item); // Simulate database or API work try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } context.getLogger().info("Successfully processed: " + item); } }
Conclusion
Global host settings and infrastructure scale limits are blunt tools that treat every function identically. By controlling concurrency directly inside your Java orchestrator with a sliding window, you gain fine-grained execution limits tailored to each downstream target. This pattern protects fragile dependencies from HTTP 429 errors and connection pool exhaustion, transforming unpredictable cloud autoscaling into a controlled, production-ready pipeline.