使用MassTransit连接Confluent Cloud失败问题求助
问题分析
错误日志显示MassTransit正在尝试连接本地RabbitMQ(rabbitmq://localhost/),但你实际要连接的是Confluent Cloud的Kafka集群。问题出在你的配置中同时启用了RabbitMQ总线和Kafka Rider,而你并没有部署本地RabbitMQ服务,导致连接被拒绝。
解决方案
移除RabbitMQ相关的配置代码,只保留Kafka Rider的配置即可。修改后的完整配置如下:
services.AddMassTransit(x => { x.AddRider(rider => { rider.AddProducer<UserEvent>(topicName: "UserCreated"); rider.AddProducer<UserEvent>(topicName: "UserUpdated"); rider.AddProducer<UserEvent>(topicName: "UserDeleted"); rider.AddConsumer<UserCreatedEventConsumer>(); rider.UsingKafka((riderContext, kafkaFactory) => { kafkaFactory.SecurityProtocol = Confluent.Kafka.SecurityProtocol.SaslSsl; kafkaFactory.Host(server: "[hided...].westeurope.azure.confluent.cloud:9092", configureHost => { configureHost.UseSasl(saslConfig => { saslConfig.Mechanism = Confluent.Kafka.SaslMechanism.Plain; saslConfig.Username = "..................."; // 替换为你的Confluent Cloud API Key saslConfig.Password = "..................."; // 替换为你的Confluent Cloud API Secret }); // 可选:添加Confluent Cloud需要的SSL配置 configureHost.UseSsl(sslConfig => { sslConfig.SslEndpointIdentificationAlgorithm = SslEndpointIdentificationAlgorithm.Https; }); }); var consumerConfig = new ConsumerConfig() { GroupId = "dotnet-example-group-1", AutoOffsetReset = AutoOffsetReset.Latest, EnableAutoCommit = false }; kafkaFactory.TopicEndpoint<UserCreatedEvent>( topicName: "UserCreated", consumerConfig, kafkaTopicReceiveEndpointConfig => { kafkaTopicReceiveEndpointConfig.ConfigureConsumer<UserCreatedEventConsumer>(riderContext); }); }); }); });
关键说明
x.UsingRabbitMq(...)是用来配置MassTransit的RabbitMQ总线,你不需要这个,直接删除即可。- 确保Sasl配置中,
Username是Confluent Cloud的API Key,Password是对应的API Secret,这两个值可在Confluent Cloud控制台的集群设置中获取。 - 若仍有连接问题,检查Confluent Cloud的集群端点是否正确,以及你的网络是否允许访问9092端口(SaslSsl端口)。
内容的提问来源于stack exchange,提问作者user20291437
相关产品推荐
相关产品推荐

