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

能否使用MassTransit实现RabbitMQ工作队列场景?请提供示例

如何用MassTransit实现RabbitMQ工作队列场景

当然可以!其实MassTransit本身就完美支持RabbitMQ的工作队列模式——只是官方文档里没直接用“工作队列”这个词来称呼它。本质上就是让多个消费者实例竞争消费同一个队列的消息,实现任务的负载均衡和故障转移,和RabbitMQ官网描述的场景完全匹配。

核心原理

工作队列的核心逻辑是:多个消费者监听同一个队列,RabbitMQ会自动将消息分发给空闲的消费者(避免重复处理),如果某个消费者崩溃,未完成的消息会重新回到队列,由其他消费者接手。MassTransit通过显式指定队列名称、配置消费者并发数来实现这个逻辑。

示例代码

1. 生产者:发送任务消息到指定队列

生产者需要将消息直接发送到目标队列(而非发布到交换器广播),确保消息只会进入工作队列等待消费:

using MassTransit;
using System;
using System.Threading.Tasks;

// 定义工作任务的消息契约
public record JobTask
{
    public Guid TaskId { get; init; }
    public string TaskContent { get; init; }
}

class JobProducer
{
    static async Task Main(string[] args)
    {
        // 创建MassTransit总线
        var bus = Bus.Factory.CreateUsingRabbitMq(cfg =>
        {
            var host = cfg.Host(new Uri("rabbitmq://localhost"), h =>
            {
                h.Username("guest");
                h.Password("guest");
            });

            // 显式指定消息对应的队列名称(确保消费者能监听同一个队列)
            cfg.Message<JobTask>(x => x.SetEntityName("worker-job-queue"));
            cfg.Publish<JobTask>(x => x.ExchangeType = "direct");
        });

        await bus.StartAsync();
        try
        {
            // 发送10个测试任务消息
            for (int i = 0; i < 10; i++)
            {
                await bus.Send<JobTask>(new 
                { 
                    TaskId = Guid.NewGuid(), 
                    TaskContent = $"待处理任务 #{i+1}" 
                });
                Console.WriteLine($"已发送任务: 任务 #{i+1}");
                await Task.Delay(500);
            }
        }
        finally
        {
            await bus.StopAsync();
        }
    }
}

2. 消费者:监听工作队列处理任务

你可以启动多个相同的消费者实例,它们会自动竞争队列中的消息:

using MassTransit;
using System;
using System.Diagnostics;
using System.Threading.Tasks;

class JobConsumer : IConsumer<JobTask>
{
    public async Task Consume(ConsumeContext<JobTask> context)
    {
        var processId = Process.GetCurrentProcess().Id;
        Console.WriteLine($"[{DateTime.Now:HH:mm:ss}] 消费者实例 {processId} 开始处理任务: {context.Message.TaskContent}");
        
        // 模拟任务处理耗时(比如业务逻辑执行)
        await Task.Delay(TimeSpan.FromSeconds(3));
        
        Console.WriteLine($"[{DateTime.Now:HH:mm:ss}] 消费者实例 {processId} 完成任务: {context.Message.TaskContent}");
    }
}

class WorkerService
{
    static async Task Main(string[] args)
    {
        var bus = Bus.Factory.CreateUsingRabbitMq(cfg =>
        {
            var host = cfg.Host(new Uri("rabbitmq://localhost"), h =>
            {
                h.Username("guest");
                h.Password("guest");
            });

            // 绑定到工作队列,确保所有消费者监听同一个队列名
            cfg.ReceiveEndpoint(host, "worker-job-queue", e =>
            {
                // 配置单个消费者实例同时处理的消息数量(控制负载)
                e.ConcurrentMessageLimit = 2;
                
                // 注册任务消费者
                e.Consumer<JobConsumer>();
            });
        });

        await bus.StartAsync();
        Console.WriteLine("消费者已启动,等待任务消息...");
        Console.ReadLine();
        await bus.StopAsync();
    }
}

关键配置说明

  • 指定队列名称:生产者和消费者必须使用完全相同的队列名(示例中为worker-job-queue),这样多个消费者才能监听同一个队列。
  • 使用Send而非Publish:Send方法会将消息直接发送到指定队列,而Publish是发布订阅模式的广播行为,不符合工作队列的点对点需求。
  • 并发数控制:ConcurrentMessageLimit可以设置单个消费者实例同时处理的消息数量,根据服务器性能调整这个值,避免过载。
  • 故障自动重试:如果消费者处理消息时抛出异常,MassTransit默认会将消息重新放回队列(可配置重试策略),确保任务不会丢失。

内容的提问来源于stack exchange,提问作者carstenj

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 11:04:55