
后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载Azure Service Bus 是微软提供的全托管企业级集成消息代理常用于解耦应用与服务并为异步数据与状态传递提供可靠、安全的平台。本文基于开源仓库 dotnetcore/CAP 官方英文文档《Azure Service Bus》系统讲解如何将 Azure Service Bus 接入 CAP 作为消息传输器Transporter覆盖安装步骤、全部配置参数、自定义 Producer、会话Sessions、异构系统兼容以及 SQL 筛选器并结合仓库源码剖析其底层实现原理。读完本文你将能够独立完成 CAP Azure Service Bus 的接入、调优与异构消息对接。Azure Service Bus 在 CAP 中的角色CAP 是一个基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。它通过抽象出ITransport、IConsumerClient等接口将消息的发布与消费委托给具体的消息中间件。Azure Service Bus 正是 CAP 支持的传输器之一负责承载 CAP 的 Outbox 消息出入站通道。仓库中 AzureServiceBusOptionsExtension 展示了接入方式它通过services.AddSingletonIConsumerClientFactory, AzureServiceBusConsumerClientFactory()与services.AddSingletonITransport, AzureServiceBusTransport()将消费者工厂与传输实现注册进 DI 容器从而让 CAP 核心以统一的接口驱动 Azure Service Bus。前提条件与安装定价层级要求!!! warning 重要前提 针对 Service Bus 的定价层级CAP 要求使用Standard标准或Premium高级以支持 Topic主题功能。Basic基础层级不支持 Topic 语义无法满足 CAP 的发布/订阅模型。NuGet 安装要使用 Azure Service Bus 作为消息传输器需要从 NuGet 安装以下扩展包PM Install-Package DotNetCore.CAP.AzureServiceBus快速接入最小配置安装完成后在Startup.cs或 Program.cs 中配置服务集合的位置的ConfigureServices方法中添加配置项public void ConfigureServices(IServiceCollection services) { // ... services.AddCap(x { x.UseAzureServiceBus(opt { //AzureServiceBusOptions }); // x.UseXXX ... }); }UseAzureServiceBus是最小接入入口。真实项目中至少要指定ConnectionString例如仓库测试代码 ServiceBusTransportTests.cs 中使用的连接字符串格式ConnectionString Endpointsb://mynamespace.servicebus.windows.net/;SharedAccessKeyNamemyPolicy;SharedAccessKeymyKey仓库提供的完整示例 Sample.AzureServiceBus.InMemory/Program.cs 演示了从appsettings读取连接字符串并同时启用内存存储与 Dashboard 的典型用法builder.Services.AddCap(c { c.UseInMemoryStorage(); c.UseAzureServiceBus(asb { asb.ConnectionString builder.Configuration.GetConnectionString(AzureServiceBus)!; // ... 自定义标头、SQL 筛选器、自定义 Producer 等 }); c.UseDashboard(); });注意连接字符串中不得包含 Topic 信息CAP 会按TopicPath配置自行创建和管理主题实体。Azure Service Bus 配置参数详解CAP 直接对外提供的 Azure Service Bus 配置参数定义于 CAP.AzureServiceBusOptions.cs如下表所示名称描述类型默认值ConnectionString终端地址Endpoint不得包含 Topic 信息stringAutoProvision自动创建主题、订阅和规则。当 Service Bus 管理 API 不可用时例如使用模拟器 Emulator可关闭此选项booltrueTopicPath主题实体路径stringcapEnableSessions启用 Service Bus 会话SessionsboolfalseMaxConcurrentSessions处理器可处理的最大并发会话数。当 EnableSessions 为 false 时不适用int8SessionIdleTimeout在会话关闭前等待新消息的最长时间。如果未指定Azure Service Bus 将使用 60 秒TimeSpannullSubscriptionAutoDeleteOnIdle在特定空闲间隔后自动删除订阅最小间隔为 5 分钟TimeSpanTimeSpan.MaxValueSubscriptionMessageLockDuration给定接收器锁定消息的时间以防止其他接收器接收相同消息最大值为 5 分钟TimeSpan60 秒SubscriptionDefaultMessageTimeToLive订阅的默认消息生存时间TTL即消息到期前的持续时间TimeSpanTimeSpan.MaxValueSubscriptionMaxDeliveryCount消息在被死信dead-lettered之前的最大传递次数int10MaxAutoLockRenewalDuration锁自动续订的最长持续时间应大于最长的消息锁定持续时间若要无限续订可设Timeout.InfiniteTimeSpanTimeSpan5 分钟TokenCredential用于身份验证的 Azure 凭据如 Managed Identity。使用此选项时还需设置NamespaceTokenCredentialnullMaxConcurrentCalls消息处理程序调用的最大并发数int1AutoCompleteMessages指示处理器在消息处理程序完成处理后是否自动完成消息。若处理程序抛出异常消息不会被自动完成boolfalseCustomHeadersBuilder为来自异构系统的传入消息添加自定义和/或强制性标头FuncServiceBusReceivedMessage, IServiceProvider, ListKeyValuePairstring, string?nullNamespaceService Bus 命名空间在使用 TokenCredential 属性时需要设置stringnullDefaultCorrelationHeaders将附加的关联属性Correlation Properties添加到所有关联筛选器IDictionarystring, stringDictionarystring, string.EmptySQLFilters在主题订阅上按名称和表达式定义的自定义 SQL 筛选器ListKeyValuePairstring, stringnull参数背后的源码逻辑对照源码可以更准确地理解这些参数的实际作用TopicPath 默认值AzureServiceBusOptions中定义了常量DefaultTopicPath capCAP.AzureServiceBusOptions.cs未显式配置时 CAP 将使用名为cap的主题。AutoProvision 与订阅属性当AutoProvision true时消费者客户端在 AzureServiceBusConsumerClient.ConnectAsync 中通过ServiceBusAdministrationClient自动创建 Topic 与 Subscription并把SubscriptionAutoDeleteOnIdle、SubscriptionMessageLockDuration、SubscriptionDefaultMessageTimeToLive、SubscriptionMaxDeliveryCount以及RequiresSession EnableSessions逐一映射到CreateSubscriptionOptions上。处理器参数映射MaxConcurrentCalls在非会话模式下映射为ServiceBusProcessorOptions.MaxConcurrentCalls在会话模式下则映射为ServiceBusSessionProcessorOptions.MaxConcurrentCallsPerSession并与MaxConcurrentSessions、SessionIdleTimeout一并生效见 AzureServiceBusConsumerClient.cs。两种认证方式设置TokenCredential后CAP 会使用new ServiceBusClient(Namespace, TokenCredential)与new ServiceBusAdministrationClient(Namespace, TokenCredential)从而支持 Azure Active Directory / Managed Identity 认证否则退化为连接字符串认证AzureServiceBusConsumerClient.cs。Commit/Reject 语义CAP 默认AutoCompleteMessages false此时消费成功由 CommitAsync 调用CompleteMessageAsync完成消息消费失败则由RejectAsync调用AbandonMessageAsync放回队列这与 CAP 自身的重试/死信机制相配合。自定义 Producer将消息发布到其他 Topic默认情况下CAP 会把所有消息发布到TopicPath即cap指定的主题。使用ConfigureCustomProducerT可以将某个消息名发布到TopicPath以外的主题泛型类型名称必须与Publish时传入的消息名称一致WithSubscription()会在AutoProvision启用时让 CAP 为该主题创建订阅WithSessions()会为此 Producer 发布的消息添加 Session ID如果提供了AzureServiceBusHeaders.SessionId标头就使用其值否则使用 CAP 消息 ID若要从启用了 Session 的订阅中消费还需启用全局EnableSessions选项。services.AddCap(cap cap.UseAzureServiceBus(asb { asb.ConnectionString ...; asb.ConfigureCustomProducerOrderCreated(producer producer.UseTopic(orders).WithSubscription()); })); await capPublisher.PublishAsync(nameof(OrderCreated), new OrderCreated(...));底层实现ConfigureCustomProducerT内部通过ServiceBusProducerDescriptorBuilderT收集配置ServiceBusProducerDescriptorBuilder.cs提供UseTopic(string)、WithSubscription()、WithSessions()三个链式方法最终构建出包含TopicPath、MessageTypeName、CreateSubscription、EnableSessions的ServiceBusProducerDescriptor。发送时AzureServiceBusTransport.CreateProducerForMessage 会先在CustomProducers集合中按消息名MessageTypeName查找匹配的 Producer 描述符找不到时回退到全局TopicPath。仓库测试 ServiceBusTransportTests.cs 分别验证了这两种路径CustomProducer_ShouldHaveCustomTopic配置了ConfigureCustomProducerEntityCreated(cfg cfg.UseTopic(entity-created).WithSubscription())后EntityCreated消息的 Producer TopicPath 为entity-createdDefaultProducer_ShouldHaveDefaultTopic未配置的消息如EntityDeleted回退到全局TopicPath。自动创建订阅时ConnectAsync会把自定义 Producer 的主题路径与默认TopicPath合并去重逐个确保 Topic 与订阅存在AzureServiceBusConsumerClient.cs。示例项目 Program.cs 中同样可以看到两个自定义 Producer 的完整用法。会话Sessions保证消息按顺序处理启用EnableSessions后每个发送的消息都会具有一个 Session ID。要控制 Session ID可在发布消息时通过额外标头AzureServiceBusHeaders.SessionId携带ICapPublisher capBus ...; string yourEventName ...; YourEventType yourEvent ...; Dictionarystring, string extraHeaders new Dictionarystring, string(); extraHeaders.Add(AzureServiceBusHeaders.SessionId, your-session-id); capBus.Publish(yourEventName, yourEvent, extraHeaders);如果头中没有 Session IDCAP 会使用消息 ID 作为 Session ID。源码细节AzureServiceBusHeaders.SessionId的常量值实际为cap-session-id见 AzureServiceBusHeaders.cs。发送端在 ITransport.AzureServiceBus.cs 中当全局EnableSessions或该 Producer 单独启用会话时会先尝试从消息头读取SessionId为空则回退到transportMessage.GetId()CAP 消息 ID。消费端在EnableSessions true时创建ServiceBusSessionProcessorAzureServiceBusConsumerClient.cs并通过ServiceBusProcessorFacade统一暴露会话与非会话两种处理器的消息与错误事件ServiceBusProcessorFacade.cs。会话模式下建议同时关注MaxConcurrentSessions默认 8与SessionIdleTimeout默认交给 Azure 的 60 秒以控制会话并发与空闲回收。会话特性在需要严格 FIFO 顺序消费如订单状态流转、同一业务实体的串行处理的场景下尤其有用。异构系统兼容自定义标头构建器有时你可能需要接收由外部系统发布的消息这些消息没有 CAP 的专用标头。此时需要添加一组两个强制标头以实现 CAP 兼容c.UseAzureServiceBus(asb { asb.ConnectionString ... asb.CustomHeadersBuilder (msg, sp) [ new(DotNetCore.CAP.Messages.Headers.MessageId, sp.GetRequiredServiceISnowflakeId().NextId().ToString()), new(DotNetCore.CAP.Messages.Headers.MessageName, msg.Subject) ]; });其中Headers.MessageId使用ISnowflakeId生成全局唯一消息 ID雪花算法 IDHeaders.MessageName取自msg.Subject即 Azure Service Bus 消息的主题Subject字段它决定 CAP 将该消息路由给哪个订阅方法。重要提示如果消息中已存在同名Key的标头则不会添加自定义标头。源码细节在 AzureServiceBusConsumerClient.ConvertMessage 中CAP 会先把 Service Bus 消息的ApplicationProperties转成标头字典并自动注入Headers.Group等于订阅名/消费者组名随后调用CustomHeadersBuilder返回的标头列表通过headers.TryAdd(key, value)逐个添加——由于使用TryAdd已存在的同名标头不会被覆盖且会记录一条 Warning 日志。示例项目 Program.cs 展示了更完整的自定义标头写法额外加入了一个业务标头IsFromSampleProject并配合下面的 SQL 筛选器一起使用用于验证筛选器是否生效。SQL 筛选器在订阅层面过滤消息可以在订阅层面设置 SQL 筛选器SQL Filters从而只接收想要的消息而无需在业务侧编写额外的过滤逻辑。SQLFilters是一个ListKeyValuePairstring, string其中Key 是规则名称Rule NameValue 是 SQL 表达式c.UseAzureServiceBus(asb { asb.ConnectionString ... asb.SQLFilters new ListKeyValuePairstring, string { new KeyValuePairstring,string(IOTFilter,FromIOTHubtrue), // 当 ApplicationProperties 包含 FromIOTHub 且值为 true 时消息才会被处理 new KeyValuePairstring,string(SequenceFilter,sys.enqueuedSequenceNumber 300) }; });上面的示例中IOTFilter仅当消息的ApplicationProperties中存在FromIOTHub且值为true时才路由到本订阅SequenceFilter基于系统属性sys.enqueuedSequenceNumber入队序号过滤序号大于等于 300 的消息才会被处理。源码细节在AutoProvision开启的情况下SubscribeAsync 会将SQLFilters中的规则名合并进主题列表对比订阅上已存在的规则后增量创建、冗余删除。具体而言若规则名命中SQLFilters中的某项则创建SqlRuleFilter(sqlExpression)SQL 规则否则创建CorrelationRuleFilter其Subject设为规则名消息名并将DefaultCorrelationHeaders中的键值对全部写入关联规则的ApplicationProperties见 AzureServiceBusConsumerClient.cs。由此可见DefaultCorrelationHeaders与SQLFilters协同工作前者为所有关联筛选器统一附加业务属性如环境标识、版本号后者则允许对单个订阅声明独立的 SQL 表达式实现按内容路由。延迟消息与标头常量除会话外AzureServiceBusHeaders还定义了ScheduledEnqueueTimeUtc cap-scheduled-enqueue-time-utcAzureServiceBusHeaders.cs。发送端 ITransport.AzureServiceBus.cs 会检查该标头若存在且可解析为DateTimeOffset则将ServiceBusMessage.ScheduledEnqueueTime设为对应时间从而实现消息延迟发送定时投递可用于延迟任务、定时提醒等场景。从发送到消费的完整调用链结合源码可以完整还原一条消息的生命周期发布业务代码调用ICapPublisher.PublishAsyncCAP 核心将消息持久化到 Outbox 表后由AzureServiceBusTransport.SendAsyncITransport.AzureServiceBus.cs构建ServiceBusMessageMessageId取 CAP 消息 ID、Subject取消息名、CorrelationId取关联 ID全部 CAP 标头写入ApplicationProperties再按需设置SessionId与ScheduledEnqueueTime最终经ServiceBusSender发送。发送失败会包装为PublisherSentFailedException并返回失败的OperateResult。自动创建实体若AutoProvision开启管理客户端会确保默认 Topic 及所有自定义 Producer 的 Topic 与其订阅存在并把订阅级参数映射到 Azure 的订阅描述上。消费AzureServiceBusConsumerClient通过ServiceBusProcessorFacade创建普通处理器或会话处理器并开始监听收到消息后ConvertMessage完成标头转换注入 Group 标头、应用自定义标头随后触发 CAP 的OnMessageCallback。确认/拒绝处理成功后CommitAsync完成消息或依赖AutoCompleteMessages自动完成处理失败则RejectAsync放弃消息交由 CAP 自身的重试机制处理。小结与进一步阅读本文完整覆盖了 CAP 接入 Azure Service Bus 的全部核心内容从定价层级要求、NuGet 安装、最小配置到 20 余项配置参数及其源码映射再到自定义 Producer、会话、异构系统标头、SQL 筛选器与延迟消息。你可以在此基础上结合仓库中的完整示例 Sample.AzureServiceBus.InMemory/Program.cs 与单元测试 ServiceBusTransportTests.cs 进一步验证行为若要深入了解传输器抽象与 CAP 核心机制可继续阅读 ITransport.cs 与 IConsumerClient.cs。赞分享后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载相关推荐MassTransit集成Azure Service Bus消息传输详解MassTransit集成Azure Service Bus消息传输详解 一、Azure Service Bus概述 Azure Service Bus是微软A后端消息队列微服务消息路由CAP消息队列集成全攻略RabbitMQ、Kafka、Azure Service Bus对比选择CAP消息队列集成全攻略RabbitMQ、Kafka、Azure Service Bus对比选择 在微服务架构中 CAP分布式事务解决方案 提供了基于最终一后端消息队列微服务消息路由Azure Service Bus Python SDK 实战指南队列、主题订阅与会话消息全解析Azure Service Bus Python SDK 实战指南队列、主题订阅与会话消息全解析 本指南基于 agentic awesome skills 仓AI 技能AI 插件上一篇Prom-Client 开源项目教程下一篇TheTNB Panel 开源项目使用指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考