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

如何使用Kafka+.Net Core+MassTransit实现请求响应模式

基于.NET Core + Kafka实现请求响应模式的实现方案

方案1:基于MassTransit实现(优先推荐)

MassTransit原生支持请求响应模式,搭配Kafka驱动可快速实现无需手动维护关联标识和响应队列。

前置依赖

安装对应版本适配的NuGet包:

  • MassTransit(.NET Core 2.1/3.1请选择7.x及以下版本,避免版本不兼容)
  • MassTransit.Kafka
  • MassTransit.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 05:15:02