using Dpz.Core.Entity.Base;
using Dpz.Core.Entity.Base.Indexes;

namespace Dpz.Core.Public.Entity;

/// <summary>
/// 消息 Outbox 数据库记录,用于追踪消息的发布和消费状态,实现数据库兜底。
/// </summary>
[BsonIgnoreExtraElements]
public class MessageOutboxRecord : BasicInformationEntity, IIndexedEntity<MessageOutboxRecord>
{
    /// <summary>
    /// 消息唯一标识,来自 <c>MessageBase.MessageId</c>(默认由 <see cref="System.Guid.NewGuid"/> 生成)。
    /// </summary>
    public required string MessageId { get; set; }

    /// <summary>
    /// 消息类型的完全限定名,用于重试时按类型路由。
    /// </summary>
    public required string MessageType { get; set; }

    /// <summary>
    /// 目标 RabbitMQ Exchange 名称,由路由约定解析后存储。
    /// </summary>
    public required string Exchange { get; set; }

    /// <summary>
    /// 路由键,由路由约定解析后存储,重试时直接使用,无需再次推导。
    /// </summary>
    public required string RoutingKey { get; set; }

    /// <summary>
    /// JSON 序列化的消息体,重试时反序列化后重新发布。
    /// </summary>
    public required string Payload { get; set; }

    /// <summary>
    /// 消息来源标识,来自 <c>MessageBase.Source</c>,可为空。
    /// </summary>
    public string? Source { get; set; }

    /// <summary>
    /// 当前 Outbox 状态。
    /// </summary>
    [BsonRepresentation(BsonType.String)]
    public OutboxMessageStatus Status { get; set; }

    /// <summary>
    /// 已尝试发布到 RabbitMQ 的次数,用于计算指数退避重试间隔。
    /// </summary>
    public int PublishAttempts { get; set; }

    /// <summary>
    /// 最后一次尝试发布的时间。
    /// </summary>
    [BsonDateTimeOptions(Kind = DateTimeKind.Local)]
    public DateTime? LastPublishAttemptAt { get; set; }

    /// <summary>
    /// 下一次允许发布重试的时间,由指数退避算法计算(<c>2^PublishAttempts</c> 分钟,上限 60 分钟)。
    /// </summary>
    [BsonDateTimeOptions(Kind = DateTimeKind.Local)]
    public DateTime? NextPublishRetryAt { get; set; }

    /// <summary>
    /// 消息成功投递到 RabbitMQ 的时间。
    /// </summary>
    [BsonDateTimeOptions(Kind = DateTimeKind.Local)]
    public DateTime? SentAt { get; set; }

    /// <summary>
    /// 最后一次发布失败的错误信息。
    /// </summary>
    public string? LastPublishError { get; set; }

    /// <summary>
    /// 已尝试消费的次数,用于计算指数退避重试间隔。
    /// </summary>
    public int ConsumeAttempts { get; set; }

    /// <summary>
    /// 最后一次尝试消费的时间。
    /// </summary>
    [BsonDateTimeOptions(Kind = DateTimeKind.Local)]
    public DateTime? LastConsumeAttemptAt { get; set; }

    /// <summary>
    /// 下一次允许消费重试的时间,由指数退避算法计算(<c>2^ConsumeAttempts</c> 分钟,上限 60 分钟)。
    /// </summary>
    [BsonDateTimeOptions(Kind = DateTimeKind.Local)]
    public DateTime? NextConsumeRetryAt { get; set; }

    /// <summary>
    /// 消息成功消费的时间。
    /// </summary>
    [BsonDateTimeOptions(Kind = DateTimeKind.Local)]
    public DateTime? ConsumedAt { get; set; }

    /// <summary>
    /// 最后一次消费失败的错误信息。
    /// </summary>
    public string? LastConsumeError { get; set; }

    /// <inheritdoc />
    public static IReadOnlyList<EntityIndexDefinition> GetIndexDefinitions() =>
        [
            new()
            {
                Fields = [EntityIndexField.Ascending<MessageOutboxRecord>(x => x.MessageId)],
                Unique = true,
            },
            new()
            {
                Fields =
                [
                    EntityIndexField.Ascending<MessageOutboxRecord>(x => x.Status),
                    EntityIndexField.Ascending<MessageOutboxRecord>(x => x.NextPublishRetryAt),
                ],
            },
            new()
            {
                Fields =
                [
                    EntityIndexField.Ascending<MessageOutboxRecord>(x => x.Status),
                    EntityIndexField.Ascending<MessageOutboxRecord>(x => x.NextConsumeRetryAt),
                ],
            },
            new()
            {
                Fields =
                [
                    EntityIndexField.Ascending<MessageOutboxRecord>(x => x.Status),
                    EntityIndexField.Ascending<MessageOutboxRecord>(x => x.ConsumedAt),
                ],
            },
        ];
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

代码解释

这是一个 消息发件箱(Message Outbox) 的实体类定义,用于实现分布式系统中的消息可靠性模式。

核心目的

实现 Outbox Pattern(发件箱模式),通过数据库记录来追踪消息的发布和消费状态,确保消息的可靠投递和处理,作为消息系统的兜底机制。

主要功能模块

1. 消息标识与路由信息

MessageId      // 全局唯一消息ID
MessageType    // 消息类型全限定名(用于反序列化和路由)
Exchange       // RabbitMQ交换机名称
RoutingKey     // 路由键
Payload        // JSON序列化的消息体
Source         // 消息来源标识

2. 发布状态追踪

Status                  // 当前状态(枚举:待发布/已发布/失败等)
PublishAttempts         // 发布重试次数
LastPublishAttemptAt    // 最后发布尝试时间
NextPublishRetryAt      // 下次重试时间(指数退避算法)
SentAt                  // 成功发送时间
LastPublishError        // 最后的发布错误信息

3. 消费状态追踪

ConsumeAttempts         // 消费重试次数
LastConsumeAttemptAt    // 最后消费尝试时间
NextConsumeRetryAt      // 下次消费重试时间(指数退避)
ConsumedAt              // 成功消费时间
LastConsumeError        // 最后的消费错误信息

4. 数据库索引定义

// 索引1: MessageId唯一索引(保证消息不重复)
// 索引2: Status + NextPublishRetryAt(快速查询待重试发布的消息)
// 索引3: Status + NextConsumeRetryAt(快速查询待重试消费的消息)
// 索引4: Status + ConsumedAt(查询已消费消息)

技术特点

  1. MongoDB 映射:使用 BsonIgnoreExtraElementsBsonRepresentationBsonDateTimeOptions 等特性
  2. 指数退避策略:重试间隔 = 2^重试次数 分钟(上限60分钟)
  3. 双向追踪:同时追踪发布端和消费端的状态
  4. 索引优化:复合索引支持高效的定时扫描重试任务

典型使用场景

  1. 事务提交后将消息写入 Outbox 表
  2. 后台任务扫描待发布消息并投递到 RabbitMQ
  3. 发布失败时记录错误,按指数退避算法重试
  4. 消费端处理后更新消费状态
  5. 支持消息的全生命周期审计和故障排查
评论加载中...