using QYZH.InteractiveMagazine.Infrastructure.RabbitMQ;
using QYZH.InteractiveMagazine.Models.Entity;
using QYZH.InteractiveMagazine.Models.Enum;
using QYZH.InteractiveMagazine.Models.Settings;
using SqlSugar;
using System.Text.Json;
namespace QYZH.InteractiveMagazine.WorkService.Jobs;
///
/// 消息Outbox派发服务。
///
public class MessageOutboxDispatchService(
IServiceScopeFactory scopeFactory,
IRabbitMQService rabbitMQService,
IConfiguration configuration,
ILogger logger) : BackgroundService
{
private const int BatchSize = 50;
private readonly RabbitMQRetrySettings retrySettings = configuration.GetSection("RabbitMQRetrySettings").Get() ?? new RabbitMQRetrySettings();
///
/// 执行Outbox派发循环。
///
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
try
{
await DispatchPendingMessagesAsync(stoppingToken);
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
break;
}
catch (Exception ex)
{
logger.LogError(ex, "Outbox消息派发循环异常");
}
await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken);
}
}
private async Task DispatchPendingMessagesAsync(CancellationToken cancellationToken)
{
using var scope = scopeFactory.CreateScope();
var db = scope.ServiceProvider.GetRequiredService();
var now = DateTime.Now;
var messages = await db.Queryable()
.Where(x => !x.IsDeleted
&& (x.Status == (int)MessageOutboxStatusEnum.Pending || x.Status == (int)MessageOutboxStatusEnum.Failed)
&& (x.NextRetryAt == null || x.NextRetryAt <= now))
.OrderBy(x => x.CreatedAt)
.Take(BatchSize)
.ToListAsync(cancellationToken);
foreach (var message in messages)
{
await DispatchMessageAsync(db, message, cancellationToken);
}
}
private async Task DispatchMessageAsync(ISqlSugarClient db, MessageOutbox message, CancellationToken cancellationToken)
{
try
{
using var payload = JsonDocument.Parse(message.Payload);
var sent = await rabbitMQService.SendAsync(new RabbitMQSendParam
{
Exchange = message.Exchange,
Queue = message.Queue,
RoutingKey = message.RoutingKey,
Data = payload.RootElement.Clone()
}, cancellationToken);
if (!sent)
{
throw new InvalidOperationException("RabbitMQ SendAsync returned false");
}
await db.Updateable()
.SetColumns(x => x.Status == (int)MessageOutboxStatusEnum.Sent)
.SetColumns(x => x.SentAt == DateTime.Now)
.SetColumns(x => x.UpdatedAt == DateTime.Now)
.Where(x => x.Id == message.Id && x.Status != (int)MessageOutboxStatusEnum.Sent)
.ExecuteCommandAsync(cancellationToken);
}
catch (Exception ex)
{
var retryCount = message.RetryCount + 1;
var abandoned = retryCount >= retrySettings.MaxRetryCount;
await db.Updateable()
.SetColumns(x => x.Status == (int)(abandoned ? MessageOutboxStatusEnum.Abandoned : MessageOutboxStatusEnum.Failed))
.SetColumns(x => x.RetryCount == retryCount)
.SetColumns(x => x.NextRetryAt == (abandoned ? null : DateTime.Now.AddMilliseconds(retrySettings.RetryDelayMilliseconds)))
.SetColumns(x => x.LastError == ex.Message)
.SetColumns(x => x.UpdatedAt == DateTime.Now)
.Where(x => x.Id == message.Id)
.ExecuteCommandAsync(cancellationToken);
logger.LogError(ex, "Outbox消息派发失败,MessageId: {MessageId}, RetryCount: {RetryCount}", message.Id, retryCount);
}
}
}