MoneyTree.Outbox 1.0.4

MoneyTree.Outbox 领域事件 Outbox 模块

📋 概述

MoneyTree.Outbox 提供领域事件的后台派发服务,将事务内持久化的领域事件异步派发到 MediatR 或自定义处理器,保障最终一致性。

MoneyTree.EFCore 协同工作,实现完整的 Transactional Outbox 模式。

属性 说明
NuGet 包 MoneyTree.Outbox
外部依赖 MoneyTree.EFCoreMediatRMicrosoft.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 负责调度,职责分离
  • OCPIEventDispatcher 可替换为任意派发实现
  • DIPOutboxDispatcherService 依赖 IEventOutbox/IEventDispatcher 抽象
  • ISPIEventOutbox(读写)与 IEventDispatcher(派发)接口分离
  • 可选依赖IOutboxTenantContextRestorer 未注册时跳过,向后兼容

No packages depend on MoneyTree.Outbox.

Version Downloads Last updated
1.0.5 1 7/21/2026
1.0.4 1 7/12/2026