MoneyTree.Outbox 1.0.4
MoneyTree.Outbox 领域事件 Outbox 模块
📋 概述
MoneyTree.Outbox 提供领域事件的后台派发服务,将事务内持久化的领域事件异步派发到 MediatR 或自定义处理器,保障最终一致性。
与 MoneyTree.EFCore 协同工作,实现完整的 Transactional Outbox 模式。
| 属性 | 说明 |
|---|---|
| NuGet 包 | MoneyTree.Outbox |
| 外部依赖 | MoneyTree.EFCore、MediatR、Microsoft.Extensions.Hosting.Abstractions |
| 定位 | 领域事件后台派发调度 |
🏗️ 职责划分
Outbox 模式由两个模块协同实现,职责清晰分离:
MoneyTree.EFCore(ORM 层,事务内持久化)
├── OutboxMessage.cs # Outbox 消息实体模型
├── OutboxMessageStatus.cs # 状态枚举
├── OutboxMessageConfiguration.cs # EF Core 模型配置(表名、索引)
├── IEventOutbox.cs # 读写接口契约
└── OutboxInterceptor.cs # 拦截器(事务提交前持久化领域事件)
MoneyTree.Outbox(运行时层,后台派发调度)
├── OutboxOptions.cs # 调度配置选项
├── DbEventOutbox.cs # IEventOutbox 实现
├── IEventDispatcher.cs # 派发接口
├── DispatchResult.cs # 派发结果(成功/失败+错误详情)
├── MediatrEventDispatcher.cs # MediatR 派发实现
├── IOutboxTenantContextRestorer.cs # 租户上下文恢复接口(多租户可选)
├── OutboxDispatcherService.cs # 后台派发服务
└── OutboxServiceExtensions.cs # DI 注册扩展
为什么拆分?
- EFCore 是类库,不应包含
BackgroundService等运行时行为 - Outbox 是运行时调度模块,需要宿主上下文
- 职责分离后,用户可仅使用持久化层(自行实现派发),或同时使用两者
🎯 核心能力
| 能力 | 实现 |
|---|---|
| 事务内持久化 | OutboxInterceptor(EFCore 层) |
| 后台轮询派发 | OutboxDispatcherService<TDbContext> |
| 原子抢占 + 多副本并发安全 | ExecuteUpdateAsync + LockToken |
| 超时回收 | ProcessingTimeout 回收卡死的 Processing 消息 |
| 指数退避重试 | NextRetryAt + RetryBaseDelaySeconds,避免快速消耗重试次数 |
| 派发器可替换 | IEventDispatcher 接口 |
| 错误详情传递 | DispatchResult 携带实际错误信息写入 LastError |
| 多租户上下文恢复 | IOutboxTenantContextRestorer(可选注入) |
| 死信重置 | IEventOutbox.ResetFailedAsync 人工干预后重新派发 |
| 过期消息清理 | ProcessedRetentionHours 定期清理 Processed 消息 |
| EF Core 自动建表索引 | OutboxMessageConfiguration(Migrations/EnsureCreated 友好) |
🚀 快速开始
两步启用
// Program.cs
using MoneyTree.EFCore.Extensions;
using MoneyTree.EFCore.PostgreSql.Extensions;
using MoneyTree.Outbox;
var builder = WebApplication.CreateBuilder(args);
// 步骤 1:EFCore 层 - 启用 Outbox 拦截器(事务内持久化领域事件)
builder.Services.AddMoneyTreeEFCore<AppDbContext>(db =>
{
db.UsePostgreSql<AppDbContext>(
builder.Configuration.GetConnectionString("Default")!);
db.UseOutbox(); // 启用持久化拦截器
});
// 步骤 2:Outbox 层 - 启用后台派发服务(异步派发到 MediatR)
builder.Services.AddMoneyTreeOutboxDispatcher<AppDbContext>(opt =>
{
opt.PollIntervalSeconds = 10; // 轮询间隔
opt.BatchSize = 50; // 单次批量
opt.MaxRetryCount = 5; // 最大重试次数
opt.RetryBaseDelaySeconds = 30; // 重试基础延迟(指数退避)
opt.ProcessedRetentionHours = 168; // Processed 消息保留 7 天
});
var app = builder.Build();
app.Run();
仅使用持久化层(自定义派发)
如果需要对接 MassTransit/Kafka 等消息中间件,可仅启用持久化层:
// 仅启用 Outbox 拦截器,不注册后台派发服务
builder.Services.AddMoneyTreeEFCore<AppDbContext>(db =>
{
db.UsePostgreSql<AppDbContext>(connectionString);
db.UseOutbox();
});
// 自定义派发:实现 IEventOutbox 读取 outbox_messages 表,通过 MassTransit 发布
多租户上下文恢复(可选)
后台 Dispatcher 不在 HTTP 请求上下文中,多租户场景下需实现 IOutboxTenantContextRestorer
在派发前根据 OutboxMessage.TenantId 恢复租户上下文:
using MoneyTree.MultiTenancy.Abstractions;
using MoneyTree.MultiTenancy.Core;
using MoneyTree.Outbox;
/// <summary>
/// 多租户上下文恢复器。
/// 在 Outbox 派发前根据消息中的 TenantId 恢复租户上下文。
/// </summary>
public class OutboxTenantContextRestorer : IOutboxTenantContextRestorer
{
private readonly ITenantContextSetter _setter;
private readonly TenantContextAccessor _accessor;
public OutboxTenantContextRestorer(
ITenantContextSetter setter,
TenantContextAccessor accessor)
{
_setter = setter;
_accessor = accessor;
}
public Task<IDisposable> RestoreAsync(long? tenantId, CancellationToken ct = default)
{
if (tenantId.HasValue)
{
// ITenantContextSetter 继承 ITenantContext,设置后可直接赋值给 accessor
_setter.SetTenant(tenantId.Value, "Outbox");
_accessor.TenantContext = _setter;
}
else
{
// 无租户消息:清除可能残留的上一条消息上下文,防止跨消息泄漏
_setter.Clear();
_accessor.TenantContext = null;
}
// 返回清理 scope,派发结束后清除 AsyncLocal,防止下一条消息继承当前租户
return Task.FromResult<IDisposable>(new TenantContextScope(_accessor));
}
/// <summary>派发完成后清除 AsyncLocal 租户上下文。</summary>
private sealed class TenantContextScope : IDisposable
{
private readonly TenantContextAccessor _accessor;
public TenantContextScope(TenantContextAccessor accessor) => _accessor = accessor;
public void Dispose() => _accessor.TenantContext = null;
}
}
// 注册(Scoped)
builder.Services.AddScoped<IOutboxTenantContextRestorer, OutboxTenantContextRestorer>();
⚠️ 关键:
RestoreAsync返回的IDisposable必须在Dispose时清除TenantContextAccessor.TenantContext(AsyncLocal)。OutboxDispatcherService在同一 Scope 内批量派发多条消息,如果不清除,上一条消息的租户上下文会泄漏到下一条消息(尤其是 TenantId 为 null 的消息),导致跨租户数据访问。 同时,tenantId为 null 时也必须清除上下文。
未注册
IOutboxTenantContextRestorer时,OutboxDispatcherService 跳过租户恢复,保持向后兼容。
📖 工作流程
应用层调用 SaveChangesAsync
│
▼
EF Core 隐式事务开始
│
├── 1. OutboxInterceptor.SavingChangesAsync
│ └── 将实体领域事件序列化写入 outbox_messages 表
│ └── 清除实体领域事件(防止 DomainEventInterceptor 重复派发)
│
├── 2. EF Core 提交事务(业务数据 + outbox 消息原子提交)
│
▼
OutboxDispatcherService 后台轮询(独立线程)
│
├── 3. GetPendingAsync(batchSize)
│ ├── 3a. 回收超时 Processing 消息(LockedAt < now - ProcessingTimeout)
│ ├── 3b. 原子抢占 Pending 消息(过滤 NextRetryAt 延迟重试)
│ └── 3c. 查询当前 Dispatcher 抢占的消息(LockToken 识别)
│
├── 4. 可选恢复租户上下文(IOutboxTenantContextRestorer)
│
├── 5. IEventDispatcher.DispatchAsync(message)
│ └── 返回 DispatchResult(含错误详情)
│ └── 反序列化领域事件 → IMediator.Publish
│
├── 6a. 成功 → MarkAsProcessedAsync(标记 Processed)
│
└── 6b. 失败 → MarkAsFailedAsync
└── 重试次数+1,指数退避设置 NextRetryAt
└── 超过 MaxRetryCount → 标记 Failed(人工介入)
│
▼
定期清理(每 6 轮)
└── PurgeProcessedAsync 删除超过保留期的 Processed 消息
🔧 配置选项
| 属性 | 类型 | 默认值 | 说明 |
|---|---|---|---|
PollIntervalSeconds |
int |
10 | Dispatcher 轮询间隔(秒),最小 1 |
BatchSize |
int |
50 | 单次轮询获取的最大消息数量,最小 1,最大 500 |
MaxRetryCount |
int |
5 | 最大重试次数,超过后标记为 Failed,最小 0 |
ProcessingTimeoutSeconds |
int |
300 | Processing 状态超时时间(秒),最小 30。超时后回退为 Pending |
RetryBaseDelaySeconds |
int |
30 | 重试基础延迟(秒),最小 1。指数退避:第 N 次 = Base × 2^(N-1) |
ProcessedRetentionHours |
int |
168 | Processed 消息保留时长(小时),0 表示永不清理。默认 7 天 |
指数退避示例
以 RetryBaseDelaySeconds = 30 为例:
| 重试次数 | 延迟 |
|---|---|
| 第 1 次 | 30 秒 |
| 第 2 次 | 60 秒 |
| 第 3 次 | 120 秒 |
| 第 4 次 | 240 秒 |
| 第 5 次 | 480 秒 |
🔄 自定义派发器
实现 IEventDispatcher 接口替换默认的 MediatR 派发:
public class MassTransitEventDispatcher : IEventDispatcher
{
private readonly IPublishEndpoint _publishEndpoint;
public MassTransitEventDispatcher(IPublishEndpoint publishEndpoint)
{
_publishEndpoint = publishEndpoint;
}
public async Task<DispatchResult> DispatchAsync(OutboxMessage message, CancellationToken ct = default)
{
try
{
var eventType = Type.GetType(message.EventType);
var domainEvent = JsonSerializer.Deserialize(message.EventPayload, eventType!);
await _publishEndpoint.Publish(domainEvent!, ct);
return DispatchResult.Ok();
}
catch (Exception ex)
{
return DispatchResult.Failure(ex.Message);
}
}
}
// 注册时替换默认实现
builder.Services.AddMoneyTreeOutboxDispatcher<AppDbContext>();
builder.Services.AddScoped<IEventDispatcher, MassTransitEventDispatcher>();
🛡️ 死信消息管理
超过 MaxRetryCount 的消息标记为 Failed,可通过 IEventOutbox.ResetFailedAsync 重置:
// 注入 IEventOutbox,重置 Failed 消息为 Pending
await outbox.ResetFailedAsync(messageId);
// 批量重置所有 Failed 消息(需自行查询后逐条重置)
📊 数据库索引
OutboxMessageConfiguration 自动配置以下索引(EF Core Migrations / EnsureCreated 友好):
| 索引名 | 列 | 用途 |
|---|---|---|
idx_status_occurred |
(Status, OccurredOn) | 拾取 Pending 消息排序 |
idx_status_locked_at |
(Status, LockedAt) | 超时回收查询 |
idx_status_next_retry |
(Status, NextRetryAt) | 延迟重试过滤 |
idx_tenant_status |
(TenantId, Status) | 多租户维度查询 |
🛡️ 设计原则
- SRP:EFCore 负责持久化,Outbox 负责调度,职责分离
- OCP:
IEventDispatcher可替换为任意派发实现 - DIP:
OutboxDispatcherService依赖IEventOutbox/IEventDispatcher抽象 - ISP:
IEventOutbox(读写)与IEventDispatcher(派发)接口分离 - 可选依赖:
IOutboxTenantContextRestorer未注册时跳过,向后兼容
No packages depend on MoneyTree.Outbox.
.NET 10.0
- MoneyTree.EFCore (>= 1.0.4)
- MediatR (>= 14.2.0)
- Microsoft.EntityFrameworkCore (>= 10.0.9)
- Microsoft.EntityFrameworkCore.Relational (>= 10.0.9)
- Microsoft.Extensions.Hosting.Abstractions (>= 10.0.9)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.9)
- Microsoft.Extensions.Options (>= 10.0.9)