В моей настройке MassTransit + RabbitMQ у меня есть издатель, который отправляет события и подписывается на них, а также отдельное потребительское приложение. Когда потребитель не работает, а издатель отправляет и обрабатывает событие, сообщение не доставляется потребителю после того, как оно возвращается в режим онлайн. Кажется, RabbitMQ считает сообщение уже обработанным, поскольку издатель его использовал.
Важное примечание: Когда оба приложения запущены, доставка сообщения работает правильно, и оба приложения использовать сообщения, как ожидалось. Проблема возникает только в сценарии восстановления, описанном выше.
Ожидаемое поведение: сообщения должны доставляться всем подписанным потребителям, даже если один потребитель (в данном случае издатель) уже обработал его. Когда клиентское приложение восстанавливается, оно должно получать сообщения, опубликованные во время его простоя.
Конфигурация одинакова для обоих приложений:
using System.Reflection;
using MassTransit;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using SharedKernel.Communication.Events;
using SharedKernel.Utils;
namespace SharedKernel.Communication.Extensions;
public static class MassTransitServicesExtensions
{
public static IServiceCollection AddMassTransit(
this IServiceCollection services,
IConfiguration config,
Assembly consumersAssembly,
string applicationName
)
{
services.AddMassTransit(busConfigurator =>
{
busConfigurator.SetKebabCaseEndpointNameFormatter();
busConfigurator.AddConsumers(consumersAssembly);
busConfigurator.UsingRabbitMq(
(context, configurator) =>
{
string host = config["RabbitMq:Host"]!;
string username = config["RabbitMq:Username"]!;
string password = config["RabbitMq:Password"]!;
if (AppEnv.IsProduction)
{
host = Environment.GetEnvironmentVariable("RABBITMQ_HOST")!;
username = Environment.GetEnvironmentVariable("RABBITMQ_USER")!;
password = Environment.GetEnvironmentVariable("RABBITMQ_PASSWORD")!;
}
configurator.Host(
new Uri(host),
hostConfigurator =>
{
hostConfigurator.Username(username);
hostConfigurator.Password(password);
}
);
configurator.UseMessageRetry(r => r.Interval(5, TimeSpan.FromSeconds(10)));
configurator.Message(e =>
{
e.SetEntityName("user-confirmed-email-event");
});
configurator.Publish(e =>
{
e.ExchangeType = "fanout";
e.Durable = true;
});
configurator.ReceiveEndpoint(
$"user-confirmed-email-{applicationName}",
e =>
{
e.Durable = true;
e.AutoDelete = false;
e.ConfigureConsumers(context);
e.Bind(
"user-confirmed-email-event",
b =>
{
b.ExchangeType = "fanout";
b.Durable = true;
}
);
}
);
}
);
});
return services;
}
}
Подробнее здесь: https://stackoverflow.com/questions/791 ... -publisher