SlimMessageBus请求响应模式详解:如何用Send()实现异步跨服务调用并等待回复
2026/8/26 15:47:14 网站建设 项目流程

SlimMessageBus请求响应模式详解:如何用Send()实现异步跨服务调用并等待回复

【免费下载链接】SlimMessageBusLightweight message bus interface for .NET (pub/sub and request-response) with transport plugins for popular message brokers.项目地址: https://gitcode.com/gh_mirrors/sl/SlimMessageBus

SlimMessageBus 是 .NET 生态中一款轻量级消息总线(Message Bus),除了经典的发布/订阅(pub/sub),它同样完整支持请求响应模式(Request-Response):只需一行Send()调用,就能把任务异步委托给远端服务,并自动等待对方回复。相比自己手写"发一条消息、建个字典、轮询回复队列"的 RPC 轮子,SlimMessageBus 把请求关联、超时清理、异常回传全部内建好了。

什么时候该用请求响应模式?

先厘清两种模式的边界,避免用错:

模式语义典型场景
发布/订阅Publish()发出去就不管了,不关心结果领域事件、通知
请求响应Send()发出去并等待回复,像异步版的远程方法调用跨服务查询、委托计算、文件处理

官方使用场景文档 RequestResponse.md 描述的就是经典案例:HTTP 接口收到请求后,把耗时计算(如生成缩略图、跑报表)委托给一组 Worker 服务,Worker 处理完把结果发回,最初的调用方收到结果后继续响应 HTTP 请求。

💡 一句话理解:Send()= 把"同步 RPC 的等待体验"装进"异步消息队列的传输管道"里。

核心三件套:消息、发送方、处理方

整个模式由三个契约接口支撑,都位于src/SlimMessageBus/RequestResponse/目录下:

  1. 请求消息IRequest<TResponse>(IRequest.cs):一个标记接口,声明"这个请求期待什么类型的回复"。
  2. 处理器IRequestHandler<TRequest, TResponse>(IRequestHandler.cs):Worker 端实现OnHandle(),处理请求并返回回复。
  3. 请求响应总线IRequestResponseBus(IRequestResponseBus.cs):调用Send()的入口。

以仓库中的图片缩略图示例(src/Samples/Sample.Images.Messages/GenerateThumbnailRequest.cs)为例,定义请求只需一行继承:

public record GenerateThumbnailRequest : IRequest<GenerateThumbnailResponse> { public string FileId { get; set; } public int Width { get; set; } public int Height { get; set; } }

IRequest<TResponse>的泛型参数就是契约:SlimMessageBus 用它在编译期保证"请求-回复"类型匹配,也用于自动路由。

Send() 的两个重载:有回复 vs 只等处理完成

IRequestResponseBus提供两组Send()(见 IRequestResponseBus.cs):

  • Task<TResponse> Send<TResponse>(...):最常用。发出请求并阻塞(异步等待)直到拿到回复,超时或取消会抛出OperationCanceledException
  • Task Send(IRequest ...):无回复类型的"确认型"调用。await它直到请求被处理完;若 Handler 抛异常,异常会回传给发送方——连"失败了"这个信息都不用手动约定。

调用方代码简洁得惊人(示例见src/Samples/Sample.Images.WebApi/Controllers/ImageController.cs):

private readonly IRequestResponseBus bus; var response = await bus.Send( new GenerateThumbnailRequest { FileId = fileId, Width = 300, Height = 200 });

Send()还支持可选的path(指定主题/队列)、headers(附加消息头)、timeout(覆盖超时)和cancellationToken

一次 Send() 背后发生了什么?

这是新手最容易好奇的部分。发送GenerateThumbnailRequest后,框架在幕后做了这些事:

  1. 自动注入关联头:请求消息携带RequestIdReplyTo两个关键消息头(定义在 ReqRespMessageHeaders.cs)——RequestId用于匹配"哪条回复对应哪个请求",ReplyTo告诉 Worker 把回复发到哪里。
  2. Worker 消费并处理:Worker 从主题中消费请求,调用你注册的IRequestHandler.OnHandle()
  3. 框架自动发回回复OnHandle()的返回值由框架序列化后发到ReplyTo指定的主题——你不需要在 Handler 里手动 publish 回复。
  4. 发送方完成等待:回复到达后,框架按RequestId找到挂起的请求并"点亮"那个Task,你的await随即拿到结果。

回复主题上会流经多种响应消息类型(不同请求类型的回复可能共享同一条回复通道),这正是 SlimMessageBus"一个主题承载多种消息类型"能力的体现:

SlimMessageBus 请求响应模式中同一主题可承载多种请求与响应消息类型

