AidopT8KpiManualRefreshService.cs 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340
  1. using Admin.NET.Core.Service;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  3. using Admin.NET.Plugin.AiDOP.FinishedWarehouse;
  4. using Admin.NET.Plugin.AiDOP.Infrastructure;
  5. using Admin.NET.Plugin.AiDOP.Manufacturing;
  6. using Admin.NET.Plugin.AiDOP.MaterialWarehouse;
  7. using Microsoft.Extensions.Logging;
  8. using SqlSugar;
  9. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  10. /// <summary>
  11. /// S5/S6/S7 T8 KPI 手动刷新统一编排(看板"刷新"按钮真实跑批入口)。
  12. /// 单次刷新语义:抢全局刷新锁 → T8 基表入站(硬门,共享 6 表)→ 本模块 RunFull(MANUAL) → 汇总逐指标状态 → 释放锁。
  13. /// 与 Cron job_s5_t8_kpi_refresh 共享同一把锁(<see cref="RefreshLockKey"/>),保证手动刷新与定时跑批互斥、并发手动刷新排他。
  14. /// 入站为 S5/S6/S7 共用(T8BaseInboundMdpSyncService 6 张基表无法按模块拆),故刷任一模块均触发一次全量入站。
  15. /// S5_L1_002 / S5_L1_004 仍 legacy 直连 T8(本轮不改),随 S5 RunFull 一并执行,逐指标标 LEGACY_SUCCESS。
  16. ///
  17. /// 锁设计(加固):用持有者无关的硬 TTL 键(SysCacheService.Set 带过期),而非 BeginCacheLock。
  18. /// 原因:BeginCacheLock 在持有者线程存活期间会续租,一旦某次运行意外挂死(如源库慢查询无超时)会永久占锁、
  19. /// 阻塞后续所有刷新。改为硬 TTL 后:正常路径 finally 主动 Remove;即使挂死,键到 <see cref="RefreshLockTtlSeconds"/> 也自动过期。
  20. /// 配合 QueryT8Async 的 T8 命令超时(60s)与主库命令超时(30s),单次运行时长恒 &lt; TTL,正常不会误过期。
  21. /// </summary>
  22. public sealed class AidopT8KpiManualRefreshService : ITransient
  23. {
  24. /// <summary>手动刷新 + Cron 共享的全局互斥锁 key(覆盖 inbound + transform 整段)。</summary>
  25. public const string RefreshLockKey = "aidop:t8kpi:refresh:lock";
  26. /// <summary>刷新锁硬 TTL(秒,持有者无关):到点自动过期,即使某次运行意外挂死也不会永久阻塞后续刷新。
  27. /// 取值须大于加固后最坏运行时长(T8 查询各 ≤60s、主库命令 ≤30s,全程数分钟内),10 分钟为安全上限。</summary>
  28. public const int RefreshLockTtlSeconds = 600;
  29. // legacy 直连 T8 的 KPI(仅剩 S5_L1_004 TVF;S5_L1_002 已中台化),逐指标状态标 LEGACY_SUCCESS 以示区分。
  30. private static readonly HashSet<string> LegacyKpiCodes = new(StringComparer.OrdinalIgnoreCase) { "S5_L1_004" };
  31. private readonly T8BaseInboundMdpSyncService _inbound;
  32. private readonly S5MdpSyncTransformService _s5;
  33. private readonly S6MdpSyncTransformService _s6;
  34. private readonly S7MdpSyncTransformService _s7;
  35. private readonly SysCacheService _cache;
  36. private readonly ISqlSugarClient _db;
  37. private readonly ILogger<AidopT8KpiManualRefreshService> _logger;
  38. public AidopT8KpiManualRefreshService(
  39. T8BaseInboundMdpSyncService inbound,
  40. S5MdpSyncTransformService s5,
  41. S6MdpSyncTransformService s6,
  42. S7MdpSyncTransformService s7,
  43. SysCacheService cache,
  44. ISqlSugarClient db,
  45. ILogger<AidopT8KpiManualRefreshService> logger)
  46. {
  47. _inbound = inbound;
  48. _s5 = s5;
  49. _s6 = s6;
  50. _s7 = s7;
  51. _cache = cache;
  52. _db = db;
  53. _logger = logger;
  54. }
  55. public Task<AidopKpiRefreshResult> RunModuleRefreshAsync(string moduleCode, CancellationToken cancellationToken) =>
  56. RunModuleRefreshAsync(moduleCode, cancellationToken, null, null);
  57. /// <summary>执行一次模块级手动刷新(S5/S6/S7)。返回统一契约,绝不抛业务异常(取消除外)。</summary>
  58. public async Task<AidopKpiRefreshResult> RunModuleRefreshAsync(
  59. string moduleCode,
  60. CancellationToken cancellationToken,
  61. MdpRebuildScope? scope,
  62. Func<ModuleProgressUpdate, Task>? report)
  63. {
  64. var mc = string.IsNullOrWhiteSpace(moduleCode) ? string.Empty : moduleCode.Trim().ToUpperInvariant();
  65. var startedAt = DateTime.Now;
  66. var result = new AidopKpiRefreshResult
  67. {
  68. ModuleCode = mc,
  69. StartedAt = startedAt.ToString("yyyy-MM-dd HH:mm:ss")
  70. };
  71. if (mc is not ("S5" or "S6" or "S7"))
  72. {
  73. result.Ok = false;
  74. result.OverallStatus = "SKIPPED";
  75. result.Message = $"模块 {mc} 不支持 T8 KPI 手动刷新";
  76. Finish(result, startedAt);
  77. return result;
  78. }
  79. // 硬 TTL 锁:ExistKey 命中即判定已有刷新在跑(手动或定时);否则 Set(带 TTL) 占锁。
  80. // 单实例 + 人工点击/每日 5 次定时下,ExistKey→Set 的微小竞态概率可忽略;即便未主动释放,TTL 到点自动过期。
  81. if (_cache.ExistKey(RefreshLockKey))
  82. {
  83. result.Ok = false;
  84. result.OverallStatus = "REFRESHING";
  85. result.Message = "T8 KPI 刷新正在进行中(手动或定时任务占用),请稍后重试";
  86. Finish(result, startedAt);
  87. return result;
  88. }
  89. _cache.Set(RefreshLockKey, $"MANUAL:{mc}:{startedAt:yyyyMMddHHmmss}", TimeSpan.FromSeconds(RefreshLockTtlSeconds));
  90. try
  91. {
  92. await ReportAsync(report, new ModuleProgressUpdate(ModuleRebuildStages.T8Inbound, 1, 15, "正在执行 T8 基表入站"));
  93. try
  94. {
  95. var tenantId = scope?.TenantId ?? 0;
  96. var registered = await CountT8ObjectsAsync(tenantId);
  97. if (registered == 0)
  98. {
  99. result.Inbound.Ok = true;
  100. result.Inbound.Message = "本租户 S5-S7 未登记 T8,跳过 T8 入站";
  101. await ReportAsync(report, new ModuleProgressUpdate(
  102. ModuleRebuildStages.T8Inbound, 1, 40, result.Inbound.Message));
  103. }
  104. else
  105. {
  106. var ztid = await LoadT8ScopeAsync(tenantId);
  107. var inbound = await _inbound.RunInboundAsync(tenantId, true, null, cancellationToken, ztid);
  108. result.Inbound.Ok = true;
  109. result.Inbound.StageRows = inbound.Tables.Sum(t => t.StgRows);
  110. result.Inbound.StandardRows = inbound.Tables.Sum(t => t.StdRows);
  111. await ReportAsync(report, new ModuleProgressUpdate(
  112. ModuleRebuildStages.T8Inbound, 1, 40, "T8 基表入站完成", result.Inbound.StageRows, ModuleRebuildStages.T8Inbound));
  113. }
  114. }
  115. catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
  116. {
  117. throw;
  118. }
  119. catch (Exception ex)
  120. {
  121. _logger.LogError(ex, "[手动刷新] {Module} T8 基表入站失败,已跳过指标计算", mc);
  122. result.Ok = false;
  123. result.OverallStatus = "FAILED";
  124. result.Inbound.Ok = false;
  125. result.Inbound.Message = ex.Message;
  126. result.Message = "源数据入站失败,看板保留原有数据";
  127. Finish(result, startedAt);
  128. return result;
  129. }
  130. // ② 本模块 KPI 转换(RunFull 为全有全无:任一 KPI 异常即整体抛错)。
  131. string batchId;
  132. int dwdRows;
  133. int kpiRows;
  134. Dictionary<string, int> perKpiRows;
  135. List<string> denomStatus;
  136. await ReportAsync(report, new ModuleProgressUpdate(ModuleRebuildStages.KpiPreparing, 2, 45, "准备 KPI 计算", result.Inbound.StandardRows, ModuleRebuildStages.KpiPreparing));
  137. await ReportAsync(report, new ModuleProgressUpdate(ModuleRebuildStages.KpiCalculating, 3, 55, "正在计算 KPI"));
  138. try
  139. {
  140. switch (mc)
  141. {
  142. case "S5":
  143. var r5 = await _s5.RunFullAsync(cancellationToken, "MANUAL",
  144. option: scope == null ? null : new S5MdpRefreshOption { TargetTenantId = scope.TenantId, TargetFactoryId = scope.FactoryId });
  145. batchId = r5.BatchId;
  146. dwdRows = r5.DwdRows;
  147. kpiRows = r5.KpiRows;
  148. perKpiRows = r5.PerKpiKpiRows;
  149. denomStatus = r5.KpiDenominatorStatus;
  150. break;
  151. case "S6":
  152. var r6 = await _s6.RunFullAsync(cancellationToken, "MANUAL",
  153. option: scope == null ? null : new S6MdpRefreshOption { TargetTenantId = scope.TenantId, TargetFactoryId = scope.FactoryId });
  154. batchId = r6.BatchId;
  155. dwdRows = r6.DwdRows;
  156. kpiRows = r6.KpiRows;
  157. perKpiRows = r6.PerKpiKpiRows;
  158. denomStatus = r6.KpiDenominatorStatus;
  159. break;
  160. default:
  161. var r7 = await _s7.RunFullAsync(cancellationToken, "MANUAL",
  162. option: scope == null ? null : new S7MdpRefreshOption { TargetTenantId = scope.TenantId, TargetFactoryId = scope.FactoryId });
  163. batchId = r7.BatchId;
  164. dwdRows = r7.DwdRows;
  165. kpiRows = r7.KpiRows;
  166. perKpiRows = r7.PerKpiKpiRows;
  167. denomStatus = r7.KpiDenominatorStatus;
  168. break;
  169. }
  170. }
  171. catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
  172. {
  173. throw;
  174. }
  175. catch (Exception ex)
  176. {
  177. _logger.LogError(ex, "[手动刷新] {Module} KPI 转换失败", mc);
  178. result.Ok = false;
  179. result.OverallStatus = "FAILED";
  180. result.Transform.Ok = false;
  181. result.Transform.Message = ex.Message;
  182. result.Message = "KPI 计算失败(源数据已入站,看板保留原有数据)";
  183. Finish(result, startedAt);
  184. return result;
  185. }
  186. result.BatchId = batchId;
  187. result.Transform.Ok = true;
  188. result.Transform.DwdRows = dwdRows;
  189. result.Transform.KpiRows = kpiRows;
  190. // ③ 逐指标状态:denom=OK → SUCCESS(legacy→LEGACY_SUCCESS);denom=NO_* → NO_DATA(真实无数据,非伪造 0)。
  191. var denomMap = ParseDenominator(denomStatus);
  192. foreach (var kv in perKpiRows.OrderBy(k => k.Key, StringComparer.Ordinal))
  193. {
  194. var code = kv.Key;
  195. var denom = denomMap.TryGetValue(code, out var d) ? d : "OK";
  196. var ok = string.Equals(denom, "OK", StringComparison.OrdinalIgnoreCase);
  197. var status = ok
  198. ? (LegacyKpiCodes.Contains(code) ? "LEGACY_SUCCESS" : "SUCCESS")
  199. : "NO_DATA";
  200. result.PerKpi.Add(new AidopKpiRefreshPerKpi
  201. {
  202. MetricCode = code,
  203. Status = status,
  204. Rows = kv.Value,
  205. Message = ok ? null : denom
  206. });
  207. }
  208. result.SuccessCount = result.PerKpi.Count(p => p.Status is "SUCCESS" or "LEGACY_SUCCESS");
  209. result.NoDataCount = result.PerKpi.Count(p => p.Status == "NO_DATA");
  210. result.FailedCount = result.PerKpi.Count(p => p.Status is "FAILED" or "LEGACY_FAILED");
  211. result.Ok = true;
  212. // RunFull 为全有全无,成功路径下 FailedCount 恒为 0;PARTIAL_SUCCESS 仅为契约保留态。
  213. result.OverallStatus = result.FailedCount > 0 ? "PARTIAL_SUCCESS" : "SUCCESS";
  214. Finish(result, startedAt);
  215. result.Message = $"刷新完成:成功 {result.SuccessCount},无数据 {result.NoDataCount},失败 {result.FailedCount},耗时 {result.DurationMs / 1000.0:0.0} 秒";
  216. await ReportAsync(report, new ModuleProgressUpdate(
  217. ModuleRebuildStages.KpiCalculating, 3, 90, "KPI 计算完成", result.Transform.KpiRows, ModuleRebuildStages.KpiCalculating));
  218. await ReportAsync(report, new ModuleProgressUpdate(ModuleRebuildStages.Finalizing, 4, 99, "正在收尾"));
  219. return result;
  220. }
  221. finally
  222. {
  223. // 无论成功/失败/取消,主动释放锁(硬 TTL 是兜底,正常路径这里即时释放)。
  224. _cache.Remove(RefreshLockKey);
  225. }
  226. }
  227. private static async Task ReportAsync(Func<ModuleProgressUpdate, Task>? report, ModuleProgressUpdate update)
  228. {
  229. if (report == null) return;
  230. await report(update);
  231. }
  232. /// <summary>把 "S5_L1_001:OK" 形式的分母状态列表解析为 metricCode → status 字典。</summary>
  233. private static Dictionary<string, string> ParseDenominator(List<string> denom)
  234. {
  235. var map = new Dictionary<string, string>(StringComparer.OrdinalIgnoreCase);
  236. foreach (var item in denom)
  237. {
  238. if (string.IsNullOrWhiteSpace(item)) continue;
  239. var idx = item.IndexOf(':');
  240. if (idx <= 0) continue;
  241. map[item.Substring(0, idx).Trim()] = item.Substring(idx + 1).Trim();
  242. }
  243. return map;
  244. }
  245. private async Task<int> CountT8ObjectsAsync(long tenantId)
  246. {
  247. if (tenantId <= 0) return 0;
  248. return await _db.Ado.GetIntAsync(
  249. """
  250. SELECT COUNT(*) FROM mdp_tenant_std_source
  251. WHERE tenant_id=@t AND source_system='T8'
  252. AND std_object IN ('INV_TRANS','WO_LINE_PROD','WO_LINE_SALES','WO_BOM','EMPLOYEE','INV_BAL_MONTHLY','S6_REPORT','FQC_TASK')
  253. """,
  254. new SugarParameter("@t", tenantId));
  255. }
  256. private async Task<string?> LoadT8ScopeAsync(long tenantId)
  257. {
  258. var rows = await _db.Ado.SqlQueryAsync<string>(
  259. """
  260. SELECT source_scope FROM mdp_tenant_std_source
  261. WHERE tenant_id=@t AND source_system='T8' AND IFNULL(source_scope,'')<>''
  262. LIMIT 1
  263. """,
  264. new SugarParameter("@t", tenantId));
  265. return rows.FirstOrDefault();
  266. }
  267. private static void Finish(AidopKpiRefreshResult result, DateTime startedAt)
  268. {
  269. var finishedAt = DateTime.Now;
  270. result.FinishedAt = finishedAt.ToString("yyyy-MM-dd HH:mm:ss");
  271. result.DurationMs = (long)(finishedAt - startedAt).TotalMilliseconds;
  272. }
  273. }
  274. /// <summary>手动刷新统一返回契约(PascalCase 属性经框架序列化为 camelCase:Ok→ok、BatchId→batchId)。</summary>
  275. public sealed class AidopKpiRefreshResult
  276. {
  277. public bool Ok { get; set; }
  278. public string ModuleCode { get; set; } = string.Empty;
  279. /// <summary>整体状态:SUCCESS / PARTIAL_SUCCESS / FAILED / REFRESHING / SKIPPED。</summary>
  280. public string OverallStatus { get; set; } = string.Empty;
  281. public string? BatchId { get; set; }
  282. public string StartedAt { get; set; } = string.Empty;
  283. public string FinishedAt { get; set; } = string.Empty;
  284. public long DurationMs { get; set; }
  285. public AidopKpiRefreshInbound Inbound { get; set; } = new();
  286. public AidopKpiRefreshTransform Transform { get; set; } = new();
  287. public int SuccessCount { get; set; }
  288. public int NoDataCount { get; set; }
  289. public int FailedCount { get; set; }
  290. public List<AidopKpiRefreshPerKpi> PerKpi { get; set; } = new();
  291. public string Message { get; set; } = string.Empty;
  292. }
  293. public sealed class AidopKpiRefreshInbound
  294. {
  295. public bool Ok { get; set; }
  296. public int StageRows { get; set; }
  297. public int StandardRows { get; set; }
  298. public string? Message { get; set; }
  299. }
  300. public sealed class AidopKpiRefreshTransform
  301. {
  302. public bool Ok { get; set; }
  303. public int DwdRows { get; set; }
  304. public int KpiRows { get; set; }
  305. public string? Message { get; set; }
  306. }
  307. /// <summary>单指标刷新状态:SUCCESS / LEGACY_SUCCESS / NO_DATA / FAILED / LEGACY_FAILED。</summary>
  308. public sealed class AidopKpiRefreshPerKpi
  309. {
  310. public string MetricCode { get; set; } = string.Empty;
  311. public string Status { get; set; } = string.Empty;
  312. public int Rows { get; set; }
  313. public string? Message { get; set; }
  314. }