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

MassTransit Kafka消费者控制台应用反复无法连接Kafka Topic端点

问题:MassTransit Kafka消费者无法连接Testcontainers Kafka容器

我正在用MassTransit开发Kafka消费者,同时用Testcontainers的Kafka容器做集成测试,但消费者始终连不上容器的Topic,收不到消息。试过各种端口配置都没用。

MassTransit配置

builder.ConfigureServices((hostContext, services) =>
{
    var topicName = configuration["Kafka:TopicName"];
    var consumerGroupName = configuration["Kafka:ConsumerGroupName"];

    services.AddLogging(c => c.AddConsole());
    services.AddSingleton(configuration);
    services.AddMediatR(cfg => cfg.RegisterServicesFromAssemblyContaining<Program>());

    services.AddMassTransit(x =>
    {
        
        x.UsingInMemory((context, config) => config.ConfigureEndpoints(context));
        x.AddRider(rider =>
        {
            rider.AddConsumer<LocationsMessageConsumer>();
            rider.UsingKafka((context, k) =>
            {
                k.Host(configuration["Kafka:Endpoint"]);
                k.TopicEndpoint<Null, string>(topicName, consumerGroupName, e =>
                {
                    e.ConfigureConsumer<LocationsMessageConsumer>(context);
                });
            });

        });
    });

消费者代码

public class KafkaMessageConsumer : IConsumer<string>
{
    public virtual Task Consume(ConsumeContext<string> context)
    {
        return Task.CompletedTask;
    }
}

Kafka容器配置

public async Task<KafkaContainer> SetupKafkaContainer(INetwork network)
{
    KafkaContainer kafkaContainer = new KafkaBuilder()
        .WithImage("confluentinc/cp-kafka:latest")
        .WithNetwork(network)
        .WithName(_kafkaContainerName)
        .WithPortBinding(_kafkaPort, 9092)
        .WithEnvironment("ALLOW_PLAINTEXT_LISTENER", "yes")
        .WithEnvironment("KAFKA_CFG_LISTENERS", $"PLAINTEXT://:9092")
        .WithEnvironment("KAFKA_CFG_ADVERTISED_LISTENERS", $"PLAINTEXT://localhost:9092")
        .Build();

    await kafkaContainer.StartAsync();

    return kafkaContainer;
}

错误信息

消费者容器错误日志

warn: MassTransit[0] Consumer [] error (Local_Transport): localhost:9092/bootstrap: Connect to ipv4#127.0.0.1:9092 failed: Connection refused (after 0ms in state CONNECT) on locations 
warn: MassTransit[0] Consumer [] error (Local_AllBrokersDown): 1/1 brokers are down on locations

Kafka容器日志(翻译后)

[2023-09-22 14:49:19,695] DEBUG [分区状态机 controllerId=1] 启动分区状态机,初始状态 -> HashMap() (kafka.controller.ZkPartitionStateMachine) 
[2023-09-22 14:49:19,695] INFO [控制器 id=1] 准备好作为新控制器服务,epoch 为1 (kafka.controller.KafkaController) 
[2023-09-22 14:49:19,695] WARN [请求发送线程 controllerId=1] 控制器1与broker 172.31.0.2:9093(id:1 rack: null)的连接失败 (kafka.controller.RequestSendThread) 
java.io.IOException: 连接172.31.0.2:9093(id:1 rack: null)失败。
	at org.apache.kafka.clients.NetworkClientUtils.awaitReady(NetworkClientUtils.java:70)
	at kafka.controller.RequestSendThread.brokerReady(ControllerChannelManager.scala:296)
	at kafka.controller.RequestSendThread.doWork(ControllerChannelManager.scala:249)
	at org.apache.kafka.server.util.ShutdownableThread.run(ShutdownableThread.java:127)
[2023-09-22 14:49:19,698] INFO [控制器 id=1, 目标BrokerId=1] 客户端请求关闭与节点1的连接 (org.apache.kafka.clients.NetworkClient)

已确认:用非MassTransit的收发器可以正常通过Topic传输消息。


解决方案

问题根源

  1. 监听配置冲突:Kafka容器日志显示尝试连接172.31.0.2:9093,但配置的监听端口仅为9092。confluentinc/cp-kafka镜像默认启用内部通信端口9093,单端口配置导致控制器无法连接broker。
  2. 通告地址不匹配:若测试代码与Kafka容器不在同一网络,localhost:9092的通告地址可能无法正确解析,尤其是测试代码也运行在容器中的场景。
  3. 消费者注册不一致:代码中注册的消费者类型LocationsMessageConsumer与实际定义的KafkaMessageConsumer不匹配。

修复步骤

1. 修正Kafka容器环境变量配置

同时配置外部访问和内部集群通信的监听地址:

public async Task<KafkaContainer> SetupKafkaContainer(INetwork network)
{
    var containerPort = 9092;
    var internalPort = 9093;
    var containerHostname = _kafkaContainerName;

    KafkaContainer kafkaContainer = new KafkaBuilder()
        .WithImage("confluentinc/cp-kafka:latest")
        .WithNetwork(network)
        .WithName(_kafkaContainerName)
        .WithPortBinding(_kafkaPort, containerPort)
        .WithEnvironment("ALLOW_PLAINTEXT_LISTENER", "yes")
        // 同时配置外部和内部监听端口
        .WithEnvironment("KAFKA_CFG_LISTENERS", $"PLAINTEXT://0.0.0.0:{containerPort},PLAINTEXT_INTERNAL://0.0.0.0:{internalPort}")
        // 外部用localhost+映射端口,内部用容器hostname+9093
        .WithEnvironment("KAFKA_CFG_ADVERTISED_LISTENERS", $"PLAINTEXT://localhost:{_kafkaPort},PLAINTEXT_INTERNAL://{containerHostname}:{internalPort}")
        .WithEnvironment("KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP", "PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT")
        .WithEnvironment("KAFKA_CFG_INTER_BROKER_LISTENER_NAME", "PLAINTEXT_INTERNAL")
        .WithEnvironment("KAFKA_CFG_BROKER_ID", "1")
        .WithEnvironment("KAFKA_CFG_OFFSETS_TOPIC_REPLICATION_FACTOR", "1")
        .Build();

    await kafkaContainer.StartAsync();

    return kafkaContainer;
}

2. 确保MassTransit端点配置正确

确认configuration["Kafka:Endpoint"]的值为localhost:{_kafkaPort},与容器端口映射保持一致。

3. 等待Kafka容器完全就绪

Kafka启动需要时间,测试代码中添加等待逻辑,确保broker就绪后再启动消费者:

// 在启动消费者前调用此方法
await WaitForKafkaReady(kafkaContainer.GetBootstrapAddress());

// 实现等待逻辑
private async Task WaitForKafkaReady(string bootstrapServers)
{
    var config = new AdminClientConfig { BootstrapServers = bootstrapServers };
    using var adminClient = new AdminClientBuilder(config).Build();
    
    for (int i = 0; i < 30; i++) // 最多等待30秒
    {
        try
        {
            var metadata = adminClient.GetMetadata(TimeSpan.FromSeconds(1));
            if (metadata.Brokers.Count > 0)
                return;
        }
        catch
        {
            // 忽略连接错误,继续等待
        }
        await Task.Delay(1000);
    }
    throw new TimeoutException("Kafka容器未在规定时间内就绪");
}

4. 修正消费者注册一致性

确保注册的消费者类型与实际定义匹配:

// 修正MassTransit注册代码
rider.AddConsumer<KafkaMessageConsumer>();
// ...
e.ConfigureConsumer<KafkaMessageConsumer>(context);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 09:47:04