Mass Transit RabbitMQ Consumer fails with long task

Viewed 453

Error Message: MassTransit.ConnectionException: The connection is stopping and cannot be used: rabbit-host

We have a long-running consumer using MassTransit with RabbitMQ. We are failing in the consumer when trying to publish our result onto a different queue after running for 20+ minutes.

We are assuming that the connection is timing out before we complete our work.

We see there is an option to use a JobConsumer for long running tasks, but were wondering if there is a way to extend our timeout on the regular Consumer when working with RabbitMQ?

We saw the MaxAutoRenewDuration option when working with Azure on this question: Masstransit - long running process and imediate response and were looking for something similar for RabbitMQ.

Is there a specific default timeout time that the connection to the rabbit host lasts?

Thank you for any help, and I can provide more details if that'd be helpful.

// Consumer Class
using System;
using System.Threading.Tasks;
using MassTransit;

namespace ExampleNameSpace
{
    public class ExampleConsumer : IConsumer<ExampleMessage>
    {
        private readonly BusinessLogicProcess _businessLogicProcess;

        public ExampleConsumer(BusinessLogicProcess businessLogicProcess)
        {
            _businessLogicProcess = businessLogicProcess;
        }

        public async Task Consume(ConsumeContext<ExampleMessage> context)
        {
            try
            {
                // Do our business logic (takes 20+ minutes)
                var businessLogicResult = await _businessLogicProcess.DoBusinessWorkAsync();
                var resultMessage = new ResultMessage { ResultValue = businessLogicResult };

                // Publish the result of our work
                // Get the following error when we publish after a long running piece of work
                // MassTransit.ConnectionException: The connection is stopping and cannot be used: rabbitmqs://our-rabbit-host/vhost-name
                await context.Publish<ResultMessage>(resultMessage);
            }
            catch (Exception e)
            {
                Console.WriteLine(e);
                throw;
            }
        }
    }

    public class BusinessLogicProcess {
        
        public async Task<int> DoBusinessWorkAsync()
        {
            // Business Logic Here 
            // Takes 20+ minutes

            return 0;
        }
    }

    public class ExampleMessage { }

    public class ResultMessage {
        public int ResultValue { get; set; }
    }

}

// Service Class - Connect to RabbitMQ
using MassTransit;
using Microsoft.Extensions.Hosting;
using System;
using System.Threading;
using System.Threading.Tasks;

namespace ExampleNameSpace
{
    public class ExampleService : BackgroundService
    {
        private readonly BusinessLogicProcess _businessLogicProcess;
        private IBusControl _bus;

        public ExampleService(BusinessLogicProcess process)
        {
            _businessLogicProcess = process;
        }

        public override async Task StartAsync(CancellationToken cancellationToken)
        {
            _bus = Bus.Factory.CreateUsingRabbitMq(
                config =>
                {
                    config.Host(new Uri("rabbitmqs://our-rabbit-host/vhost-name"), hostConfig =>
                    {
                        hostConfig.Username("guest");
                        hostConfig.Password("guest");
                    });

                    config.ReceiveEndpoint("exampleendpoint",
                        endpointConfigurator =>
                        {
                            endpointConfigurator.Consumer(() => new ExampleConsumer(_businessLogicProcess),
                                config => config.UseConcurrentMessageLimit(1));
                        });
                });

            await _bus.StartAsync(cancellationToken);
            await base.StopAsync(cancellationToken);
        }

        protected override Task ExecuteAsync(CancellationToken stoppingToken)
        {
            return Task.CompletedTask;
        }

        public override async Task StopAsync(CancellationToken cancellationToken)
        {
            await _bus.StopAsync(cancellationToken);
            await base.StopAsync(cancellationToken);
        }
    }
}
0 Answers
Related