ASP.NET Core  

Building Asynchronous Messaging Pipelines with MassTransit and RabbitMQ

Monolithic architectures often struggle with high concurrency workloads. Transitioning to event-driven architectures via asynchronous message brokers (like RabbitMQ) decouples services, ensures fault tolerance, and handles traffic spikes gracefully. MassTransit is the premier .NET distributed application framework that abstracts away low-level broker complexities.

Step 1: Install MassTransit and RabbitMQ Packages

Add MassTransit along with the RabbitMQ transport package via the CLI:

Bash

dotnet add package MassTransit.RabbitMQ

Step 2: Define Shared Contracts (Messages and Events)

Create a shared data contract library or namespace containing record types that represent events or commands flowing through your message broker.

C#

namespace SharedContracts
{
    // Event published when a new order is submitted
    public record OrderSubmittedEvent(Guid OrderId, string CustomerEmail, decimal TotalAmount, DateTime Timestamp);
}

Step 3: Implement a Consumer for Processing Messages

Create a consumer class that listens for OrderSubmittedEvent messages and executes business logic asynchronously.

C#

using MassTransit;
using SharedContracts;
using Microsoft.Extensions.Logging;

public class OrderSubmittedConsumer : IConsumer<OrderSubmittedEvent>
{
    private readonly ILogger<OrderSubmittedConsumer> _logger;

    public OrderSubmittedConsumer(ILogger<OrderSubmittedConsumer> logger)
    {
        _logger = logger;
    }

    public async Task Consume(ConsumeContext<OrderSubmittedEvent> context)
    {
        var message = context.Message;

        _logger.LogInformation("Processing Order ID: {OrderId} for Customer: {Email}, Amount: {Amount}", 
            message.OrderId, message.CustomerEmail, message.TotalAmount);

        // Execute downstream business tasks (e.g., generate invoice, update inventory...)
        await Task.Delay(500); 

        _logger.LogInformation("Successfully processed Order ID: {OrderId}", message.OrderId);
    }
}

Step 4: Configure MassTransit and RabbitMQ in Program.cs

Wire up MassTransit into your dependency injection pipeline, configuring endpoints, retry policies, and connection configurations to your RabbitMQ broker.

C#

using MassTransit;
using SharedContracts;

var builder = WebApplication.CreateBuilder(args);

// Configure MassTransit with RabbitMQ transport
builder.Services.AddMassTransit(x =>
{
    // Register the consumer
    x.AddConsumer<OrderSubmittedConsumer>();

    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.Host("localhost", "/", h =>
        {
            h.Username("guest");
            h.Password("guest");
        });

        // Automatically configure endpoints with standard queue naming conventions
        cfg.ConfigureEndpoints(context);

        // Configure resilient retry policy (e.g., retry 3 times with interval)
        cfg.UseMessageRetry(r => r.Interval(3, TimeSpan.FromSeconds(5)));
    });
});

builder.Services.AddControllers();
var app = builder.Build();

app.MapControllers();
app.Run();

Step 5: Publishing Messages from an API Endpoint

Inject IPublishEndpoint or ISendEndpointProvider into your API controllers to publish events onto the message bus asynchronously.

C#

using MassTransit;
using Microsoft.AspNetCore.Mvc;
using SharedContracts;

[ApiController]
[Route("api/orders")]
public class OrdersController : ControllerBase
{
    private readonly IPublishEndpoint _publishEndpoint;

    public OrdersController(IPublishEndpoint publishEndpoint)
    {
        _publishEndpoint = publishEndpoint;
    }

    [HttpPost]
    public async Task<IActionResult> SubmitOrder([FromBody] SubmitOrderModel model)
    {
        var orderId = Guid.NewGuid();

        // Construct the integration event
        var orderEvent = new OrderSubmittedEvent(
            OrderId: orderId,
            CustomerEmail: model.CustomerEmail,
            TotalAmount: model.TotalAmount,
            Timestamp: DateTime.UtcNow
        );

        // Publish event to the message broker exchange; decoupled from direct handlers
        await _publishEndpoint.Publish(orderEvent);

        return Accepted(new { OrderId = orderId, Status = "Order submitted and queued for processing." });
    }
}

public record SubmitOrderModel(string CustomerEmail, decimal TotalAmount);