MdpOutboxPushWorker.cs 1.8 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849
  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. try
  25. {
  26. using var scope = _scopeFactory.CreateScope();
  27. var dispatcher = scope.ServiceProvider.GetRequiredService<MdpTargetPushDispatcher>();
  28. var (success, failed, skipped) = await dispatcher.PushPendingAsync(
  29. MdpTargetPushDispatcher.DefaultTake, stoppingToken);
  30. if (success + failed + skipped > 0)
  31. _logger.LogInformation(
  32. "[MdpOutboxPushWorker] success={Success} failed={Failed} retrying={Skipped}",
  33. success, failed, skipped);
  34. }
  35. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  36. {
  37. break;
  38. }
  39. catch (Exception ex)
  40. {
  41. _logger.LogWarning(ex, "[MdpOutboxPushWorker] push failed");
  42. }
  43. }
  44. }
  45. }