Implementando Azure Service Bus con Clean Architecture
Implementa Azure Service Bus con Clean Architecture, dos proyectos desacoplados y ejemplos reales

Hey devs, este artículo está buenísimo pero primero entendamos el contexto.

Introducción

En los sistemas modernos, la comunicación asíncrona es esencial para desacoplar los microservicios y mejorar la escalabilidad.

Cuando una parte del sistema necesita notificar a otra que algo ocurrió, como por ejemplo "se creó un pedido", lo ideal no es hacerlo mediante una llamada HTTP clásica, sino publicar un mensaje que otra parte del sistema (otro microservicio, otra API) pueda consumir cuando esté lista.

Ahí es donde entra en juego Azure Service Bus, la solución de mensajería empresarial de Microsoft.

¿Qué es Azure Service Bus?

Es un message broker (intermediario de mensajes) de la nube de Azure.

Permite que las aplicaciones envíen y reciban mensajes de forma confiable y asíncrona, asegurando la entrega garantizada incluso si un servicio está caído temporalmente.

Soporta dos modos de mensajería:

TipoDescripciónPatrón
QueuesUn mensaje es recibido por un sólo consumidorPoint-to-point
Topics y SubscriptionsUn mensaje puede ser recibido por múltiples consumidoresPub/Sub o publish-subscribe

En este artículo, estimados devs, usaremos queues para mantenerlo simple y didáctico, pero con calidad de producción, eso sí. Ya en un futuro haremos Pub/Sub con topics y Masstransit.

Escenario real

Tenemos una empresa de e-commerce llamada MyToyStore con dos microservicios (tiene más pero aquí nos interesan sólo dos):

  • Pedidos API (Publisher)
    • Cuando se crea un pedido, publica un mensaje en Service Bus.
  • Facturacion API (Subscriber)
    • Escucha los mensajes de "pedidos" y genera la factura correspondiente.

Este flujo permite que Pedidos y Facturación estén completamente desacoplados, lo cual quiere decir que si el servicio de facturación está caído, los mensajes quedan en la cola hasta que vuelva a estar disponible.

Estructura de la solución (Clean Architecture)

Tenemos dos proyectos independientes con arquitectura limpia:

Proyecto 1: Pedidos API (Publisher)

Proyecto 2: Facturacion API (Subscriber)

Publisher: Enviando mensajes desde Pedidos.API

Application layer

Aquí tenemos el código que corresponde a IMessageBus, CreatePedidoCommand y CreatePedidoHandler.

public interface IMessageBus
{
    Task PublishAsync(string queueName, object message);
}
public record CreatePedidoCommand(string Cliente, decimal Monto);
public class CreatePedidoHandler
{
    private readonly IMessageBus bus;
    private readonly ILogger<CreatePedidoHandler> logger;

    public CreatePedidoHandler(IMessageBus bus, ILogger<CreatePedidoHandler> logger)
    {
        this.bus = bus;
        this.logger = logger;
    }

    public async Task<Guid> Handle(CreatePedidoCommand command)
    {
        var pedidoId = Guid.NewGuid();

        // Aquí guardarías en base de datos (omitido por simplicidad)

        await bus.PublishAsync("pedidos", new
        {
            PedidoId = pedidoId,
            Cliente = command.Cliente,
            Monto = command.Monto
        });

        logger.LogInformation($"Pedido publicado: {pedidoId}");
        return pedidoId;
    }
}

Infrastructure layer

Aquí está la implementación del IMessageBus, es decir la clase AzureServiceBus

using Azure.Messaging.ServiceBus;
using System.Text.Json;

public class AzureServiceBus : IMessageBus, IAsyncDisposable
{
    private readonly ServiceBusClient client;

    public AzureServiceBus(string connectionString)
    {
        client = new ServiceBusClient(connectionString);
    }

    public async Task PublishAsync(string queueName, object message)
    {
        var sender = client.CreateSender(queueName);
        var body = JsonSerializer.Serialize(message);
        await sender.SendMessageAsync(new ServiceBusMessage(body));
    }

    public ValueTask DisposeAsync() => client.DisposeAsync();
}

Ahora toca registrar lo necesario en Program.cs

var builder = WebApplication.CreateBuilder(args);

builder.Services.AddSingleton<IMessageBus>(
    _ => new AzureServiceBus(builder.Configuration["AzureServiceBus:ConnectionString"])
);

builder.Services.AddScoped<CreatePedidoHandler>();

var app = builder.Build();

// Opcionalmente añadimos endpoint con minimal apis
app.MapPost("/pedidos", async (CreatePedidoCommand command, CreatePedidoHandler handler) =>
{
    var id = await handler.Handle(command);
    return Results.Ok(new { PedidoId = id });
});

app.Run();

En este punto, Pedidos API puede publicar mensajes en una cola llamada pedidos.

