a consumer consumes the same message many times - masstransit

Viewed 310

there are 2 apps which are work as consumers and producers. app1 and app2. both of them are asp.net core web api.

app1 or app2 publish a message and both app1 and app2 consume the message. but sometimes app1 consumes the same message many times ( from 2 to 478) but the same message is consuming by app2 just one time as expected. the consume method is very simple just one select and one insert operation.

at first, I thought after insert operation some exceptions occurred so masstransit enqueued it and tried to consume it again. but there is no code to make exceptions after insert operation. even if there is, we send exceptions to the sentry but there is nothing about exception in sentry.

according to aws sqs logs there are abnormal message counts in a specific period of time

enter image description here

app2's queue is normal

there are my configurations and codes below. so what could cause this?

configurations
enter image description here

app1

masstransit configuration

public static IServiceCollection AddEventBus(this IServiceCollection services, IConfiguration configuration)
{
    services.AddTransient<OrderCompletedIntegrationEvent>();

    var options = configuration.GetSection("AmazonSqs").Get<AmazonSqsOptions>();
    var massOptions = configuration.GetSection("Masstransit").Get<MasstransitOptions>();
    IBusControl CreateBus(IServiceProvider serviceProvider)
    {
        var secretManagerService = serviceProvider.GetService<ISecretManagerService>();
        var secrets = secretManagerService.GetSecrets();
        var bus = Bus.Factory.CreateUsingAmazonSqs(cfg =>
        {
            cfg.Host(options.Region, h =>
            {
                h.AccessKey(secrets.AccessKeyId);
                h.SecretKey(secrets.SecretKey);
            });
            cfg.Durable = true;
            cfg.WaitTimeSeconds = 10;
            cfg.ReceiveEndpoint(options.QueueName, q =>
            {
                // RetryCount: 1
                // RetryWaitIntervalInSeconds: 1
                q.UseMessageRetry(r => r.Interval(massOptions.RetryCount, TimeSpan.FromSeconds(massOptions.RetryWaitIntervalInSeconds)));
                q.ConfigureConsumer<OrderCompletedIntegrationEventHandler>(serviceProvider);

            });
        });                

        return bus;
    }

    // local function to configure consumers
    void ConfigureMassTransit(IServiceCollectionConfigurator configurator)
    {
        configurator.AddConsumer<OrderCompletedIntegrationEventHandler>();
    }

    // configures MassTransit to integrate with the built-in dependency injection
    services.AddMassTransit(configurator =>
    {
        configurator.AddBus(CreateBus);
        ConfigureMassTransit(configurator);
    });
    services.AddMassTransitHostedService();


    return services;
}

consumer

public class OrderCompletedIntegrationEventHandler : IConsumer<OrderCompletedIntegrationEvent>
{
    private readonly MasstransitOptions options;
    private readonly IServiceProvider serviceProvider;
    private readonly ICoinService coinService;
    private readonly IAccountService accountService;
    private readonly ILogger<OrderCompletedIntegrationEventHandler> _logger;

    public OrderCompletedIntegrationEventHandler(
        IConfiguration configuration,
        IServiceProvider serviceProvider,
        ICoinService coinService,
        IAccountService accountService)
    {
        options = configuration.GetSection("Masstransit").Get<MasstransitOptions>();
        this.serviceProvider = serviceProvider;
        this.coinService = coinService;
        this.accountService = accountService;
    }

    public async Task Consume(ConsumeContext<OrderCompletedIntegrationEvent> context)
    {
        var @event = context.Message;
        var retry = context.GetRetryAttempt();
        try
        {
            await Task.Run(() =>
            {
                var coin = coinService.TlToCoin(@event.Total);
                var coinToAdd = coin * 25 / 1000;
                var user = accountService.FindByEmail(@event.Email);
                if (user != null) coinService.AddBonus(user.Id, coinToAdd);
            });
        }
        catch (Exception ex)
        {
            if (retry == options.RetryCount)
            {
                using (var scope = serviceProvider.CreateScope())
                {
                    var hub = scope.ServiceProvider.GetService<IHub>();
                    hub.CaptureException(new IntegrationEventException(data: @event, ex));
                }
            }
            throw ex;
        }
    }
}

