Files
Academic-Affairs-System/src/Jiaowu.Api/Infrastructure/BackgroundJobs/RabbitMqBackgroundJobs.cs
T
2026-08-09 10:13:12 +08:00

440 lines
16 KiB
C#

using System.Text.Json;
using Jiaowu.Api.Domain.System;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
namespace Jiaowu.Api.Infrastructure.BackgroundJobs;
public sealed class RabbitMqBackgroundJobTransport(
BackgroundJobOptions jobOptions,
RabbitMqOptions rabbitOptions,
ILogger<RabbitMqBackgroundJobTransport> logger)
: IBackgroundJobTransport, IAsyncDisposable
{
private readonly SemaphoreSlim _gate = new(1, 1);
private IConnection? _connection;
private IChannel? _channel;
public bool IsDurable => true;
public async ValueTask PublishAsync(
BackgroundJobEnvelope message,
CancellationToken cancellationToken)
{
var body = JsonSerializer.SerializeToUtf8Bytes(message);
await _gate.WaitAsync(cancellationToken);
try
{
var channel = await GetChannelAsync(cancellationToken);
var properties = new BasicProperties
{
Persistent = true,
ContentType = "application/json",
MessageId = message.OutboxMessageId.ToString("D"),
Type = message.JobKind.ToString(),
AppId = "jiaowu-api"
};
await channel.BasicPublishAsync(
jobOptions.Exchange,
RabbitMqBackgroundJobTopology.RoutingKey(message.JobKind),
mandatory: true,
properties,
body,
cancellationToken);
}
catch
{
await ResetConnectionAsync();
throw;
}
finally
{
_gate.Release();
}
}
public async Task<bool> CheckHealthAsync(CancellationToken cancellationToken)
{
await _gate.WaitAsync(cancellationToken);
try
{
var channel = await GetChannelAsync(cancellationToken);
return channel.IsOpen;
}
catch (Exception exception)
{
logger.LogWarning(exception, "RabbitMQ health check failed.");
await ResetConnectionAsync();
return false;
}
finally
{
_gate.Release();
}
}
private async Task<IChannel> GetChannelAsync(CancellationToken cancellationToken)
{
if (_channel is { IsOpen: true })
return _channel;
await ResetConnectionAsync();
_connection = await RabbitMqBackgroundJobTopology.CreateConnectionAsync(
rabbitOptions,
"jiaowu-background-job-publisher",
consumerDispatchConcurrency: 1,
cancellationToken);
_channel = await _connection.CreateChannelAsync(
new CreateChannelOptions(
publisherConfirmationsEnabled: true,
publisherConfirmationTrackingEnabled: true),
cancellationToken);
await RabbitMqBackgroundJobTopology.DeclareAsync(
_channel,
jobOptions,
cancellationToken);
return _channel;
}
private async Task ResetConnectionAsync()
{
if (_channel is not null)
{
try
{
await _channel.DisposeAsync();
}
catch
{
// The broker may already have closed the channel.
}
_channel = null;
}
if (_connection is not null)
{
try
{
await _connection.DisposeAsync();
}
catch
{
// The broker may already have closed the connection.
}
_connection = null;
}
}
public async ValueTask DisposeAsync()
{
await _gate.WaitAsync();
try
{
await ResetConnectionAsync();
}
finally
{
_gate.Release();
_gate.Dispose();
}
}
}
public sealed class RabbitMqBackgroundJobWorker(
BackgroundJobOptions jobOptions,
RabbitMqOptions rabbitOptions,
BackgroundJobRunner runner,
ILogger<RabbitMqBackgroundJobWorker> logger) : BackgroundService
{
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
IConnection? connection = null;
var channels = new List<IChannel>();
try
{
var consumerCount = RabbitMqBackgroundJobTopology.JobKinds.Sum(
jobOptions.ConsumerConcurrency);
connection = await RabbitMqBackgroundJobTopology.CreateConnectionAsync(
rabbitOptions,
"jiaowu-background-job-worker",
(ushort)consumerCount,
stoppingToken);
foreach (var kind in RabbitMqBackgroundJobTopology.JobKinds)
{
for (var workerIndex = 0;
workerIndex < jobOptions.ConsumerConcurrency(kind);
workerIndex++)
{
var channel = await connection.CreateChannelAsync(
new CreateChannelOptions(
publisherConfirmationsEnabled: false,
publisherConfirmationTrackingEnabled: false),
stoppingToken);
channels.Add(channel);
await RabbitMqBackgroundJobTopology.DeclareAsync(
channel,
jobOptions,
stoppingToken);
await channel.BasicQosAsync(
0,
jobOptions.PrefetchCount,
global: false,
stoppingToken);
var consumer = new AsyncEventingBasicConsumer(channel);
consumer.ReceivedAsync += async (_, eventArgs) =>
{
using var deliveryCancellation =
CancellationTokenSource.CreateLinkedTokenSource(
eventArgs.CancellationToken,
stoppingToken);
var deliveryToken = deliveryCancellation.Token;
try
{
var message = JsonSerializer
.Deserialize<BackgroundJobEnvelope>(
eventArgs.Body.Span);
if (message is null || message.JobKind != kind)
{
await channel.BasicNackAsync(
eventArgs.DeliveryTag,
multiple: false,
requeue: false,
deliveryToken);
return;
}
var result = await runner.RunAsync(
message,
deliveryToken);
if (result.Outcome ==
BackgroundJobRunOutcome.Completed)
{
await channel.BasicAckAsync(
eventArgs.DeliveryTag,
multiple: false,
deliveryToken);
return;
}
await Task.Delay(
result.RetryAfter,
deliveryToken);
await channel.BasicNackAsync(
eventArgs.DeliveryTag,
multiple: false,
requeue: true,
deliveryToken);
}
catch (OperationCanceledException)
{
// Closing the channel requeues unacknowledged deliveries.
}
catch (Exception exception)
{
logger.LogError(
exception,
"RabbitMQ delivery for {JobKind} failed and will be requeued.",
kind);
if (!channel.IsOpen)
return;
await channel.BasicNackAsync(
eventArgs.DeliveryTag,
multiple: false,
requeue: true,
CancellationToken.None);
}
};
await channel.BasicConsumeAsync(
RabbitMqBackgroundJobTopology.QueueName(
jobOptions,
kind),
autoAck: false,
consumer,
stoppingToken);
}
}
logger.LogInformation(
"{ConsumerCount} RabbitMQ background job consumers are connected to {HostName}:{Port}.",
consumerCount,
rabbitOptions.HostName,
rabbitOptions.Port);
while (connection.IsOpen && !stoppingToken.IsCancellationRequested)
await Task.Delay(TimeSpan.FromSeconds(1), stoppingToken);
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
break;
}
catch (Exception exception)
{
logger.LogError(
exception,
"RabbitMQ background job consumer connection failed; retrying.");
}
finally
{
foreach (var channel in channels)
{
try
{
await channel.DisposeAsync();
}
catch
{
// The connection may already have disposed its channels.
}
}
if (connection is not null)
{
try
{
await connection.DisposeAsync();
}
catch
{
// The broker may already have closed the connection.
}
}
}
if (!stoppingToken.IsCancellationRequested)
await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken);
}
}
}
internal static class RabbitMqBackgroundJobTopology
{
public static readonly BackgroundJobKind[] JobKinds =
[
BackgroundJobKind.AutomaticSchedule,
BackgroundJobKind.SchedulePublish,
BackgroundJobKind.MakeupExamAuto,
BackgroundJobKind.ExamArrangement,
BackgroundJobKind.ExamSignInExport,
BackgroundJobKind.ExamPublish,
BackgroundJobKind.CourseGradeStatisticsRefresh
];
public static async Task<IConnection> CreateConnectionAsync(
RabbitMqOptions options,
string clientName,
ushort consumerDispatchConcurrency,
CancellationToken cancellationToken)
{
var factory = new ConnectionFactory
{
HostName = options.HostName,
Port = options.Port,
UserName = options.UserName,
Password = options.Password,
VirtualHost = options.VirtualHost,
ClientProvidedName = clientName,
AutomaticRecoveryEnabled = false,
TopologyRecoveryEnabled = false,
RequestedHeartbeat = TimeSpan.FromSeconds(30),
ConsumerDispatchConcurrency = consumerDispatchConcurrency
};
if (options.UseTls)
{
factory.Ssl = new SslOption
{
Enabled = true,
ServerName = string.IsNullOrWhiteSpace(options.TlsServerName)
? options.HostName
: options.TlsServerName
};
}
return await factory.CreateConnectionAsync(cancellationToken);
}
public static async Task DeclareAsync(
IChannel channel,
BackgroundJobOptions options,
CancellationToken cancellationToken)
{
var deadExchange = options.Exchange + ".dead";
await channel.ExchangeDeclareAsync(
options.Exchange,
ExchangeType.Direct,
durable: true,
autoDelete: false,
cancellationToken: cancellationToken);
await channel.ExchangeDeclareAsync(
deadExchange,
ExchangeType.Direct,
durable: true,
autoDelete: false,
cancellationToken: cancellationToken);
foreach (var kind in JobKinds)
{
var routingKey = RoutingKey(kind);
var queueName = QueueName(options, kind);
var deadQueueName = queueName + ".dead";
var queueArguments = QueueArguments(options);
queueArguments["x-dead-letter-exchange"] = deadExchange;
queueArguments["x-dead-letter-routing-key"] = routingKey;
await channel.QueueDeclareAsync(
queueName,
durable: true,
exclusive: false,
autoDelete: false,
arguments: queueArguments,
cancellationToken: cancellationToken);
await channel.QueueBindAsync(
queueName,
options.Exchange,
routingKey,
cancellationToken: cancellationToken);
await channel.QueueDeclareAsync(
deadQueueName,
durable: true,
exclusive: false,
autoDelete: false,
arguments: QueueArguments(options),
cancellationToken: cancellationToken);
await channel.QueueBindAsync(
deadQueueName,
deadExchange,
routingKey,
cancellationToken: cancellationToken);
}
}
public static string RoutingKey(BackgroundJobKind kind) => kind switch
{
BackgroundJobKind.AutomaticSchedule => "schedule.automatic",
BackgroundJobKind.SchedulePublish => "schedule.publish",
BackgroundJobKind.MakeupExamAuto => "makeup-exam.automatic",
BackgroundJobKind.ExamArrangement => "exam.arrangement",
BackgroundJobKind.ExamSignInExport => "exam.sign-in-export",
BackgroundJobKind.ExamPublish => "exam.publish",
BackgroundJobKind.CourseGradeStatisticsRefresh => "grade-statistics.refresh",
_ => throw new ArgumentOutOfRangeException(nameof(kind), kind, null)
};
public static string QueueName(
BackgroundJobOptions options,
BackgroundJobKind kind) =>
$"{options.QueuePrefix}.{RoutingKey(kind)}";
private static Dictionary<string, object?> QueueArguments(
BackgroundJobOptions options)
{
var arguments = new Dictionary<string, object?>();
if (options.UseQuorumQueues)
arguments["x-queue-type"] = "quorum";
return arguments;
}
}