MdpOutboxPushWorker.cs 2.1 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253
  1. using Microsoft.Extensions.DependencyInjection;
  2. using Microsoft.Extensions.Hosting;
  3. using Microsoft.Extensions.Logging;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  5. /// <summary>事件驱动出站:收到 <see cref="MdpOutboxWakeSignal"/> 后立即推一批。</summary>
  6. public sealed class MdpOutboxPushWorker : BackgroundService
  7. {
  8. private readonly IServiceScopeFactory _scopeFactory;
  9. private readonly MdpOutboxWakeSignal _wake;
  10. private readonly ILogger _logger;
  11. public MdpOutboxPushWorker(
  12. IServiceScopeFactory scopeFactory,
  13. MdpOutboxWakeSignal wake,
  14. ILoggerFactory loggerFactory)
  15. {
  16. _scopeFactory = scopeFactory;
  17. _wake = wake;
  18. _logger = loggerFactory.CreateLogger(nameof(MdpOutboxPushWorker));
  19. }
  20. protected override async Task ExecuteAsync(CancellationToken stoppingToken)
  21. {
  22. await foreach (var _ in _wake.Reader.ReadAllAsync(stoppingToken))
  23. {
  24. // 逐次求值而非启动时一次:本 Worker 无启动延时,宿主配置未必就绪,
  25. // 而它是事件驱动的,空闲时不产生任何查询,逐次求值的代价可以忽略。
  26. if (!AidopJobGate.ShouldRun(nameof(MdpOutboxPushWorker), _logger)) continue;
  27. try
  28. {
  29. using var scope = _scopeFactory.CreateScope();
  30. var dispatcher = scope.ServiceProvider.GetRequiredService<MdpTargetPushDispatcher>();
  31. var (success, failed, skipped) = await dispatcher.PushPendingAsync(
  32. MdpTargetPushDispatcher.DefaultTake, stoppingToken);
  33. if (success + failed + skipped > 0)
  34. _logger.LogInformation(
  35. "[MdpOutboxPushWorker] success={Success} failed={Failed} retrying={Skipped}",
  36. success, failed, skipped);
  37. }
  38. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  39. {
  40. break;
  41. }
  42. catch (Exception ex)
  43. {
  44. _logger.LogWarning(ex, "[MdpOutboxPushWorker] push failed");
  45. }
  46. }
  47. }
  48. }