You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RabbitMQ+MassTransit架构中消费者恢复后丢失离线期间消息的问题求助

RabbitMQ+MassTransit架构中消费者恢复后丢失离线期间消息的问题求助

我现在在搭建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<UserConfirmedEmailEvent>(e =>
                    {
                        e.SetEntityName("user-confirmed-email-event");
                    });
                    configurator.Publish<UserConfirmedEmailEvent>(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;
    }
}

备注:内容来源于stack exchange,提问作者Egor Bobrov

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.15 09:50:29