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传输消息。
解决方案
问题根源
- 监听配置冲突:Kafka容器日志显示尝试连接
172.31.0.2:9093,但配置的监听端口仅为9092。confluentinc/cp-kafka镜像默认启用内部通信端口9093,单端口配置导致控制器无法连接broker。 - 通告地址不匹配:若测试代码与Kafka容器不在同一网络,
localhost:9092的通告地址可能无法正确解析,尤其是测试代码也运行在容器中的场景。 - 消费者注册不一致:代码中注册的消费者类型
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
相关产品推荐
相关产品推荐

