Files
Academic-Affairs-System/src/Jiaowu.Api/Infrastructure/BackgroundJobs/RabbitMqBackgroundJobs.cs
T
biss 3ec0f117d3 参照现有的 ExamArrangementJob / SchedulePublishJob 后台任务模式,将发布改为通过 RabbitMQ
队列异步执行,拆分查询消除笛卡尔积。

  修改的文件(共 13 个)

  ┌───────────────────────────────────────────────────────────────┬──────────────────────────────────────────────────┐
  │                             文件                              │                       变更                       │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Domain/Academic/ExamEntities.cs                               │ 新增 ExamPublishJob 实体 + ExamPublishJobStatus  │
  │                                                               │ / ExamPublishJobKind 枚举                        │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Domain/System/BackgroundJobOutboxMessage.cs                   │ BackgroundJobKind 新增 ExamPublish = 6           │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Infrastructure/BackgroundJobs/BackgroundJobOptions.cs         │ 新增 ExamPublishConcurrency 配置项               │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Infrastructure/BackgroundJobs/BackgroundJobRunner.cs          │ RunAsync 和 MarkJobRetryLimitExceeded 添加       │
  │                                                               │ ExamPublish 分支                                 │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Infrastructure/BackgroundJobs/RabbitMqBackgroundJobs.cs       │ JobKinds 数组和 RoutingKey 添加 ExamPublish →    │
  │                                                               │ "exam.publish"                                   │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Infrastructure/BackgroundJobs/BackgroundJobOutboxPublisher.cs │ 启动恢复逻辑添加 ExamPublishJobs                 │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Infrastructure/Persistence/AppDbContext.cs                    │ 新增 ExamPublishJobs DbSet                       │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Controllers/OperationsController.cs                           │ CountFailedJobsAsync / GetFailedJobs             │
  │                                                               │ 添加考试发布失败统计和筛选                       │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Program.cs                                                    │ 校验 ExamPublishConcurrency + 注册               │
  │                                                               │ ExamPublishJobProcessor                          │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Infrastructure/Exams/ExamPublishJobs.cs                       │ 新文件 —                                         │
  │                                                               │ ExamPublishJobProcessor,拆分查询校验后发布      │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Controllers/ExamsController.cs                                │ Publish 改为创建后台任务 + 202 返回;新增 GET    │
  │                                                               │ publish-jobs/{id} / GET plans/{id}/publish-job   │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ Controllers/MakeupExamsController.cs                          │ 同上改造                                         │
  ├───────────────────────────────────────────────────────────────┼──────────────────────────────────────────────────┤
  │ tests/.../TeachingWorkflowRosterTests.cs                      │ 更新测试适配新的异步发布模式                     │
  └───────────────────────────────────────────────────────────────┴──────────────────────────────────────────────────┘

  笛卡尔积消除

  之前:一个 Include 链拉全部 → EF Core 生成 Sessions × Invigilators × RoomLinks × Seats 笛卡尔积

  之后:
  - 场次计数:db.ExamSessions.CountAsync(无 JOIN)
  - 场次摘要:Select new { Id, ClassroomId, InvigilatorCount, RoomLinkCount }(只查所需列)
  - 容量超限:db.ExamRooms.Select(r => new { SeatCount = r.Seats.Count, Capacity })(单表 JOIN)
  - 课程冲突:db.ExamRoomSessions.Where(link => ...CourseId != link.ExamRoom!.CourseId)(独立查询)
  - 每个查询只做自己需要的 JOIN,互不干扰

  测试结果

  213 通过,0 失败,0 跳过

  配置方式

  - BackgroundJobs__Transport=RabbitMq → 走 RabbitMQ 队列 jiaowu.background-jobs.exam.publish
  - BackgroundJobs__Transport=InMemory(默认) → 走内存 Channel
  - BackgroundJobs__ExamPublishConcurrency=1(默认,可调 1-16)
2026-07-28 08:10:22 +08:00

438 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
];
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",
_ => 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;
}
}