Curiosus.RequestProcessing.RabbitMQ 3.0.0

dotnet add package Curiosus.RequestProcessing.RabbitMQ --version 3.0.0
                    
NuGet\Install-Package Curiosus.RequestProcessing.RabbitMQ -Version 3.0.0
                    
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="Curiosus.RequestProcessing.RabbitMQ" Version="3.0.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Curiosus.RequestProcessing.RabbitMQ" Version="3.0.0" />
                    
Directory.Packages.props
<PackageReference Include="Curiosus.RequestProcessing.RabbitMQ" />
                    
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 Curiosus.RequestProcessing.RabbitMQ --version 3.0.0
                    
#r "nuget: Curiosus.RequestProcessing.RabbitMQ, 3.0.0"
                    
#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 Curiosus.RequestProcessing.RabbitMQ@3.0.0
                    
#: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=Curiosus.RequestProcessing.RabbitMQ&version=3.0.0
                    
Install as a Cake Addin
#tool nuget:?package=Curiosus.RequestProcessing.RabbitMQ&version=3.0.0
                    
Install as a Cake Tool

Curiosus.RequestProcessing.RabbitMQ

NuGet Downloads Coverage

RabbitMQ transport for Curiosus.RequestProcessing: messages from a durable RabbitMQ queue are dispatched to a pool of workers and acknowledged after successful processing or rejected after a failure.

Installation

dotnet add package Curiosus.RequestProcessing.RabbitMQ

Usage

Node options must implement IRabbitMQRequestProcessorNodeOptions:

public class MyNodeOptions : RequestProcessorNodeOptions, IRabbitMQRequestProcessorNodeOptions
{
    public RabbitMQEventReceiverOptions RabbitMQEventReceiver { get; } = new();
}
RequestProcessor:
  Name: my-consumer
  WorkersCount: 4
  RabbitMQEventReceiver:
    HostName: localhost    # default localhost
    Port: 5672             # default 5672
    UserName: guest        # required
    Password: guest        # required
    ClientName: my-app     # connection name prefix, machine name by default
    QueueName: requests    # required, declared as durable on start
    QosMultiplier: 1       # prefetch = WorkersCount * QosMultiplier, 1..20

The dispatcher turns received messages (ReceivedEvents) into RabbitMQRequestWrapper<TRequest>; the base class confirms or rejects each message in HandleRequestProcessingCompletionAsync:

public class MyDispatcher : RabbitMQRequestDispatcherBase<
    MyRequest, MyWorker, WorkerBasicExtraParams, MyProcessingInfo, MyNodeOptions>
{
    public MyDispatcher(
        MyNodeOptions nodeOptions,
        EventWaitHandle manualResetEvent,
        IReadOnlyList<MyWorker> workers,
        ILogger logger,
        ConcurrentQueue<RabbitMQEvent> receivedEvents)
        : base(nodeOptions, manualResetEvent, workers, logger, receivedEvents)
    {
    }

    protected override Task<IReadOnlyList<RabbitMQRequestWrapper<MyRequest>>?> GetRequestsAsync(
        int maxRequestsCount,
        CancellationToken cancellationToken = default)
    {
        var result = new List<RabbitMQRequestWrapper<MyRequest>>(maxRequestsCount);
        while (result.Count < maxRequestsCount && ReceivedEvents.TryDequeue(out var rabbitMQEvent))
        {
            var request = JsonSerializer.Deserialize<MyRequest>(rabbitMQEvent.Payload)!;
            result.Add(new RabbitMQRequestWrapper<MyRequest>(request.Id, CultureInfo.InvariantCulture, request, rabbitMQEvent));
        }

        return Task.FromResult<IReadOnlyList<RabbitMQRequestWrapper<MyRequest>>?>(result);
    }
}

The bootstrapper derives from RabbitMQRequestProcessorBootstrapperBase<...> and implements CreateDispatcher (pass RabbitMQReceivedEvents to the dispatcher) and CreateWorkerParams. The worker derives from WorkerBase<RabbitMQRequestWrapper<MyRequest>, ...> and reads the message from SourceRequest.

services.AddRabbitMQRequestProcessor<
    MyRequest, MyWorker, WorkerBasicExtraParams, MyBootstrapper, MyNodeOptions, MyDispatcher, MyProcessingInfo>(
    nodeOptions);

RabbitMQEvent.Payload is a copy of the message body. RabbitMQEvent.ReceivedData is the BasicDeliverEventArgs of RabbitMQ.Client 7: read BasicProperties (for example, CorrelationId) and DeliveryTag from it, but not Body, which is valid only while the message is being delivered.

The receiver uses RabbitMQ automatic recovery and additionally reconnects manually (up to 10 attempts) after channel failures such as consumer timeouts. Rejected messages are not requeued. StopAsync sends the confirmations and rejections made before it and closes the connection; DisposeAsync stops the receiver if it was not stopped.

A complete consumer and producer is in samples/RequestProcessing/RabbitMQ.

See also

Product Compatible and additional computed target framework versions.
.NET net9.0 is compatible.  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 is compatible.  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
3.0.0 42 10/3/2026
2.0.0 63 9/28/2026