资讯动态

CAP 配置完全指南:从 AddCap 基础注册到订阅器与 CapOptions 高级调优

发布时间:2026/9/29 3:06:34 来源:尧图企业网站定制
后端消息队列微服务消息路由【免费下载链接】CAPDistributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern项目地址https://gitcode.com/gh_mirrors/ca/CAP点击查看免费下载CAPDotNetCore.CAP是基于最终一致性思想的分布式事务解决方案同时也是一个采用 Outbox 模式的事件总线。本文聚焦于 CAP 的配置体系从AddCap的注册方式、[CapSubscribe]订阅器参数到CapOptions的每一项核心选项逐一讲解其作用、默认值、源码依据与生产环境中的调优建议。读完本文你将能够独立完成 CAP 服务的最小可用配置并针对重试、并发、消息过期等场景进行精准的参数定制。图片说明无强相关图片不配图本文涉及的配置均为代码级参数仓库图片目录下没有与配置主题直接相关的架构图或运行截屏因此不插入图片以免误导。一、注册入口AddCap 与 DI 容器CAP 的所有配置都从 ASP.NET Core 的依赖注入容器开始。在Program.cs中通过IServiceCollection.AddCap()扩展方法注册并配置 CAP 服务services.AddCap(config { // config.XXX });其中services是Microsoft.Extensions.DependencyInjection包提供的IServiceCollection接口。从源码看AddCap方法src/DotNetCore.CAP/CAP.ServiceCollectionExtensions.cs内部完成了一整套核心服务注册注册消息发布者ICapPublisher、订阅者选择器IConsumerServiceSelector、订阅调用器ISubscribeInvoker与方法匹配缓存MethodMatcherCache注册 4 个后台处理器重试处理器MessageNeedToRetryProcessor、传输检查处理器TransportCheckProcessor、延迟消息处理器MessageDelayedProcessor和过期消息收集器CollectorProcessor注册发送器IMessageSender、默认序列化器JsonUtf8Serializer与分发器IDispatcher创建CapOptions实例并执行你传入的配置委托setupAction随后调用注册到options.Extensions中的存储与传输扩展如UseSqlServer、UseRabbitMQ注册Bootstrapper作为IHostedService在应用启动时完成建表、初始化等引导工作并返回一个支持链式调用的CapBuilder例如继续调用.AddSubscribeFilterT()注册订阅过滤器或.AddSubscriberAssembly(...)指定扫描程序集。1.1 最小必需配置CAP 要求至少配置一个消息传输Transport和一个存储Storage。如果想快速跑通可使用内存传输与内存存储services.AddCap(capOptions { capOptions.UseInMemoryQueue(); // 需要 Savorboard.CAP.InMemoryMessageQueue NuGet 包 capOptions.UseInMemoryStorage(); });不同传输与存储的完整配置选项见仓库文档 Transports 总览 与 Storage 总览例如 RabbitMQ、Kafka、Azure Service Bus、NATS、Redis Streams 等传输以及 SQL Server、PostgreSQL、MySQL、MongoDB、内存存储等存储。二、订阅器配置[CapSubscribe] 的三个核心参数订阅器通过[CapSubscribe]特性标记可位于 ASP.NET Core 的 Controller 或 Service 中。源码层面CapSubscribeAttribute继承自TopicAttributesrc/DotNetCore.CAP/CAP.Attribute.cs、src/DotNetCore.CAP/Internal/TopicAttribute.cs其中TopicAttribute定义了两个订阅配置属性public string Group { get; set; } // 订阅组 public byte GroupConcurrent { get; set; } // 组内并发度2.1 Name必填类型string必填Name参数指定订阅的消息名称与发布方调用_cap.Publish(Name, ...)时的名称一一对应发布接口定义见 src/DotNetCore.CAP/ICapPublisher.cs支持Publish/PublishAsync/PublishDelay等重载。该名称在不同消息中间件中映射为不同概念消息中间件Name 对应RabbitMQRouting Key路由键KafkaTopicAzure Service BusSubjectNATSSubjectRedis StreamsStream示例[CapSubscribe(order.created)] public void OnOrderCreated(OrderCreatedEvent event) { // 处理订单创建事件 }2.2 Group可选类型string可选Group参数将订阅器放入独立的消费组概念类似 Kafka 的 Consumer Group。若不指定默认使用当前程序集名即DefaultGroupName格式为cap.queue.{程序集名小写}.v1。这一点在源码中有明确体现——ConsumerServiceSelector的SetSubscribeAttribute方法会把组名拼接为{GroupNamePrefix.}{Group ?? DefaultGroupName}.{Version}src/DotNetCore.CAP/Internal/IConsumerServiceSelector.Default.cs。Group 的匹配行为相同 Name、不同 Group所有订阅器都会收到消息相同 Name、相同 Group只有一个订阅器会收到消息竞争消费不同 Name、不同 Group各订阅器拥有独立的执行线程不同 Name、相同 Group共享消费线程。Group 在不同中间件中的映射消息中间件Group 对应RabbitMQQueue队列名KafkaConsumer GroupAzure Service BusSubscription NameNATSQueue GroupRedis StreamsConsumer Group2.3 GroupConcurrent可选类型byte可选GroupConcurrent设置订阅器并发执行的并行度。由于并发执行需要运行在独立线程上如果未指定GroupCAP 会自动用Name值创建一个组。源码印证了这一行为src/DotNetCore.CAP/Internal/IConsumerServiceSelector.Default.csif (attribute.Group null attribute.GroupConcurrent 0) { attribute.Group ${prefix}{attribute.Name}.{_capOptions.Version}; }消费端在创建消费者客户端时会读取该组的并发限制GetGroupConcurrentLimit(groupKey)并以此限定并发消费数量src/DotNetCore.CAP/Internal/IConsumerRegister.Default.cs。注意事项若多个订阅器配置了相同的 Group 且各自设置了GroupConcurrent组内的并行度是所有值的总和该设置只作用于新消息重试消息不受并发限制约束。三、全局配置CapOptions 核心参数详解CapOptions类src/DotNetCore.CAP/CAP.Options.cs集中存放全局配置所有选项在构造函数中都有默认值。下面按功能域逐组讲解。3.1 组名与主题名前缀DefaultGroupName默认值cap.queue.{程序集名}默认消费组名其默认值由构造函数动态生成cap.queue. Assembly.GetEntryAssembly()?.GetName().Name!.ToLower()src/DotNetCore.CAP/CAP.Options.cs。可自定义该值以统一不同传输中的组名便于监控排查。传输DefaultGroupName 映射RabbitMQQueue NameApache KafkaConsumer Group IdAzure Service BusSubscription NameNATSQueue Group NameRedis StreamsConsumer GroupGroupNamePrefix默认值null为所有消费组名添加统一前缀。在订阅器组名生成时前缀会被拼接到最前面src/DotNetCore.CAP/Internal/IConsumerServiceSelector.Default.csvar prefix !string.IsNullOrEmpty(_capOptions.GroupNamePrefix) ? ${_capOptions.GroupNamePrefix}. : string.Empty;TopicNamePrefix默认值null为所有主题/队列名添加统一前缀。发布端在发送前会执行name ${_capOptions.TopicNamePrefix}.{name}src/DotNetCore.CAP/Internal/ICapPublisher.Default.cs订阅端在解析订阅描述符时同样会把前缀拼进主题名src/DotNetCore.CAP/Internal/ConsumerExecutorDescriptor.cs从而保证发布与订阅两侧的名称一致。3.2 版本隔离Version默认值v1Version用于为消息指定版本号实现不同服务实例间消息的版本隔离适用于 A/B 测试或多版本共存场景。组名生成时版本号会拼在最后{Group}.{Version}。典型应用场景业务迭代与向后兼容业务迭代快时消息数据结构可能变化。新系统无碍但已上生产的环境若直接发布不兼容的新结构会引发严重问题按版本隔离后无需清空全部队列和持久化消息再重启应用多服务版本并存服务端需同时为不同客户端版本提供多套接口同一交互的数据结构可能不同不同版本可通过不同路由/组名隔离多实例共享同一存储表/集合多个服务实例共享同一数据库时可通过不同的表名前缀隔离数据表——注意这里是通过设置不同的表名/前缀实现而不是靠Version本身。3.3 重试机制FailedRetryInterval、FailedRetryCount、FallbackWindowLookbackSeconds、FailedThresholdCallback、UseStorageLockFailedRetryInterval默认值60 秒消息发送失败与消费失败时CAP 都会进行重试该参数指定每次重试之间的间隔。⚠️重试节奏说明默认情况下发送或消费失败后重试会在4 分钟即FallbackWindowLookbackSeconds后才开始以避免消息状态延迟带来的潜在问题发送与消费失败前 3 次会立即重试源码中var retryCount Math.Min(_options.Value.FailedRetryCount, 3)见 src/DotNetCore.CAP/Internal/IMessageSender.Default.cs 与 src/DotNetCore.CAP/Internal/ISubscribeExector.Default.cs3 次之后进入轮询式重试此时FailedRetryInterval才真正生效。⚠️多实例并发重试自 7.1.0 版本起CAP 引入基于数据库的分布式锁以减少多实例重试时对数据库的重复拉取必须显式设置UseStorageLock true才能启用。需要明确UseStorageLock协调的是各重试处理器的数据库拉取动作并不提供集群级别的单条毒消息节流。多实例部署时一条持续失败的消息可能在不同轮询周期被不同实例拾取因此观察到的重试节奏会受到副本数量的影响。UseStorageLock默认值false置为 true 后将使用基于数据库的分布式锁处理多实例间重试进程的并发数据拉取同时会在数据库中生成cap.lock表。该锁只保护重试消息的拉取不保证每条失败消息在整个集群范围内以每个FailedRetryInterval至多重试一次。FailedRetryCount默认值50最大重试次数达到该值后停止重试。达到阈值时源码会触发FailedThresholdCallback并记录日志src/DotNetCore.CAP/Internal/IMessageSender.Default.cs订阅端在消息重试超过阈值后同样会记录ConsumerExecutedAfterThreshold日志src/DotNetCore.CAP/Internal/IConsumerRegister.Default.cs。另外若消费抛出的异常是SubscriberNotFoundException订阅器不存在消息会被直接置为失败态而不重试src/DotNetCore.CAP/Internal/ISubscribeExector.Default.cs。FallbackWindowLookbackSeconds默认值240 秒4 分钟配置重试处理器在回看时间窗口内拾取状态为Scheduled或Failed的消息。它保证了即使时钟略有偏差消息仍能被正确处理同时也意味着消费者端积压超过 4 分钟的消息可能被重试线程再次拾取并重复执行详见 3.6 节。FailedThresholdCallback默认值null 类型ActionFailedInfo失败阈值回调。当重试次数达到FailedRetryCount设定值时调用可用于接收通知并人工干预例如发送邮件或告警。FailedInfo位于 src/DotNetCore.CAP/Messages/FailedInfo.cs。services.AddCap(options { options.FailedThresholdCallback failed { // 记录到日志 / 发送邮件 / 推送告警 Console.WriteLine($消息 {failed.Message} 已重试达上限); }; });3.4 清理机制CollectorCleaningInterval、SucceedMessageExpiredAfter、FailedMessageExpiredAfterCollectorCleaningInterval默认值300 秒过期消息收集器的清理间隔即每隔多久删除一次已过期的消息。对应后台处理器CollectorProcessor注册于 src/DotNetCore.CAP/CAP.ServiceCollectionExtensions.cs。SucceedMessageExpiredAfter默认值24×3600 秒1 天成功发送或消费的消息在数据库中的保留时长。消息成功处理后将在该秒数之后从数据库移除。FailedMessageExpiredAfter默认值15×24×3600 秒15 天失败消息在数据库中的保留时长。消息发送或消费失败后将在该秒数之后从数据库移除方便保留足够长的排查窗口。3.5 调度与消费并发SchedulerBatchSize、ConsumerThreadCountSchedulerBatchSize默认值1000每个调度周期内最多拉取的延迟/排队消息数量。调大可提升吞吐但占用更多内存调小则相反。ConsumerThreadCount默认值1消费线程数。当该值大于 1 时无法保证消息执行的顺序。注册逻辑中每个消费组按ConsumerThreadCount启动对应数量的后台任务src/DotNetCore.CAP/Internal/IConsumerRegister.Default.cs。3.6 并行执行与预取EnableSubscriberParallelExecute、SubscriberParallelExecuteThreadCount、SubscriberParallelExecuteBufferFactorEnableSubscriberParallelExecute默认值false 7.0 版本前的同名选项为EnableConsumerPrefetch默认 true现已改名并作废置为 true 时CAP 会从 Broker 预取一批消息缓冲到内存中然后执行订阅方法执行完后再取下一批。底层通过有界 Channel 实现src/DotNetCore.CAP/Processor/IDispatcher.Default.cs缓冲区容量 SubscriberParallelExecuteThreadCount × SubscriberParallelExecuteBufferFactor。⚠️注意事项开启此选项可能带来问题——如果订阅方法执行缓慢且耗时较长重试线程可能拾取尚未执行的消息。重试线程默认拾取 4 分钟FallbackWindowLookbackSeconds前的消息若消费端积压超过 4 分钟这些消息会被再次拾取并重复执行。SubscriberParallelExecuteThreadCount默认值Environment.ProcessorCount逻辑处理器数量启用EnableSubscriberParallelExecute时并行任务执行使用的线程数。分发器会按该值启动对应数量的处理任务src/DotNetCore.CAP/Processor/IDispatcher.Default.cs。SubscriberParallelExecuteBufferFactor默认值1并行订阅执行时用于计算缓冲容量的乘数因子。缓冲容量 SubscriberParallelExecuteThreadCount × SubscriberParallelExecuteBufferFactor。3.7 发布并行发送EnablePublishParallelSend默认值false在 7.2 ≤ 版本 8.1 区间内默认值为 true默认情况下待发送消息被放入单个内存 Channel 中线性处理置为 true 后消息发送任务将由 .NET 线程池并行处理大幅提升发送性能。分发器构造时会读取该开关src/DotNetCore.CAP/Processor/IDispatcher.Default.cs发布通道大小默认设为Environment.ProcessorCount * 500。3.8 已移除与已废弃的选项[Removed] UseDispatchingPerGroup默认值false已在 8.2 版本移除其行为已成为默认行为曾经的语义同一组内有多个消费者时每个消费组将收到的消息推入各自独立的分发管道通道每个通道的线程数设置为ConsumerThreadCount的值。[Obsolete] EnableConsumerPrefetch默认值false7.0 版本之前默认 true该选项已更名为EnableSubscriberParallelExecute请改用新选项。四、综合配置示例与调优建议4.1 一个覆盖核心选项的完整配置services.AddCap(options { // 传输与存储二选一组合即可运行 options.UseRabbitMQ(rabbit { rabbit.HostName localhost; rabbit.Port 5672; rabbit.UserName guest; rabbit.Password guest; }); options.UseSqlServer(Serverlocalhost;DatabaseCapDb;User Idsa;Passwordyour_password;); // 命名与版本 options.DefaultGroupName cap.queue.order-service; options.GroupNamePrefix prod; options.TopicNamePrefix biz; options.Version v2; // 重试 options.FailedRetryCount 30; // 最大重试 30 次 options.FailedRetryInterval 120; // 轮询重试间隔 120 秒 options.FallbackWindowLookbackSeconds 180; options.UseStorageLock true; // 多实例部署时开启数据库分布式锁 options.FailedThresholdCallback failed { /* 告警通知 */ }; // 过期清理 options.SucceedMessageExpiredAfter 24 * 3600; // 成功消息保留 1 天 options.FailedMessageExpiredAfter 15 * 24 * 3600; // 失败消息保留 15 天 options.CollectorCleaningInterval 300; // 并发 options.ConsumerThreadCount 4; options.EnableSubscriberParallelExecute true; options.SubscriberParallelExecuteThreadCount Environment.ProcessorCount; options.SubscriberParallelExecuteBufferFactor 1; options.EnablePublishParallelSend true; });4.2 按场景的调优速查追求发送吞吐将EnablePublishParallelSend置为 true消息发送将在线程池中并行执行消费侧提速适当调大ConsumerThreadCount或开启EnableSubscriberParallelExecute并配合SubscriberParallelExecuteThreadCount注意顺序性会丢失且警惕超过FallbackWindowLookbackSeconds积压导致的重复消费多实例部署开启UseStorageLock减少重试拉取对数据库的竞争同时意识到它不提供单条毒消息的集群级节流控制失败消息堆积根据业务容忍度调整FailedRetryCount与FailedMessageExpiredAfter并用FailedThresholdCallback及时告警多版本并存 / A-B 测试利用Version隔离不同版本的订阅组避免数据结构的兼容性事故。五、进一步阅读传输配置Transports 总览、RabbitMQ、Kafka、Azure Service Bus、NATS、Redis Streams存储配置Storage 总览、SQL Server、PostgreSQL、MySQL、MongoDB消息与事务messaging、transactions、idempotence序列化与过滤器serialization、filter监控仪表盘 dashboard、opentelemetry、kubernetes快速上手quick-start、introduction赞分享后端消息队列微服务消息路由【免费下载链接】CAPDistributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern项目地址https://gitcode.com/gh_mirrors/ca/CAP点击查看免费下载相关推荐Genkit JS 中间件Middleware实战指南用 use 数组为生成、提示词与 Agent 注入重试、回退与工具增强Genkit JS 中间件Middleware实战指南用 use 数组为生成、提示词与 Agent 注入重试、回退与工具增强 中间件Middleware后端消息队列微服务Dendrite配置完全手册从基础设置到高级调优Dendrite配置完全手册从基础设置到高级调优 Dendrite是第二代Matrix家庭服务器采用Go语言编写提供了高性能、可扩展的实时通信解决方案。本后端即时通讯XiaomiGateway3多网关部署指南如何构建稳定的大规模Zigbee网络XiaomiGateway3多网关部署指南如何构建稳定的大规模Zigbee网络 在智能家居生态中Zigbee网络的稳定性和扩展性至关重要。XiaomiGat物联网智能家居上一篇如何快速构建个人离线音频库喜马拉雅下载工具的终极指南下一篇3个核心步骤解锁《Honey Select 2》完整游戏体验HS2-HF_Patch深度指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