Event-driven document processing
Replaced synchronous flows with Azure Service Bus topics and subscriptions. Throughput went up by about 40%.
Context
An insurance and agents B2B platform takes in documents in bulk. Uploads arrive as a hierarchy, and each document passes through several levels of validation before the platform can use it.
The platform runs on .NET and Azure, with an Angular front end.
Problem
Document processing ran as synchronous calls. With synchronous calls, the request that receives a document also does the work on it. Every slow step makes the caller wait, and a burst of uploads lines up behind whatever is already running.
There is also no natural place to retry. If one step fails part-way through, the whole request fails with it.
What I did
I moved the processing onto Azure Service Bus. Receiving a document now means publishing a message and returning. The work happens in consumers that read from the bus.
[HttpPost("documents")] public async Task<IActionResult> Upload(DocumentUploaded document) {- // The caller waits while every step runs, one after another.- await validator.ValidateAsync(document);- await indexer.IndexAsync(document);- await notifier.NotifyAsync(document);- return Ok();+ // Publish once and return. Each subscription on the topic does its part in parallel.+ await sender.SendMessageAsync(new ServiceBusMessage(BinaryData.FromObjectAsJson(document)));+ return Accepted(); }Topics and subscriptions let one uploaded document feed several independent steps. Each subscription gets its own copy of the message and its own consumers, so a slow step does not hold up a fast one.
Consumers run in parallel, so throughput grows with the number of consumers. Service Bus redelivers a message that fails. One that keeps failing moves to the dead-letter queue, where it can be inspected instead of lost.
On the .NET side the design follows CQRS. Commands that change state go through the bus, and queries read along their own path.
await using var client = new ServiceBusClient(connectionString);
var processor = client.CreateProcessor("documents", "validation", new ServiceBusProcessorOptions
{
MaxConcurrentCalls = 4,
AutoCompleteMessages = false,
});
processor.ProcessMessageAsync += async args =>
{
var document = args.Message.Body.ToObjectFromJson<DocumentUploaded>();
try
{
await validator.ValidateAsync(document, args.CancellationToken);
await args.CompleteMessageAsync(args.Message);
}
catch (ValidationException ex)
{
// A document that fails validation will never pass: set it aside with the reason.
await args.DeadLetterMessageAsync(args.Message, "ValidationFailed", ex.Message);
}
catch (Exception)
{
// Anything else may be transient: release it so Service Bus redelivers it.
// After MaxDeliveryCount attempts it moves to the dead-letter queue.
await args.AbandonMessageAsync(args.Message);
throw;
}
};
processor.ProcessErrorAsync += args =>
{
logger.LogError(args.Exception, "Service Bus error on {Entity}", args.EntityPath);
return Task.CompletedTask;
};
await processor.StartProcessingAsync();Result
Throughput went up by about 40%. Ingestion now scales by adding consumers. A failing document is retried or set aside, instead of failing the upload.
Stack
- Azure Service Bus
- .NET
- CQRS
- Angular
Try it
A recreation with synthetic documents. Change the arrival rate, the number of consumers and the failure rate. Switch retries off to send failures straight to the dead-letter queue. Below the simulator, one topic fans out to three subscriptions.
Synchronous
One worker handles one document at a time. The rest wait.
Event-driven
A queue feeds three workers in parallel. Failures go back on the queue and are retried.
Static preview. The interactive version needs JavaScript.