如何使用Kafka+.Net Core+MassTransit实现请求响应模式
基于.NET Core + Kafka实现请求响应模式的实现方案
方案1:基于MassTransit实现(优先推荐)
MassTransit原生支持请求响应模式,搭配Kafka驱动可快速实现无需手动维护关联标识和响应队列。
前置依赖
安装对应版本适配的NuGet包:
MassTransit(.NET Core 2.1/3.1请选择7.x及以下版本,避免版本不兼容)MassTransit.KafkaMassTransit.Extensions.DependencyInjection
步骤1:定义消息契约
// 请求消息契约 public interface IOrderQueryRequest { Guid OrderId { get; } } // 响应消息契约 public interface IOrderQueryResponse { Guid OrderId { get; } string OrderStatus { get; } decimal TotalAmount { get; } }
步骤2:响应端(服务端)配置与实现
在Startup.cs的ConfigureServices中注册服务:
services.AddMassTransit(x => { x.UsingInMemory(); x.AddRequestClient<IOrderQueryRequest>(); x.AddRider(rider => { // 注册请求消费者 rider.AddConsumer<OrderQueryRequestConsumer>(); rider.UsingKafka((context, config) => { // 替换为实际Kafka集群地址 config.Host("localhost:9092"); // 配置请求Topic监听 config.TopicEndpoint<IOrderQueryRequest>("order-query-request", "order-query-group", e => { e.ConfigureConsumer<OrderQueryRequestConsumer>(context); }); }); }); }); // 启动MassTransit后台服务 services.AddMassTransitHostedService();
消费者实现代码:
public class OrderQueryRequestConsumer : IConsumer<IOrderQueryRequest> { public async Task Consume(ConsumeContext<IOrderQueryRequest> context) { // 执行业务逻辑,比如查询订单数据库 var orderInfo = await GetOrderInfoAsync(context.Message.OrderId); // 直接返回响应,MassTransit自动处理关联匹配 await context.RespondAsync<IOrderQueryResponse>(new { context.Message.OrderId, orderInfo.OrderStatus, orderInfo.TotalAmount }); } }
步骤3:请求端(客户端)调用
在业务服务/控制器中注入请求客户端直接调用:
public class OrderController : ControllerBase { private readonly IRequestClient<IOrderQueryRequest> _requestClient; public OrderController(IRequestClient<IOrderQueryRequest> requestClient) { _requestClient = requestClient; } [HttpGet("{orderId}")] public async Task<IActionResult> GetOrderInfo(Guid orderId) { try { // 发送请求并等待响应,默认超时时间30秒,可通过参数自定义 var response = await _requestClient.GetResponse<IOrderQueryResponse>(new { OrderId = orderId }); return Ok(response.Message); } catch (RequestTimeoutException) { return StatusCode(504, "请求处理超时"); } } }
注意事项
- MassTransit自动维护*关联ID(CorrelationId)*和响应路由,无需手动创建响应Topic
- .NET Core 2.x版本需要额外适配依赖包版本,避免出现API不兼容问题
方案2:基于Confluent.Kafka原生实现
如果不使用MassTransit,可通过原生Kafka客户端手动实现请求响应逻辑。
前置依赖
安装NuGet包:Confluent.Kafka(选择兼容.NET Core版本的发行版)
通用模型定义
public class KafkaRequestMessage { // 请求唯一标识,用于匹配响应 public Guid CorrelationId { get; set; } // 响应回发的Topic名称 public string ReplyTopic { get; set; } // 请求业务数据 public string Payload { get; set; } }
响应端(服务端)实现
public class KafkaResponseHandler { private readonly IConsumer<Null, string> _consumer; private readonly IProducer<Null, string> _producer; public KafkaResponseHandler() { var consumerConfig = new ConsumerConfig { BootstrapServers = "localhost:9092", GroupId = "response-handler-group", AutoOffsetReset = AutoOffsetReset.Earliest }; _consumer = new ConsumerBuilder<Null, string>(consumerConfig).Build(); _consumer.Subscribe("business-request-topic"); var producerConfig = new ProducerConfig { BootstrapServers = "localhost:9092" }; _producer = new ProducerBuilder<Null, string>(producerConfig).Build(); } public async Task StartAsync(CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { var consumeResult = _consumer.Consume(cancellationToken); var request = JsonSerializer.Deserialize<KafkaRequestMessage>(consumeResult.Message.Value); // 处理业务逻辑 var responsePayload = await ProcessBusinessLogic(request.Payload); // 回发响应 await _producer.ProduceAsync(request.ReplyTopic, new Message<Null, string> { Headers = new Headers { {"CorrelationId", Encoding.UTF8.GetBytes(request.CorrelationId.ToString())} }, Value = responsePayload }, cancellationToken); } } }
请求端(客户端)实现
public class KafkaRequestClient { private readonly IProducer<Null, string> _producer; private readonly IConsumer<Null, string> _consumer; private readonly string _replyTopic; private readonly ConcurrentDictionary<Guid, TaskCompletionSource<string>> _pendingRequests = new(); public KafkaRequestClient() { _replyTopic = $"reply-topic-{Guid.NewGuid():N}"; var producerConfig = new ProducerConfig { BootstrapServers = "localhost:9092" }; _producer = new ProducerBuilder<Null, string>(producerConfig).Build(); var consumerConfig = new ConsumerConfig { BootstrapServers = "localhost:9092", GroupId = $"request-client-group-{Guid.NewGuid():N}", AutoOffsetReset = AutoOffsetReset.Earliest }; _consumer = new ConsumerBuilder<Null, string>(consumerConfig).Build(); _consumer.Subscribe(_replyTopic); // 后台监听响应 _ = Task.Run(ListenResponseAsync); } public async Task<string> SendRequestAsync(string payload, int timeoutMs = 30000, CancellationToken cancellationToken = default) { var correlationId = Guid.NewGuid(); var tcs = new TaskCompletionSource<string>(); _pendingRequests.TryAdd(correlationId, tcs); try { await _producer.ProduceAsync("business-request-topic", new Message<Null, string> { Value = JsonSerializer.Serialize(new KafkaRequestMessage { CorrelationId = correlationId, ReplyTopic = _replyTopic, Payload = payload }) }, cancellationToken); using var timeoutCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); timeoutCts.CancelAfter(timeoutMs); using var registration = timeoutCts.Token.Register(() => tcs.TrySetCanceled()); return await tcs.Task; } finally { _pendingRequests.TryRemove(correlationId, out _); } } private async Task ListenResponseAsync() { while (true) { var result = _consumer.Consume(); var correlationIdHeader = result.Message.Headers.First(x => x.Key == "CorrelationId"); var correlationId = Guid.Parse(Encoding.UTF8.GetString(correlationIdHeader.GetValueBytes())); if (_pendingRequests.TryRemove(correlationId, out var tcs)) { tcs.TrySetResult(result.Message.Value); } } } }
内容的提问来源于stack exchange,提问作者KDS
相关产品推荐
相关产品推荐

