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); } } }