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/目录下:
- 请求消息
IRequest<TResponse>(IRequest.cs):一个标记接口,声明"这个请求期待什么类型的回复"。 - 处理器
IRequestHandler<TRequest, TResponse>(IRequestHandler.cs):Worker 端实现OnHandle(),处理请求并返回回复。 - 请求响应总线
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后,框架在幕后做了这些事:
- 自动注入关联头:请求消息携带
RequestId和ReplyTo两个关键消息头(定义在 ReqRespMessageHeaders.cs)——RequestId用于匹配"哪条回复对应哪个请求",ReplyTo告诉 Worker 把回复发到哪里。 - Worker 消费并处理:Worker 从主题中消费请求,调用你注册的
IRequestHandler.OnHandle()。 - 框架自动发回回复:
OnHandle()的返回值由框架序列化后发到ReplyTo指定的主题——你不需要在 Handler 里手动 publish 回复。 - 发送方完成等待:回复到达后,框架按
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),仅供参考