Flash.Middleware.ReliableRabbitMQPublisher.Dotnet 1.0.3

dotnet add package Flash.Middleware.ReliableRabbitMQPublisher.Dotnet --version 1.0.3
                    
NuGet\Install-Package Flash.Middleware.ReliableRabbitMQPublisher.Dotnet -Version 1.0.3
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="Flash.Middleware.ReliableRabbitMQPublisher.Dotnet" Version="1.0.3" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Flash.Middleware.ReliableRabbitMQPublisher.Dotnet" Version="1.0.3" />
                    
Directory.Packages.props
<PackageReference Include="Flash.Middleware.ReliableRabbitMQPublisher.Dotnet" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add Flash.Middleware.ReliableRabbitMQPublisher.Dotnet --version 1.0.3
                    
#r "nuget: Flash.Middleware.ReliableRabbitMQPublisher.Dotnet, 1.0.3"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package Flash.Middleware.ReliableRabbitMQPublisher.Dotnet@1.0.3
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=Flash.Middleware.ReliableRabbitMQPublisher.Dotnet&version=1.0.3
                    
Install as a Cake Addin
#tool nuget:?package=Flash.Middleware.ReliableRabbitMQPublisher.Dotnet&version=1.0.3
                    
Install as a Cake Tool

Reliable RabbitMQ Publisher for .NET

High-throughput, backpressure-aware, circuit-breaking, metrics-enabled RabbitMQ publisher for .NET 8+

This library implements a high-performance, adaptive RabbitMQ publisher designed for extreme throughput and reliability.
It automatically handles:

  • Worker pool scaling
  • Backpressure
  • Circuit-breakers
  • Dead-letter storage
  • Structured telemetry & Prometheus metrics
  • Configurable batch publishing
  • Graceful shutdown with backlog preservation
  • Pluggable hooks for failure/crash/backlog events

It is packaged as a reusable class library, published to NuGet, and consumed by any .NET app via dependency injection.


✨ Features

🔧 Adaptive Worker Pool

Automatically scales publisher workers between MinWorkers and MaxWorkers depending on backlog pressure.

🧠 Circuit Breaker

When RabbitMQ is unavailable or slow, the publisher opens a circuit to protect the application.

🚦 Backpressure

If the backlog grows too large, TryEnqueueAsync returns Backpressure to protect upstream systems.

📦 Batch Publishing

Messages are batched to reduce broker round-trips and increase throughput.

📝 Dead-Letter Support

Failures are stored using an IDeadLetterStore implementation of your choosing (File, Redis, DB, S3, etc.).

📊 Unified Metrics

Prometheus metrics and internal stats tracking via:

  • IRabbitPublisherMetrics
  • IRabbitPublisherStatsSink

🛠️ Configurable via Builder Pattern

Allows a fully fluent experience during DI registration.

🚀 Hosted Background Service

The publisher starts immediately when the host starts. No manual warm-up required.


📦 Installation: (requires RabbitMQ.Client 6.8.1)

dotnet add package RabbitMQ.Client --version 6.8.1

dotnet add package Flash.Middleware.ReliableRabbitMQPublisher.Dotnet

⚙️ Usage example:

Configuration:

NOTE: You MUST configure either appsettings.json with a "RabbitPublisher" section or in code using helper method:

    rabbit.Configure(settings =>
    {
        settings.ConnectionString = "amqp://rabbitdemo:rabbit@localhost:5672/Master" 
        settings.MinWorkers = 1;
        settings.MaxWorkers = 4;
        settings.BatchSize = 100,
        settings.MaxBacklog = 1000,
        //etc
    })

OR

  "RabbitPublisher": {
    "ConnectionString": "amqp://rabbitdemo:rabbit@localhost:5672/Master",
    "Exchange": "test-exchange",
    "RoutingKey": "",
    "MaxBacklog": 1000,
    "MinWorkers": 1,
    "MaxWorkers": 4,
    "BatchSize": 100,
    "HighBacklogThreshold": 500,
    "LowBacklogTicksToScaleDown": 15,
    "ScaleIntervalSeconds": 10,
    "ConfirmTimeoutSeconds": 2
  }

Dependency Injection Registration:

// Register the RabbitMQ publisher using the new builder pattern
builder.Services.AddRabbitPublisher((rabbit, sp) =>
{
    rabbit.Configure(settings =>
    {
        settings.MaxWorkers = 4;
        settings.BatchSize = 500;
    });

    rabbit.EnableMetrics();
    rabbit.EnableStats();

    // Register dead letter store
    rabbit.UseDeadLetterStore(sp =>
        new FileDeadLetterStore(
            Path.Combine(AppContext.BaseDirectory, "rabbitmq-publisher", "deadletters", "deadletters.log"))
    );

    rabbit.OnPublishFailed(async (evt, sp) =>
    {
        var logger = sp.GetRequiredService<ILogger<AdaptiveRabbitPublisher>>();
        var dead = sp.GetService<IDeadLetterStore>();
        try
        {
            await dead.StoreAsync(evt.MsgEnvelope, evt.Exception, evt.Reason);
        }
        catch (Exception ex)
        {
            logger.LogError(ex, "Failed to store dead-letter for {MessageId}", evt.MsgEnvelope.MessageId);
        }
    });

    rabbit.OnShutdownBacklog(async (evt, sp) =>
    {
        var logger = sp.GetRequiredService<ILogger<AdaptiveRabbitPublisher>>();
        var dead = sp.GetService<IDeadLetterStore>();
        try
        {
            var batch = evt.BufferedMessages
                .Select(env => (env, (Exception?)null, "Shutdown backlog"));

            await dead.StoreBatchAsync(batch);
        }
        catch (Exception ex)
        {
            logger.LogError(ex, "Failed to store shutdown backlog");
        }
    });

    rabbit.OnWorkerCrashed(async (evt, sp) =>
    {   
        var logger = sp.GetRequiredService<ILogger<AdaptiveRabbitPublisher>>();
        logger.LogError(evt.Exception, "Publisher worker crashed: {Details}", evt.Details);
        await Task.CompletedTask;
    });

    rabbit.OnBatchPublished(async (evt, sp) =>
    {
        await Task.CompletedTask;
    });
});

Publishing messages:

Use Dependency injection to resolve the publisher instance and enqueue messages, will return 'PublishResult' immediately:

see Model.cs in the library:

/// <summary>
/// The message publish result.
/// </summary>
/// <param name="MessageId"></param>
/// <param name="Status"></param>
public sealed record PublishResult(String MessageId, PublishStatus Status, String? Details = null)
{
    public static PublishResult Create(String messageId, PublishStatus status, String? details = null)
        => new(messageId, status, details);
}
PublishStatus.cs:
/// <summary>
/// Publish status when enqueing message.
/// </summary>
public enum PublishStatus
{
    /// <summary>
    /// When the message is accepted into the processing backlog.
    /// </summary>
    Enqueued,
    
    /// <summary>
    /// When the backlog is currently full and no new messages are accepted into the backlog.
    /// </summary>
    Backpressure,
    
    /// <summary>
    /// When the circuit break is currently isolated due to network or publishing failures.
    /// </summary>
    CircuitOpen,
    
    /// <summary>
    /// When the publisher being stopped and closing down workers.
    /// </summary>
    Disposed,
    
    /// <summary>
    /// If there's an actual internal error in the publisher.
    /// </summary>
    Error
}

Enqueue message example with builder:

// High-perf publish endpoint
app.MapPost("/publish", async (HttpRequest req, AdaptiveRabbitPublisher publisher, CancellationToken ct) =>
{
    var body = await ThreadLocalBodyReader.ReadAsync(req, ct);

    if (body.Length == 0)
        return Results.BadRequest("Empty body");

    // Create the message envelope. (example)
    var messageId = Guid.NewGuid().ToString();
    var traceId = Guid.NewGuid().ToString();
    var envelope = new PublishEnvelopeBuilder()
        .WithPayload(body) // Required*
        .WithMessageId(messageId) // Optional message ID, generates internal message ID when not supplied
        .WithHeader("TraceId", traceId) //Optional
        .WithRoutingKey("orders.created") //Optional
        .Build();
    
    // Fast enqueue
    var result = await publisher.TryEnqueueAsync(envelope);

    return result.Status switch
    {
        PublishStatus.Enqueued => Results.Accepted(value: result),
        PublishStatus.Backpressure => Results.StatusCode(429).WithRetryAfter(5),
        PublishStatus.CircuitOpen => Results.StatusCode(503),
        PublishStatus.Disposed => Results.StatusCode(410),
        _ => Results.StatusCode(500)
    };
});


Contact details (if changes, assistance required):

waldu.coetzee@flash.co.za
torsten.genis@flash.co.za

Change log:

Version: 1.0.3:

High performance hot path improvements to over 75k messages per second
Optimize conditional path for adding message headers
Remove unnecessary headers for publishing throughput
Socket options added to enable NAGLE's algorithm for TCP socket packet batching
Refactor of Worker.cs to make it much more performant.
Version: 1.0.2:

Add message envelope builder (see: PublishEnvelopeBuilder) to allow publishers to add headers
+ override the default or configured routingkey per message.
Product Compatible and additional computed target framework versions.
.NET net8.0 is compatible.  net8.0-android was computed.  net8.0-browser was computed.  net8.0-ios was computed.  net8.0-maccatalyst was computed.  net8.0-macos was computed.  net8.0-tvos was computed.  net8.0-windows was computed.  net9.0 was computed.  net9.0-android was computed.  net9.0-browser was computed.  net9.0-ios was computed.  net9.0-maccatalyst was computed.  net9.0-macos was computed.  net9.0-tvos was computed.  net9.0-windows was computed.  net10.0 was computed.  net10.0-android was computed.  net10.0-browser was computed.  net10.0-ios was computed.  net10.0-maccatalyst was computed.  net10.0-macos was computed.  net10.0-tvos was computed.  net10.0-windows was computed. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

NuGet packages

This package is not used by any NuGet packages.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
1.0.3 1,698 12/19/2025