Estamos suponiendo que a estas alturas ya tienes una cadena de conexión de Azure Service Bus y la estás configurando en appsettings.development.json, el cual podrías usar para desarrollo local y que cuando subas a producción tendrás que usar algo más seguro como environment variables de Azure AppService.

Subscriber: Recibiendo mensajes desde Facturacion API

A continuación te muestro la clase Consumer y un DTO:

using Azure.Messaging.ServiceBus;
using Microsoft.Extensions.Logging;
using System.Text.Json;

public class PedidoCreatedConsumer
{
    private readonly ServiceBusProcessor processor;
    private readonly ILogger<PedidoCreatedConsumer> logger;

    public PedidoCreatedConsumer(ServiceBusClient client, ILogger<PedidoCreatedConsumer> logger)
    {
        processor = client.CreateProcessor("pedidos", new ServiceBusProcessorOptions());
        processor.ProcessMessageAsync += ProcessMessageHandler;
        processor.ProcessErrorAsync += ErrorHandler;
        this.logger = logger;
    }

    private async Task ProcessMessageHandler(ProcessMessageEventArgs args)
    {
        var body = args.Message.Body.ToString();
        var evento = JsonSerializer.Deserialize<PedidoCreadoDto>(body);

        logger.LogInformation($"Pedido recibido: {evento.PedidoId} - Cliente: {evento.Cliente}");

        // Aquí iría la lógica real, por ejemplo:
        // await facturaService.GenerarFactura(evento.PedidoId);

        await args.CompleteMessageAsync(args.Message);
    }

    private Task ErrorHandler(ProcessErrorEventArgs args)
    {
        logger.LogInformation($"Error en Service Bus: {args.Exception.Message}");
        return Task.CompletedTask;
    }

    public async Task StartAsync() => await processor.StartProcessingAsync();
    public async Task StopAsync() => await processor.StopProcessingAsync();
}

public record PedidoCreadoDto(Guid PedidoId, string Cliente, decimal Monto);

Ahora registramos lo necesario y además creamos el background worker

var builder = Host.CreateDefaultBuilder(args)
    .ConfigureServices((context, services) =>
    {
        var connection = context.Configuration["AzureServiceBus:ConnectionString"];
        services.AddSingleton(new ServiceBusClient(connection));
        services.AddSingleton<PedidoCreatedConsumer>();
        services.AddHostedService<Worker>();
    });

await builder.RunConsoleAsync();

Y aquí está la clase que hará background worker:

public class Worker : BackgroundService
{
    private readonly PedidoCreatedConsumer consumer;

    public Worker(PedidoCreatedConsumer consumer)
    {
        this.consumer = consumer;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        await consumer.StartAsync();
        await Task.Delay(Timeout.Infinite, stoppingToken);
    }

    public override async Task StopAsync(CancellationToken cancellationToken)
    {
        await consumer.StopAsync();
        await base.StopAsync(cancellationToken);
    }
}

Un background worker es un proceso que se ejecuta continuamente o de manera programada en el backend sin intervención del usuario, para realizar tareas automáticas como por ejemplo, escuchar mensajes de una cola de Azure Service Bus.

Beneficios de este enfoque

Está claro que como el servicio de Pedidos no sabe si Facturación está vivo o no, hay desacoplamiento total.

Es escalable porque puedes añadir más consumers sin tocar el publisher.

Es resiliente porque el service bus almacena los mensajes hasta que se procesen correctamente.

Al estar implementado con arquitectura limpia, Application no depende de Azure, la implementación del bus está en Infrastructure y el Domain sigue estando puro.

Escenarios donde este enfoque es ideal

EscenarioMotivo
Comunicación entre microservicios internosReduce dependencias directas HTTP
En procesos asíncronosEvita bloquear el hilo principal
Integración con sistemas externosPermite confiabilidad y reintentos
Sistemas on-premise que quieren volverse híbridos al integrarlos con cloudConecta un on-premise con entorno cloud moderno fácilmente

Mis conclusiones...

En conclusión, estimado dev, implementar Azure Service Bus sin frameworks adicionales te permite tener:

  • Un control más fino del flujo
  • Un código más transparente y fácil de depurar
  • Una base sólida para escalar tu arquitectura más adelante

Sin embargo, cuando tu ecosistema crece como cuando tienes múltiples microservicios, colas, eventos y SAGAs, manejar manualmente todos los procesadores y conexiones puede volverse complejo y tedioso, es allí en donde la estrategia debe incluir un framework como MassTransit.

En un próximo artículo veremos cómo MassTransit simplifica esta misma arquitectura, agregando Pub/Sub, reintentos automáticos y resiliencia avanzada, con mucho menos código repetitivo.

Si esta entrada te ha gustado, compártela! 🐿️💪🏼

Créditos de imagen de portada: Basado en Foto de Anton Kalinin en Unsplash

Deja una respuesta

Tu dirección de correo electrónico no será publicada. Los campos obligatorios están marcados con *