107 lines
4.2 KiB
C#
107 lines
4.2 KiB
C#
|
|
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;
|
|||
|
|
|
|||
|
|
/// <summary>
|
|||
|
|
/// 消息Outbox派发服务。
|
|||
|
|
/// </summary>
|
|||
|
|
public class MessageOutboxDispatchService(
|
|||
|
|
IServiceScopeFactory scopeFactory,
|
|||
|
|
IRabbitMQService rabbitMQService,
|
|||
|
|
IConfiguration configuration,
|
|||
|
|
ILogger<MessageOutboxDispatchService> logger) : BackgroundService
|
|||
|
|
{
|
|||
|
|
private const int BatchSize = 50;
|
|||
|
|
private readonly RabbitMQRetrySettings retrySettings = configuration.GetSection("RabbitMQRetrySettings").Get<RabbitMQRetrySettings>() ?? new RabbitMQRetrySettings();
|
|||
|
|
|
|||
|
|
/// <summary>
|
|||
|
|
/// 执行Outbox派发循环。
|
|||
|
|
/// </summary>
|
|||
|
|
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<ISqlSugarClient>();
|
|||
|
|
var now = DateTime.Now;
|
|||
|
|
var messages = await db.Queryable<MessageOutbox>()
|
|||
|
|
.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<MessageOutbox>()
|
|||
|
|
.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<MessageOutbox>()
|
|||
|
|
.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);
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
}
|