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

使用MassTransit多总线Kafka Rider消费多主机消息异常求助

解决MassTransit多总线Kafka消费仅单实例生效的问题

本地Docker部署两个Kafka实例(地址分别为localhost:9094、localhost:9095),通过MassTransit创建两个独立总线分别对应消费,但仅localhost:9094能正常消费消息,localhost:9095无消费行为且无任何报错信息。

问题根源

代码中仅启动了第一个默认总线(IBusControl),而标记为ISecondBus的第二个总线未执行启动操作,导致其关联的Kafka Rider从未初始化运行,自然无法消费对应Kafka实例的消息。

修复后的代码

namespace KafkaConsumer;
using MassTransit;
using MassTransit.Kafka.Consumer.Consumers;
using MassTransit.Kafka.Contracts;
using Microsoft.Extensions.DependencyInjection;

public class Program
{
    public static async Task Main()
    {
        var services = new ServiceCollection();
        services.AddMassTransit(x =>
        {
            const string topicName = "test_topic";
            const string consumerGroup = "consumer-group-medium";
            const string kafkaBrokerServers = "localhost:9094";
            
            x.UsingInMemory((context, cfg) =>
            {
                cfg.ConfigureEndpoints(context);
            });

            x.AddRider(rider =>
            {
                rider.AddConsumer<KafkaMessageConsumer>();
                rider.UsingKafka((context, k) =>
                {
                    k.Host(kafkaBrokerServers);

                    k.TopicEndpoint<IMessage>(topicName, consumerGroup, e =>
                    {
                        e.ConfigureConsumer<KafkaMessageConsumer>(context);
                        e.CreateIfMissing();
                    });
                });
            });
        });
        
        services.AddMassTransit<ISecondBus>(x =>
        {
            const string topicName1 = "test_topic";
            const string consumerGroup1 = "consumer-group-medium";
            const string kafkaBrokerServers1 = "localhost:9095";
            
            x.UsingInMemory((context, cfg) =>
            {
                cfg.ConfigureEndpoints(context);
            });
            
            x.AddConsumer<KafkaMessageConsumer1>();

            x.AddRider(rider =>
            {
                rider.AddConsumer<KafkaMessageConsumer1>();
                rider.UsingKafka((context, k) =>
                {
                    k.Host(kafkaBrokerServers1);
                    k.TopicEndpoint<IMessage>(topicName1, consumerGroup1, e =>
                    {
                        e.ConfigureConsumer<KafkaMessageConsumer1>(context);
                        e.CreateIfMissing();
                    });
                });
            });
        });
        
        services.AddScoped<KafkaMessageConsumer>();
        services.AddScoped<KafkaMessageConsumer1>();
        
        var provider = services.BuildServiceProvider();
        var cancellationToken = new CancellationTokenSource(TimeSpan.FromSeconds(60)).Token;
        
        // 启动第一个总线
        var busControl = provider.GetRequiredService<IBusControl>();
        await busControl.StartAsync(cancellationToken);
        
        // 启动第二个总线(关键:必须显式启动命名总线)
        var secondBusControl = provider.GetRequiredService<IBusControl, ISecondBus>();
        await secondBusControl.StartAsync(cancellationToken);
        
        Console.WriteLine("Started...");
        Console.ReadKey();
        
        // 停止两个总线,释放资源
        await busControl.StopAsync(cancellationToken);
        await secondBusControl.StopAsync(cancellationToken);
    }
}

// 确保ISecondBus接口已定义(如果未定义需添加)
public interface ISecondBus : IBus
{
}

关键修改说明

  1. 启动第二个总线:通过provider.GetRequiredService<IBusControl, ISecondBus>()获取命名总线的控制实例,并调用StartAsync启动它,这是让第二个Kafka Rider运行的核心操作。
  2. 添加总线停止逻辑:程序退出时停止两个总线,避免Kafka消费者连接等资源泄漏。
  3. 确保ISecondBus接口存在:需定义ISecondBus接口继承自IBus,作为第二个总线的标识。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 10:05:19