跳到主要内容
版本:3.0.0

分布式事件总线(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 使用同一个数据库,确保本地消息表和业务表在同一个事务中。

第三步:确认生效

启动应用后:

  • 数据库中自动创建 cap schema 下的消息表(cap.publishedcap.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.BrokerJsonExtensions)使用 Newtonsoft.Json。这两者互不影响——CAP 的序列化只作用于消息体的传输和持久化,不影响 HTTP API 的输入输出格式。

如果需要自定义消息体的 JSON 格式(如驼峰命名),可在 AddCap 中配置 x.JsonSerializerOptions

配置项说明

CAP 的完整配置选项请参考 CAP 官方文档

常用选项:

选项默认值说明
DefaultGroupName程序集名默认消费者组名
TopicNamePrefixnullTopic 名称前缀
GroupNamePrefixnullGroup 名称前缀
Version"v1"消息版本号,用于隔离
FailedRetryCount50最大重试次数
FailedRetryInterval60重试间隔(秒)
UseStorageLockfalse启用分布式锁防并发重试

关键边界

  • 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.RabbitMQDotNetCore.CAP.PostgreSql

和 Aegis.EventBus 能一起用吗

可以。两者完全独立:EventBus 处理进程内事件(Channel 内存队列),CAP 处理分布式事件(RabbitMQ/Kafka + 数据库持久化)。

配套阅读