EtlInstanceStore.cs 2.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475
  1. using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
  2. namespace Admin.NET.Plugin.AiDOP.Infrastructure;
  3. /// <summary>
  4. /// ETL 执行机实例注册表的写入侧。
  5. ///
  6. /// <para><b>「最多一台执行机」这条不变式由本类保证</b>,不要在控制器、前端或数据库约束里另做一份:
  7. /// MySQL 没有部分索引,无法用唯一索引表达「最多一行 <c>is_runner=1</c>」;
  8. /// 而前端保证等于没有保证(两个超管同时点两台就会破)。
  9. /// 故 <see cref="AssignRunnerAsync"/> 在**同一事务内**先清后置。</para>
  10. /// </summary>
  11. public interface IEtlInstanceStore
  12. {
  13. /// <summary>把指定实例指派为唯一执行机。同一事务内先把其余行置 0。</summary>
  14. /// <returns>false 表示该 instanceId 不存在(进程可能已退出并被清理)。</returns>
  15. Task<bool> AssignRunnerAsync(string instanceId, CancellationToken ct = default);
  16. /// <summary>取消所有指派。此后无人跑定时 ETL——这是允许的状态。</summary>
  17. Task ClearRunnerAsync(CancellationToken ct = default);
  18. Task<List<AdoEtlInstance>> ListAsync(CancellationToken ct = default);
  19. }
  20. public sealed class EtlInstanceStore : IEtlInstanceStore, ITransient
  21. {
  22. private readonly ISqlSugarClient _db;
  23. public EtlInstanceStore(ISqlSugarClient db) => _db = db;
  24. public async Task<bool> AssignRunnerAsync(string instanceId, CancellationToken ct = default)
  25. {
  26. if (string.IsNullOrWhiteSpace(instanceId))
  27. return false;
  28. var id = instanceId.Trim();
  29. var assigned = false;
  30. // 事务边界不可省:先清后置若被拆成两个自动提交语句,
  31. // 中间那一瞬会出现「零台执行机」或(并发指派时)「两台执行机」。
  32. await _db.Ado.UseTranAsync(async () =>
  33. {
  34. var now = DateTime.Now;
  35. await _db.Updateable<AdoEtlInstance>()
  36. .SetColumns(x => new AdoEtlInstance { IsRunner = false, UpdateTime = now })
  37. .Where(x => x.IsRunner && x.InstanceId != id)
  38. .ExecuteCommandAsync(ct);
  39. var n = await _db.Updateable<AdoEtlInstance>()
  40. .SetColumns(x => new AdoEtlInstance { IsRunner = true, UpdateTime = now })
  41. .Where(x => x.InstanceId == id)
  42. .ExecuteCommandAsync(ct);
  43. assigned = n > 0;
  44. });
  45. return assigned;
  46. }
  47. public async Task ClearRunnerAsync(CancellationToken ct = default)
  48. {
  49. var now = DateTime.Now;
  50. await _db.Updateable<AdoEtlInstance>()
  51. .SetColumns(x => new AdoEtlInstance { IsRunner = false, UpdateTime = now })
  52. .Where(x => x.IsRunner)
  53. .ExecuteCommandAsync(ct);
  54. }
  55. public Task<List<AdoEtlInstance>> ListAsync(CancellationToken ct = default) =>
  56. _db.Queryable<AdoEtlInstance>()
  57. .OrderBy(x => x.MachineName)
  58. .OrderBy(x => x.ProcessId)
  59. .ToListAsync(ct);
  60. }