KpiDimensionRunService.cs 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208
  1. using System.Security.Cryptography;
  2. using System.Text;
  3. using Admin.NET.Core;
  4. using Admin.NET.Plugin.AiDOP.Entity;
  5. using SqlSugar;
  6. namespace Admin.NET.Plugin.AiDOP.SmartOps;
  7. /// <summary>维度执行结果概要(供调度/接口回执)。</summary>
  8. public sealed class KpiDimensionRunResult
  9. {
  10. public string Status { get; set; } = string.Empty;
  11. public int RowCount { get; set; }
  12. public string? ErrorCode { get; set; }
  13. public string? ErrorMessage { get; set; }
  14. }
  15. /// <summary>
  16. /// KPI 维度执行通用入口:读当前生效维度配置 → 校验绑定汇总版本一致 → 只读执行 DIMENSION_SQL
  17. /// → 事务 FULL REPLACE 写维度结果表 → 写运行日志。禁止 per-KPI 分支。
  18. /// </summary>
  19. public sealed class KpiDimensionRunService : ITransient
  20. {
  21. private readonly ISqlSugarClient _db;
  22. private readonly KpiDimensionSqlExecutor _executor;
  23. private readonly AdoSmartOpsKpiDimensionConfigService _dimensionConfig;
  24. private readonly AdoSmartOpsKpiCalcConfigService _summaryConfig;
  25. public KpiDimensionRunService(
  26. ISqlSugarClient db, KpiDimensionSqlExecutor executor,
  27. AdoSmartOpsKpiDimensionConfigService dimensionConfig, AdoSmartOpsKpiCalcConfigService summaryConfig)
  28. {
  29. _db = db;
  30. _executor = executor;
  31. _dimensionConfig = dimensionConfig;
  32. _summaryConfig = summaryConfig;
  33. }
  34. /// <summary>
  35. /// 执行某 KPI 当前生效维度配置。SUMMARY_ONLY / 无配置 / 版本不匹配均只记日志、不写结果。
  36. /// </summary>
  37. public async Task<KpiDimensionRunResult> RunDimensionAsync(
  38. string metricCode, string moduleCode, long tenantId, DateTime valueDate,
  39. string batchId, string triggerType, CancellationToken ct)
  40. {
  41. var startedAt = DateTime.Now;
  42. var bizDate = valueDate.Date;
  43. var active = await _dimensionConfig.GetActiveAsync(tenantId, metricCode);
  44. if (active == null)
  45. {
  46. await WriteRunLogAsync(null, metricCode, moduleCode, tenantId, batchId, triggerType, bizDate, startedAt,
  47. "NOT_CONFIGURED", 0, "NOT_CONFIGURED", "无生效维度配置", null);
  48. return new KpiDimensionRunResult { Status = "NOT_CONFIGURED", ErrorCode = "NOT_CONFIGURED" };
  49. }
  50. if (active.AggregationType == "SUMMARY_ONLY")
  51. {
  52. await WriteRunLogAsync(active, metricCode, moduleCode, tenantId, batchId, triggerType, bizDate, startedAt,
  53. "NOT_CONFIGURED", 0, "SUMMARY_ONLY", "维度配置为 SUMMARY_ONLY,不产生维度明细", null);
  54. return new KpiDimensionRunResult { Status = "NOT_CONFIGURED", ErrorCode = "SUMMARY_ONLY" };
  55. }
  56. // 绑定汇总版本一致性校验
  57. var summary = await _summaryConfig.GetActiveAsync(tenantId, metricCode);
  58. if (summary == null || summary.Id != active.SummaryConfigId || summary.VersionNo != active.SummaryConfigVersion)
  59. {
  60. await WriteRunLogAsync(active, metricCode, moduleCode, tenantId, batchId, triggerType, bizDate, startedAt,
  61. "VERSION_MISMATCH", 0, "VERSION_MISMATCH", "维度配置绑定的汇总版本与当前生效汇总版本不一致", null);
  62. return new KpiDimensionRunResult { Status = "VERSION_MISMATCH", ErrorCode = "VERSION_MISMATCH" };
  63. }
  64. var pars = AdoSmartOpsKpiDimensionConfigService.BuildRunParams(tenantId, moduleCode, metricCode, bizDate);
  65. var exec = await _executor.ExecuteAsync(active.DataSourceCode, active.SqlScript ?? "", active.TimeoutSeconds, false, pars, ct);
  66. if (exec.Status == "FAILED")
  67. {
  68. await WriteRunLogAsync(active, metricCode, moduleCode, tenantId, batchId, triggerType, bizDate, startedAt,
  69. "FAILED", exec.RowCount, exec.ErrorCode, exec.ErrorMessage, exec.SqlHash);
  70. return new KpiDimensionRunResult { Status = "FAILED", RowCount = exec.RowCount, ErrorCode = exec.ErrorCode, ErrorMessage = exec.ErrorMessage };
  71. }
  72. // FULL REPLACE:同 (tenant, metric, dimVersion) + 本次返回的各 value_date 范围内先删后插,事务提交。
  73. var now = DateTime.Now;
  74. var rows = exec.Rows.Select(r => ToEntity(r, active, tenantId, metricCode, moduleCode, batchId, now)).ToList();
  75. var dates = rows.Select(x => x.ValueDate.Date).Distinct().ToList();
  76. if (dates.Count == 0) dates.Add(bizDate); // NO_DATA 也清空当日该版本旧数据
  77. var tran = await _db.AsTenant().UseTranAsync(async () =>
  78. {
  79. await _db.Deleteable<AdoSmartOpsKpiDimensionValueDay>()
  80. .Where(x => x.TenantId == tenantId && x.MetricCode == metricCode
  81. && x.DimensionConfigVersion == active.DimensionConfigVersion
  82. && dates.Contains(x.ValueDate.Date))
  83. .ExecuteCommandAsync();
  84. if (rows.Count > 0)
  85. await _db.Insertable(rows).ExecuteCommandAsync();
  86. });
  87. if (!tran.IsSuccess)
  88. {
  89. await WriteRunLogAsync(active, metricCode, moduleCode, tenantId, batchId, triggerType, bizDate, startedAt,
  90. "FAILED", rows.Count, "WRITE_FAILED", tran.ErrorException?.Message, exec.SqlHash);
  91. return new KpiDimensionRunResult { Status = "FAILED", ErrorCode = "WRITE_FAILED", ErrorMessage = tran.ErrorException?.Message };
  92. }
  93. var status = rows.Count == 0 ? "NO_DATA" : "SUCCESS";
  94. await WriteRunLogAsync(active, metricCode, moduleCode, tenantId, batchId, triggerType, bizDate, startedAt,
  95. status, rows.Count, null, null, exec.SqlHash);
  96. return new KpiDimensionRunResult { Status = status, RowCount = rows.Count };
  97. }
  98. private AdoSmartOpsKpiDimensionValueDay ToEntity(
  99. KpiDimensionRow r, AdoSmartOpsKpiDimensionConfig cfg, long tenantId, string metricCode, string moduleCode,
  100. string batchId, DateTime now)
  101. {
  102. var valueDate = r.ValueDate == DateTime.MinValue ? DateTime.Today : r.ValueDate.Date;
  103. var e = new AdoSmartOpsKpiDimensionValueDay
  104. {
  105. TenantId = tenantId,
  106. ModuleCode = moduleCode,
  107. MetricCode = metricCode,
  108. SummaryConfigId = cfg.SummaryConfigId,
  109. SummaryConfigVersion = cfg.SummaryConfigVersion,
  110. DimensionConfigId = cfg.Id,
  111. DimensionConfigVersion = cfg.DimensionConfigVersion,
  112. DataSourceCode = cfg.DataSourceCode,
  113. ValueDate = valueDate,
  114. DimensionType = r.DimensionType,
  115. DimensionCode = r.DimensionCode,
  116. DimensionName = r.DimensionName,
  117. OrgId = r.OrgId,
  118. FactoryId = r.FactoryId ?? 1,
  119. MaterialCode = r.MaterialCode,
  120. WorkOrderNo = r.WorkOrderNo,
  121. CategoryCode = r.CategoryCode,
  122. WarehouseCode = r.WarehouseCode,
  123. OrderNo = r.OrderNo,
  124. CustomerCode = r.CustomerCode,
  125. ProductCode = r.ProductCode,
  126. EquipmentCode = r.EquipmentCode,
  127. MetricValue = r.MetricValue,
  128. Numerator = r.Numerator,
  129. Denominator = r.Denominator,
  130. SumValue = r.SumValue,
  131. SampleCount = r.SampleCount,
  132. SourceKey = r.SourceKey,
  133. BatchId = batchId,
  134. CalcTime = now,
  135. CreatedTime = now,
  136. };
  137. e.RowHash = ComputeRowHash(tenantId, metricCode, cfg.DimensionConfigVersion, valueDate, r.DimensionType, r.DimensionCode, r.SourceKey);
  138. return e;
  139. }
  140. private static string ComputeRowHash(long tenantId, string metricCode, int dimVersion, DateTime valueDate,
  141. string dimType, string dimCode, string? sourceKey)
  142. {
  143. var raw = $"{tenantId}|{metricCode}|{dimVersion}|{valueDate:yyyyMMdd}|{dimType}|{dimCode}|{sourceKey}";
  144. return Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(raw)));
  145. }
  146. private async Task WriteRunLogAsync(
  147. AdoSmartOpsKpiDimensionConfig? cfg, string metricCode, string moduleCode, long tenantId,
  148. string batchId, string triggerType, DateTime bizDate, DateTime startedAt,
  149. string status, int rowCount, string? errorCode, string? errorMessage, string? sqlHash)
  150. {
  151. try
  152. {
  153. var finishedAt = DateTime.Now;
  154. var log = new AdoSmartOpsKpiDimensionRunLog
  155. {
  156. TenantId = tenantId,
  157. BatchId = batchId,
  158. ModuleCode = moduleCode,
  159. MetricCode = metricCode,
  160. DimensionConfigId = cfg?.Id,
  161. DimensionConfigVersion = cfg?.DimensionConfigVersion,
  162. SummaryConfigVersion = cfg?.SummaryConfigVersion,
  163. DataSourceCode = cfg?.DataSourceCode,
  164. BizDate = bizDate.Date,
  165. StartedAt = startedAt,
  166. FinishedAt = finishedAt,
  167. DurationMs = (long)(finishedAt - startedAt).TotalMilliseconds,
  168. Status = status,
  169. RowCount = rowCount,
  170. ErrorCode = errorCode,
  171. ErrorMessage = errorMessage == null ? null : (errorMessage.Length > 480 ? errorMessage.Substring(0, 480) : errorMessage),
  172. SqlHash = sqlHash,
  173. TriggerType = triggerType,
  174. CreateTime = finishedAt,
  175. };
  176. var existing = await _db.Queryable<AdoSmartOpsKpiDimensionRunLog>().ClearFilter<ITenantIdFilter>()
  177. .Where(x => x.MetricCode == metricCode && x.BatchId == batchId).FirstAsync();
  178. if (existing != null)
  179. {
  180. log.Id = existing.Id;
  181. await _db.Updateable(log).ExecuteCommandAsync();
  182. }
  183. else
  184. {
  185. await _db.Insertable(log).ExecuteCommandAsync();
  186. }
  187. }
  188. catch
  189. {
  190. // 运行日志失败不影响主执行链路
  191. }
  192. }
  193. }