超时与清理同样自动化:PendingRequestManager.cs 会定期扫描挂起的请求,把超时或被取消的请求从等待表中清除并触发OperationCanceledException,保证内存不会因"永远等不到回复的请求"而泄漏。

三步跑通:发送方、Worker 与消息定义

下面以 Kafka 传输为例,给出最小可用的配置思路(完整可运行代码在src/Samples/Sample.Images.WebApi/src/Samples/Sample.Images.Worker/)。

第 1 步:发送方配置Produce+ExpectRequestResponses

services.AddSlimMessageBus(mbb => mbb .Produce<GenerateThumbnailRequest>(x => x.DefaultTopic("thumbnail-generation")) // 请求发到哪个主题 .ExpectRequestResponses(x => { x.ReplyToTopic("webapi-1-response"); // 我的回复请发到这个主题 x.DefaultTimeout(TimeSpan.FromSeconds(30)); // 默认等待超时 }) .WithProviderKafka(cfg => cfg.BrokerList = "localhost:9092") .AddJsonSerializer());

ReplyToTopic是关键:它相当于告诉 Worker"回话地址",因此每个发送方实例都应使用自己的回复主题,避免多个实例互相抢回复。

第 2 步:Worker 注册 Handler

mbb.Handle<GenerateThumbnailRequest, GenerateThumbnailResponse>(s => s.Topic("thumbnail-generation", t => t .WithHandler<GenerateThumbnailRequestHandler>() .KafkaGroup("workers") // 共享消费组 .Instances(3)));

第 3 步:实现 Handler 并直接 return 回复

public class GenerateThumbnailRequestHandler : IRequestHandler<GenerateThumbnailRequest, GenerateThumbnailResponse> { public Task<GenerateThumbnailResponse> OnHandle( GenerateThumbnailRequest request, CancellationToken ct) { // 处理图片…… return Task.FromResult(new GenerateThumbnailResponse { FileId = "thumb-xxx" }); } }

注意第 2 步中的KafkaGroup("workers"):同一消费组内多个 Worker 实例会分摊消息,保证每个请求恰好被一个实例处理、只产生一条回复——这是分布式下"不会收到重复回复"的关键。

生产环境必看的 3 个细节 ⏱️

① 超时策略要分层设置全局用ExpectRequestResponses(...).DefaultTimeout(...)兜底,个别慢请求在Send()时传timeout:覆盖。超时会抛出OperationCanceledException,而不是静默挂起。

② 熟悉异常传播链Send()可能抛出的异常定义在src/SlimMessageBus/Exceptions/:发送失败对应SendMessageBusException/ProducerMessageBusException;而 Worker 端 Handler 抛出的业务异常,会被框架通过Error消息头回传给发送方,包装成RequestHandlerFaultedMessageBusException——你在发送端就能直接catch到对端的失败原因。

③ 回复主题的可见性ReplyTo主题只会被发起请求的那个实例消费,通常无需对外暴露。如果部署多实例,请像示例那样用InstanceId拼出独立的回复主题(如webapi-1-response)。

进阶:用拦截器插桩请求响应流程

SlimMessageBus 的拦截器机制(src/SlimMessageBus.Host.Interceptor/)同样覆盖请求响应流程:生产端的IProducerInterceptor可以在Send()前后记录日志、打点、注入通用消息头,消费端拦截器则可统一处理重试与熔断。

常见问题 FAQ

Q:它和直接用 gRPC/HTTP 远程调用有什么区别?A:传输是异步且可削峰的,Worker 可以水平扩展、可以跨语言(配合序列化插件);代价是多了一跳"回复路由"。适合计算密集、需要弹性扩容的跨服务调用。

Q:Worker 挂了怎么办?A:请求会留在队列中由组内其他实例接管;若最终无人回复,发送方会按DefaultTimeout超时并抛出取消异常,不会永久挂起。

Q:如何调试"等不到回复"?A:优先检查三点:ReplyToTopic是否配置、Worker 的KafkaGroup是否与发送方期望一致、Handler 返回类型是否与IRequest<TResponse>TResponse一致。

小结

SlimMessageBus 的请求响应模式把分布式 RPC 里最繁琐的部分——请求关联、回复路由、超时清理、异常回传——全部内建在消息总线里。你只需要定义一个IRequest<TResponse>消息、注册一个 Handler,然后一行Send()就实现了真正的异步跨服务调用并等待回复。想深入动手,仓库src/Samples/下的图片缩略图示例(WebApi + Worker)是最佳起点。

【免费下载链接】SlimMessageBusLightweight message bus interface for .NET (pub/sub and request-response) with transport plugins for popular message brokers.项目地址: https://gitcode.com/gh_mirrors/sl/SlimMessageBus

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询