app2

masstransit configuration

public static IServiceCollection AddEventBus(this IServiceCollection services, IConfiguration configuration)
{
    services.AddScoped<OrderCompletedIntegrationEventHandler>();

    var options = configuration.GetSection("AmazonSqs").Get<AmazonSqsOptions>();
    var massOptions = configuration.GetSection("Masstransit").Get<MasstransitOptions>();

    IBusControl CreateBus(IServiceProvider serviceProvider)
    {
        return Bus.Factory.CreateUsingAmazonSqs(cfg =>
        {
            cfg.UseAmazonSqsMessageScheduler();

            cfg.Host(options.Region, h =>
            {
                h.AccessKey(options.AccessKeyId);
                h.SecretKey(options.SecretKey);
            });

            cfg.Durable = true;
            cfg.WaitTimeSeconds = 20;

            cfg.ReceiveEndpoint(options.QueueName, q =>
            {
                // RetryCount: 5
                // RetryWaitIntervalInSeconds: 300
                q.UseMessageRetry(r => r.Interval(massOptions.RetryCount, TimeSpan.FromSeconds(massOptions.RetryWaitIntervalInSeconds)));
                q.ConfigureConsumer<OrderCompletedIntegrationEventHandler>(serviceProvider);
            });
        });
    }

    // local function to configure consumers
    void ConfigureMassTransit(IServiceCollectionConfigurator configurator)
    {
        configurator.AddConsumer<OrderCompletedIntegrationEventHandler>();
    }

    //configures MassTransit to integrate with the built -in dependency injection
    services.AddMassTransit(x =>
    {
        x.AddBus(CreateBus);
        ConfigureMassTransit(x);
    });
    services.AddMassTransitHostedService();

    return services;
}

consumer

public class OrderCompletedIntegrationEventHandler : IConsumer<OrderCompletedIntegrationEvent>
{
    private readonly MasstransitOptions options;
    private readonly IServiceProvider serviceProvider;
    private IUserService _userService;
    private IScoreService _scoreService;
    public OrderCompletedIntegrationEventHandler(
        IConfiguration configuration,
        IServiceProvider serviceProvider, IUserService userService, IScoreService scoreService)
    {
        options = configuration.GetSection("Masstransit").Get<MasstransitOptions>();
        this.serviceProvider = serviceProvider;
        this._userService = userService;
        this._scoreService = scoreService;
    }

    public async Task Consume(ConsumeContext<OrderCompletedIntegrationEvent> context)
    {
        var @event = context.Message;
        var retry = context.GetRetryAttempt();
        try
        {
            var user = _userService.GetUserByEmail(@event.Email);
            if (user != null)
            {

                var activity = _scoreService.GetActivityTypeByName("buy-book");
                var scoreToAdd = Convert.ToInt32(@event.Total * 150);
                var coinToAdd = Convert.ToInt32(@event.Total * 100 * 2.5M / 100);
                var scoreModel = new ScoreModel
                {
                    ActivityType = new ActivityType
                    {
                        Id = activity.Id
                    },
                    Name = activity.Name,
                    Description = $"{@event.Total:F2} TL'lik bir sipariş verildi. Sipariş no: {@event.OrderId}",
                    UserId = (int)user.ID,
                    BranchId = null,
                    GainedScore = scoreToAdd,
                    GainedCoins = coinToAdd,
                    CreateDate = DateTime.UtcNow.TurkeyTime()
                };
                await _scoreService.AddActivity(scoreModel);
            }
        }
        catch (System.Exception ex)
        {
            if (retry == options.RetryCount)
            {
                using (var scope = serviceProvider.CreateScope())
                {
                    var hub = scope.ServiceProvider.GetService<IHub>();
                    hub.CaptureException(ex);
                }
            }

            throw ex;
        }
    }
}
0 Answers
Related