S1MdpFullRunLock.cs 5.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154
  1. namespace Admin.NET.Plugin.AiDOP.DataPlatform.S1Refresh;
  2. public static class S1DashboardRebuildStatus
  3. {
  4. public const string Queued = "QUEUED";
  5. public const string Running = "RUNNING";
  6. public const string Success = "SUCCESS";
  7. public const string Failed = "FAILED";
  8. public const string Cancelled = "CANCELLED";
  9. }
  10. public sealed class S1MdpAlreadyRunningException : Exception
  11. {
  12. public S1MdpAlreadyRunningException()
  13. : base("S1 数据重算正在执行,请勿重复提交")
  14. {
  15. }
  16. }
  17. public interface IS1MdpFullRunLease : IAsyncDisposable
  18. {
  19. Task HeartbeatAsync(CancellationToken ct = default);
  20. }
  21. public interface IS1MdpFullRunLock
  22. {
  23. Task<IS1MdpFullRunLease?> TryAcquireAsync(string lockName, string holderId, long? jobId = null, CancellationToken ct = default);
  24. }
  25. public sealed class S1MdpFullRunLock : IS1MdpFullRunLock, ITransient
  26. {
  27. public const string LegacyLockName = "S1_MDP_FULL";
  28. public static readonly TimeSpan StaleAfter = TimeSpan.FromMinutes(30);
  29. private readonly ISqlSugarClient _db;
  30. public S1MdpFullRunLock(ISqlSugarClient db) => _db = db;
  31. public async Task<IS1MdpFullRunLease?> TryAcquireAsync(string lockName, string holderId, long? jobId = null, CancellationToken ct = default)
  32. {
  33. if (string.IsNullOrWhiteSpace(lockName))
  34. throw new InvalidOperationException("S1 全量锁名不能为空");
  35. var now = DateTime.Now;
  36. var staleBefore = now - StaleAfter;
  37. await _db.Ado.ExecuteCommandAsync(
  38. """
  39. INSERT INTO ado_s1_mdp_run_lock (lock_name, status)
  40. SELECT @LockName, 'FREE'
  41. FROM DUAL
  42. WHERE NOT EXISTS (SELECT 1 FROM ado_s1_mdp_run_lock WHERE lock_name=@LockName)
  43. """,
  44. new SugarParameter("@LockName", lockName));
  45. var n = await _db.Ado.ExecuteCommandAsync(
  46. """
  47. UPDATE ado_s1_mdp_run_lock
  48. SET holder_id=@HolderId, job_id=@JobId, status='HELD',
  49. acquired_at=@Now, heartbeat_at=@Now
  50. WHERE lock_name=@LockName
  51. AND (status='FREE' OR heartbeat_at IS NULL OR heartbeat_at < @StaleBefore)
  52. """,
  53. new SugarParameter("@HolderId", holderId),
  54. new SugarParameter("@JobId", (object?)jobId ?? DBNull.Value),
  55. new SugarParameter("@Now", now),
  56. new SugarParameter("@LockName", lockName),
  57. new SugarParameter("@StaleBefore", staleBefore));
  58. if (n <= 0)
  59. return null;
  60. return new DbLease(_db, lockName, holderId);
  61. }
  62. private sealed class DbLease : IS1MdpFullRunLease
  63. {
  64. private readonly ISqlSugarClient _db;
  65. private readonly string _lockName;
  66. private readonly string _holderId;
  67. private int _released;
  68. public DbLease(ISqlSugarClient db, string lockName, string holderId)
  69. {
  70. _db = db;
  71. _lockName = lockName;
  72. _holderId = holderId;
  73. }
  74. public async Task HeartbeatAsync(CancellationToken ct = default)
  75. {
  76. using var db = _db.CopyNew();
  77. await db.Ado.ExecuteCommandAsync(
  78. """
  79. UPDATE ado_s1_mdp_run_lock
  80. SET heartbeat_at=@Now
  81. WHERE lock_name=@LockName AND holder_id=@HolderId AND status='HELD'
  82. """,
  83. new SugarParameter("@Now", DateTime.Now),
  84. new SugarParameter("@LockName", _lockName),
  85. new SugarParameter("@HolderId", _holderId));
  86. }
  87. public async ValueTask DisposeAsync()
  88. {
  89. if (Interlocked.Exchange(ref _released, 1) == 1)
  90. return;
  91. try
  92. {
  93. using var db = _db.CopyNew();
  94. await db.Ado.ExecuteCommandAsync(
  95. """
  96. UPDATE ado_s1_mdp_run_lock
  97. SET status='FREE', holder_id=NULL, job_id=NULL
  98. WHERE lock_name=@LockName AND holder_id=@HolderId
  99. """,
  100. new SugarParameter("@LockName", _lockName),
  101. new SugarParameter("@HolderId", _holderId));
  102. }
  103. catch
  104. {
  105. // 释放失败留给租约过期
  106. }
  107. }
  108. }
  109. }
  110. public sealed class InMemoryS1MdpFullRunLock : IS1MdpFullRunLock
  111. {
  112. private readonly Dictionary<string, string> _holders = new(StringComparer.Ordinal);
  113. public Task<IS1MdpFullRunLease?> TryAcquireAsync(string lockName, string holderId, long? jobId = null, CancellationToken ct = default)
  114. {
  115. lock (this)
  116. {
  117. if (_holders.ContainsKey(lockName))
  118. return Task.FromResult<IS1MdpFullRunLease?>(null);
  119. _holders[lockName] = holderId;
  120. return Task.FromResult<IS1MdpFullRunLease?>(new MemoryLease(this, lockName));
  121. }
  122. }
  123. private sealed class MemoryLease : IS1MdpFullRunLease
  124. {
  125. private readonly InMemoryS1MdpFullRunLock _owner;
  126. private readonly string _lockName;
  127. public MemoryLease(InMemoryS1MdpFullRunLock owner, string lockName)
  128. {
  129. _owner = owner;
  130. _lockName = lockName;
  131. }
  132. public Task HeartbeatAsync(CancellationToken ct = default) => Task.CompletedTask;
  133. public ValueTask DisposeAsync()
  134. {
  135. lock (_owner) _owner._holders.Remove(_lockName);
  136. return ValueTask.CompletedTask;
  137. }
  138. }
  139. }