分布式事件总线(Aegis.CAP)
当你需要跨服务通信、消息不丢失、或把消息发布和业务操作绑在同一个数据库事务里时,可以使用 Aegis.CAP。它基于 DotNetCore.CAP,提供本地消息表(Outbox Pattern)、消息重试和 FreeSql 事务集成。
Aegis.CAP 和 事件总线(Aegis.EventBus) 不是同一层:前者面向分布式场景,后者面向进程内异步解耦。两者可以同时使用。
组件概览
| 字段 | 说明 |
|---|---|
| 组件名称 | 分布式事件总线 |
| 真实类库 | Aegis.CAP |
| 组件定位 | 基于 DotNetCore.CAP 的分布式事件总线,提供 FreeSql 事务集成 |
| 引入方式 | 安装 NuGet,在 Startup 中手动调用 services.AddCap()。纯代码注册组件,不走 Component.deps.json |
| 组件声明 | 无 |
| 核心能力 | 消息发布/订阅、本地消息表持久化、延迟消息、消费者组、FreeSql 事务集成 |
| 额外依赖 | 使用者需自行引入 Transport 包(如 DotNetCore.CAP.RabbitMQ)和 Storage 包(如 DotNetCore.CAP.PostgreSql) |
| 典型配套 | 事件总线、Aegis.Core.FreeSql |
什么时候要用它
适合场景:
- 微服务之间需要通过事件通知进行通信
- 消息不能丢失,需要持久化到数据库后再投递
- 消息发布需要和业务操作在同一个数据库事务中完成
- 消费失败后需要自动重试,直到成功或达到上限
如果你只是做应用内部的异步解耦(比如主链路完成后触发日志、通知等后置动作),优先用 事件总线。
最小接入方式
第一步:安装 NuGet 包
<PackageReference Include="Aegis.CAP" Version="3.1.0" />
<PackageReference Include="DotNetCore.CAP.RabbitMQ" Version="10.0.1" />
<PackageReference Include="DotNetCore.CAP.PostgreSql" Version="10.0.1" />
Transport 和 Storage 包需要根据实际中间件选择,Aegis.CAP 本身不硬引用它们。
第二步:在 Startup 中注册
CAP 需要绑定数据源和消息中间件,必须在 Startup 中手动注册:
using DotNetCore.CAP;
public void ConfigureServices(IServiceCollection services)
{
var connectionString = ConfigManager.Get("PostgreConnection");
services.AddDbSource<AegisDb>(x =>
{
x.ConnectionString = connectionString;
x.DataType = "PostgreSQL";
});
services.AddDbRepositories<AegisDb>();
services.AddCap(x =>
{
x.UseRabbitMQ(opt =>
{
opt.HostName = "localhost";
opt.Port = 5672;
opt.UserName = "guest";
opt.Password = "guest";
});
x.UsePostgreSql(opt =>
{
opt.ConnectionString = connectionString;
opt.Schema = "cap";
});
x.DefaultGroupName = "aegis.cap";
});
}
Storage 的连接字符串应和 FreeSql 使用同一个数据库,确保本地消息表和业务表在同一个事务中。
第三步:确认生效
启动应用后:
- 数据库中自动创建
capschema 下的消息表(cap.published、cap.received) - 可以注入
ICapPublisher发送消息 - 标记了
[CapSubscribe]的方法能收到消息
命名规范
建议在项目中定义两个常量类,所有 Publish 和 Subscribe 统一引用常量,避免拼写不一致。
消息名称
public static class CapTopics
{
// 订单
public const string OrderCreated = "order.created";
public const string OrderCancelled = "order.cancelled";
public const string OrderTimeoutCheck = "order.timeout.check";
// 库存
public const string StockDeducted = "stock.deducted";
}
命名规则:{聚合根}.{过去式动词},全小写,点号分隔。
消费者组
public static class CapGroups
{
// 按服务名分组(广播消费)
public const string OrderService = "order-service";
public const string NotificationService = "notification-service";
// 同服务内竞争消费时,细化组名
public const string OrderTimeoutProcessor = "order-timeout-processor";
}
命名规则:{服务名}-{用途},kebab-case。同组竞争消费,不同组广播消费。
发布消息
基本发布
注入 ICapPublisher,调用 PublishAsync:
public class OrderController : ControllerBase
{
private readonly ICapPublisher _capPublisher;
public OrderController(ICapPublisher capPublisher)
{
_capPublisher = capPublisher;
}
[HttpPost]
public async Task<IActionResult> CreateOrder(OrderDto order)
{
await _capPublisher.PublishAsync(CapTopics.OrderCreated, new OrderCreatedEvent
{
OrderId = order.Id,
CustomerId = order.CustomerId
});
return Ok();
}
}
延迟消息
不需要消息中间件的延迟队列,CAP 原生支持:
await _capPublisher.PublishDelayAsync(
TimeSpan.FromSeconds(30),
CapTopics.OrderTimeoutCheck,
new { OrderId = order.Id });
订阅消息
在 Controller 中
Controller 中的方法直接用 [CapSubscribe] 标注,不需要实现额外接口:
[CapSubscribe(CapTopics.OrderCreated)]
public void HandleOrderCreated(OrderCreatedEvent evt)
{
// 处理事件
}
在 Service 中
非 Controller 类必须实现 ICapSubscribe:
public class OrderEventSubscriber : ICapSubscribe
{
[CapSubscribe(CapTopics.OrderCreated)]
public async Task HandleAsync(OrderCreatedEvent evt)
{
// 处理事件
}
}
记得注册到 DI 容器:services.AddScoped<OrderEventSubscriber>();
消费者组
广播消费
不同服务各自收到一份消息:
// 订单服务
[CapSubscribe(CapTopics.PaymentSucceeded, Group = CapGroups.OrderService)]
public async Task UpdateOrderStatus(PaymentSucceededEvent evt) { }
// 通知服务
[CapSubscribe(CapTopics.PaymentSucceeded, Group = CapGroups.NotificationService)]
public async Task SendNotification(PaymentSucceededEvent evt) { }
竞争消费
同组只有一个实例收到消息,适合多实例部署:
[CapSubscribe(CapTopics.OrderTimeoutCheck, Group = CapGroups.OrderTimeoutProcessor)]
public async Task CheckTimeout(TimeoutCheckMessage msg) { }
并行消费
通过 GroupConcurrent 设置单个组的消费并行度:
[CapSubscribe(CapTopics.OrderCreated, Group = CapGroups.OrderService, GroupConcurrent = 4)]
public async Task HandleAsync(OrderCreatedEvent evt) { }
FreeSql 事务集成
CAP 的核心价值在于消息发布和业务操作在同一事务中完成。事务必须由 CAP 的 Storage 包主导创建,确保 CAP 能正确写入本地消息表并在提交后推送消息。
标准用法
使用 Storage 包提供的 connection.BeginTransaction(capPublisher, autoCommit: false) 创建事务,通过 WithTransaction() 将事务传递给 FreeSql:
using System.Data;
using Aegis.Core.FreeSql;
public class OrderService
{
private readonly IFreeSql<AegisDb> _fsql;
private readonly ICapPublisher _capPublisher;
public async Task CreateOrderAsync(Order order)
{
using var connection = _fsql.Ado.MasterPool.Get().Value;
// 由 CAP 创建事务
using var transaction = connection.BeginTransaction(_capPublisher, autoCommit: false);
try
{
await _fsql.Insert(order)
.WithTransaction((IDbTransaction)transaction.DbTransaction)
.ExecuteAffRowsAsync();
await _capPublisher.PublishAsync(CapTopics.OrderCreated, new OrderCreatedEvent
{
OrderId = order.Id,
CustomerId = order.CustomerId
});
transaction.Commit();
}
catch
{
transaction.Rollback();
throw;
}
}
}
必须使用
connection.BeginTransaction(capPublisher, autoCommit: false)而非connection.BeginTransaction()。前者由 CAP 的 Storage 包提供(如DotNetCore.CAP.PostgreSql),会创建正确的ICapTransaction实例来管理消息缓冲和推送。
补偿事务
消费者可以返回结果,发布者通过 callbackName 接收回调:
// 发布者:扣减库存,指定回调
await _capPublisher.PublishAsync(
CapTopics.StockDeducted,
contentObj: new StockDeductRequest { OrderId = 1234, ProductId = 23255 },
callbackName: CapTopics.StockDeductFailed);
// 发布者:接收回调
[CapSubscribe(CapTopics.StockDeductFailed)]
public void HandleDeductResult(JsonElement param)
{
var isSuccess = param.GetProperty("IsSuccess").GetBoolean();
// 标记订单状态
}
// 消费者:返回结果
[CapSubscribe(CapTopics.StockDeducted)]
public object DeductStock(StockDeductRequest req)
{
var success = _inventoryService.TryDeduct(req.ProductId, req.Qty);
return new { req.OrderId, IsSuccess = success };
}
消息重试
发送或消费失败时自动重试:先立即重试 3 次,4 分钟后按 FailedRetryInterval 间隔继续重试,达到 FailedRetryCount 后停止。
重试耗尽时可通过 FailedThresholdCallback 收到通知:
services.AddCap(x =>
{
x.FailedThresholdCallback = failed =>
{
// 发送告警
};
});
序列化说明
CAP 内部使用 System.Text.Json,Aegis 的其他组件(如 Aegis.Net.Broker、JsonExtensions)使用 Newtonsoft.Json。这两者互不影响——CAP 的序列化只作用于消息体的传输和持久化,不影响 HTTP API 的输入输出格式。
如果需要自定义消息体的 JSON 格式(如驼峰命名),可在 AddCap 中配置 x.JsonSerializerOptions。
配置项说明
CAP 的完整配置选项请参考 CAP 官方文档。
常用选项:
| 选项 | 默认值 | 说明 |
|---|---|---|
DefaultGroupName | 程序集名 | 默认消费者组名 |
TopicNamePrefix | null | Topic 名称前缀 |
GroupNamePrefix | null | Group 名称前缀 |
Version | "v1" | 消息版本号,用于隔离 |
FailedRetryCount | 50 | 最大重试次数 |
FailedRetryInterval | 60 | 重试间隔(秒) |
UseStorageLock | false | 启用分布式锁防并发重试 |
关键边界
- CAP 需要在 Startup 中手动注册,不走
Component.deps.json,和Aegis.RocketMQ.Client一样是纯代码注册组件 - Transport 和 Storage 包由使用者按需引入,
Aegis.CAP不硬引用 - CAP 内部使用
System.Text.Json,与 Aegis 的Newtonsoft.Json互不影响
常见问题
为什么要在 Startup 手动注册
CAP 的 Storage 需要绑定具体的数据库连接字符串,这个连接字符串应和 FreeSql 使用同一个数据库。框架无法自动推断使用者的数据源配置。
已经注册了 CAP,但启动报错找不到程序集
Aegis.CAP 不硬引用 Transport 和 Storage 包。使用者需要在项目中安装对应的 NuGet 包,例如 DotNetCore.CAP.RabbitMQ 和 DotNetCore.CAP.PostgreSql。
和 Aegis.EventBus 能一起用吗
可以。两者完全独立:EventBus 处理进程内事件(Channel 内存队列),CAP 处理分布式事件(RabbitMQ/Kafka + 数据库持久化)。