| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849 |
- using Microsoft.Extensions.DependencyInjection;
- using Microsoft.Extensions.Hosting;
- using Microsoft.Extensions.Logging;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
- /// <summary>事件驱动出站:收到 <see cref="MdpOutboxWakeSignal"/> 后立即推一批。</summary>
- public sealed class MdpOutboxPushWorker : BackgroundService
- {
- private readonly IServiceScopeFactory _scopeFactory;
- private readonly MdpOutboxWakeSignal _wake;
- private readonly ILogger _logger;
- public MdpOutboxPushWorker(
- IServiceScopeFactory scopeFactory,
- MdpOutboxWakeSignal wake,
- ILoggerFactory loggerFactory)
- {
- _scopeFactory = scopeFactory;
- _wake = wake;
- _logger = loggerFactory.CreateLogger(nameof(MdpOutboxPushWorker));
- }
- protected override async Task ExecuteAsync(CancellationToken stoppingToken)
- {
- await foreach (var _ in _wake.Reader.ReadAllAsync(stoppingToken))
- {
- try
- {
- using var scope = _scopeFactory.CreateScope();
- var dispatcher = scope.ServiceProvider.GetRequiredService<MdpTargetPushDispatcher>();
- var (success, failed, skipped) = await dispatcher.PushPendingAsync(
- MdpTargetPushDispatcher.DefaultTake, stoppingToken);
- if (success + failed + skipped > 0)
- _logger.LogInformation(
- "[MdpOutboxPushWorker] success={Success} failed={Failed} retrying={Skipped}",
- success, failed, skipped);
- }
- catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
- {
- break;
- }
- catch (Exception ex)
- {
- _logger.LogWarning(ex, "[MdpOutboxPushWorker] push failed");
- }
- }
- }
- }
|