using Dpz.Core.MessageQueue.Models;
namespace Dpz.Core.MessageQueue.Abstractions;
/// <summary>
/// 消息 Outbox 存储接口,用于持久化消息投递状态,实现发布和消费的数据库兜底。
/// </summary>
public interface IMessageOutboxStore
{
/// <summary>
/// 创建一条 Outbox 记录,在消息发布前调用,初始状态为"待发布"。
/// </summary>
/// <param name="messageId">消息唯一标识</param>
/// <param name="messageType">消息类型的完全限定名</param>
/// <param name="exchange">目标 Exchange 名称</param>
/// <param name="routingKey">路由键</param>
/// <param name="payload">JSON 序列化的消息体</param>
/// <param name="source">消息来源标识,可为空</param>
/// <param name="cancellationToken">取消令牌</param>
Task CreateAsync(
string messageId,
string messageType,
string exchange,
string routingKey,
string payload,
string? source,
CancellationToken cancellationToken = default
);
/// <summary>
/// 将消息标记为已成功发布到 RabbitMQ,累加发布尝试次数,状态变为"已发布",
/// 并记录本次尝试时间和发布时间。
/// </summary>
/// <param name="messageId">消息唯一标识</param>
/// <param name="cancellationToken">取消令牌</param>
Task MarkSentAsync(string messageId, CancellationToken cancellationToken = default);
/// <summary>
/// 将消息标记为发布失败,累加发布尝试次数并按指数退避计算下次重试时间,
/// 状态变为"发布失败"。
/// </summary>
/// <param name="messageId">消息唯一标识</param>
/// <param name="error">失败原因描述</param>
/// <param name="cancellationToken">取消令牌</param>
Task MarkPublishFailedAsync(
string messageId,
string error,
CancellationToken cancellationToken = default
);
/// <summary>
/// 将消息标记为已成功消费,累加消费尝试次数,状态变为"已消费",
/// 并记录本次尝试时间和消费时间。
/// </summary>
/// <param name="messageId">消息唯一标识</param>
/// <param name="cancellationToken">取消令牌</param>
Task MarkConsumedAsync(string messageId, CancellationToken cancellationToken = default);
/// <summary>
/// 记录一次消费失败尝试,累加消费尝试次数并记录错误,但不改变消息当前状态。
/// 用于消费者仍会立即重发消息体的短重试路径。
/// </summary>
/// <param name="messageId">消息唯一标识</param>
/// <param name="error">失败原因描述</param>
/// <param name="cancellationToken">取消令牌</param>
Task RecordConsumeAttemptFailedAsync(
string messageId,
string error,
CancellationToken cancellationToken = default
);
/// <summary>
/// 将消息标记为消费失败,累加消费尝试次数并按指数退避计算下次重试时间,
/// 状态变为"消费失败"。
/// </summary>
/// <param name="messageId">消息唯一标识</param>
/// <param name="error">失败原因描述</param>
/// <param name="cancellationToken">取消令牌</param>
Task MarkConsumeFailedAsync(
string messageId,
string error,
CancellationToken cancellationToken = default
);
/// <summary>
/// 获取当前可进行发布重试的 Outbox 条目(状态为"发布失败"且退避时间已到)。
/// </summary>
/// <param name="batchSize">每批最大条数</param>
/// <param name="cancellationToken">取消令牌</param>
/// <returns>待重试的 Outbox 条目列表</returns>
Task<IReadOnlyList<MessageOutboxEntry>> GetPendingPublishRetryAsync(
int batchSize,
CancellationToken cancellationToken = default
);
/// <summary>
/// 获取当前可进行消费重试的 Outbox 条目(状态为"消费失败"且退避时间已到)。
/// </summary>
/// <param name="batchSize">每批最大条数</param>
/// <param name="cancellationToken">取消令牌</param>
/// <returns>待重试的 Outbox 条目列表</returns>
Task<IReadOnlyList<MessageOutboxEntry>> GetPendingConsumeRetryAsync(
int batchSize,
CancellationToken cancellationToken = default
);
}
⚠⚠ 以下内容为AI分析的结果,请根据实际情况进行判断。
代码解释:消息 Outbox 存储接口
这是一个用于实现 Outbox 模式 的接口定义,Outbox 模式是分布式系统中保证消息可靠投递的常用模式。
核心目的
该接口定义了消息持久化存储的标准操作,用于:
- 持久化消息状态:在数据库中记录消息的生命周期
- 实现消息可靠投递:防止消息丢失
- 支持重试机制:通过数据库兜底实现发布和消费的重试
主要方法说明
1. CreateAsync - 创建消息记录
Task CreateAsync(string messageId, string messageType, ...)
- 在消息发布前创建一条记录
- 初始状态为"待发布"
- 记录消息的元数据(ID、类型、路由信息、消息体等)
2. MarkSentAsync - 标记发布成功
- 消息成功发送到 RabbitMQ 后调用
- 状态变更为"已发布"
- 累加发布尝试次数
3. MarkPublishFailedAsync - 标记发布失败
- 发布失败时调用
- 记录失败原因
- 使用指数退避算法计算下次重试时间
4. MarkConsumedAsync - 标记消费成功
- 消费者成功处理消息后调用
- 状态变更为"已消费"
- 累加消费尝试次数
5. RecordConsumeAttemptFailedAsync - 记录消费尝试失败
- 记录一次失败尝试,但不改变状态
- 用于短重试路径(立即重试的场景)
6. MarkConsumeFailedAsync - 标记消费失败
- 消费失败且需要延迟重试时调用
- 状态变更为"消费失败"
- 使用指数退避算法计算下次重试时间
7. GetPendingPublishRetryAsync - 获取待重试的发布消息
- 查询状态为"发布失败"且退避时间已到的消息
- 支持批量获取(batchSize)
8. GetPendingConsumeRetryAsync - 获取待重试的消费消息
- 查询状态为"消费失败"且退避时间已到的消息
- 支持批量获取
设计亮点
- 完整的生命周期管理:覆盖从创建到发布、消费的完整流程
- 区分短重试和长重试:
RecordConsumeAttemptFailedAsyncvsMarkConsumeFailedAsync - 指数退避:失败重试采用指数退避算法,避免系统压力
- 异步设计:所有方法都是异步的,支持高并发场景
- 取消令牌支持:所有方法都支持
CancellationToken,便于优雅关闭
典型使用场景
这个接口通常配合数据库事务使用,确保业务操作和消息发布的原子性,是实现最终一致性的关键组件。
AI 正在分析代码…
评论加载中...