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
<PackageReference Include="Flash.Middleware.ReliableRabbitMQPublisher.Dotnet" Version="1.0.3" />
<PackageVersion Include="Flash.Middleware.ReliableRabbitMQPublisher.Dotnet" Version="1.0.3" />
<PackageReference Include="Flash.Middleware.ReliableRabbitMQPublisher.Dotnet" />
paket add Flash.Middleware.ReliableRabbitMQPublisher.Dotnet --version 1.0.3
#r "nuget: Flash.Middleware.ReliableRabbitMQPublisher.Dotnet, 1.0.3"
#:package Flash.Middleware.ReliableRabbitMQPublisher.Dotnet@1.0.3
#addin nuget:?package=Flash.Middleware.ReliableRabbitMQPublisher.Dotnet&version=1.0.3
#tool nuget:?package=Flash.Middleware.ReliableRabbitMQPublisher.Dotnet&version=1.0.3
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:
IRabbitPublisherMetricsIRabbitPublisherStatsSink
🛠️ 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 | Versions 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. |
-
net8.0
- Microsoft.Extensions.Configuration.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Configuration.Binder (>= 10.0.0)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Hosting.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Options (>= 10.0.0)
- Microsoft.Extensions.Options.ConfigurationExtensions (>= 10.0.0)
- Polly (>= 8.0.0)
- prometheus-net (>= 8.2.1)
- RabbitMQ.Client (>= 6.8.1)
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 |