KpiCalcDispatcher.cs 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144
  1. using Admin.NET.Core;
  2. using Admin.NET.Plugin.AiDOP.Entity;
  3. using SqlSugar;
  4. namespace Admin.NET.Plugin.AiDOP.SmartOps;
  5. /// <summary>分发结果:引擎、是否应写 KPI 值、指标值、分母状态。</summary>
  6. public sealed class KpiCalcDispatchResult
  7. {
  8. public string EngineType { get; set; } = "LEGACY_CODE";
  9. /// <summary>是否写 KPI 日值。FAILED 时 false → 不写、保留上一成功值。</summary>
  10. public bool ShouldUpsert { get; set; } = true;
  11. public decimal? MetricValue { get; set; }
  12. /// <summary>分母/结果状态:OK / NO_DATA / NO_NUMERATOR / FAILED。</summary>
  13. public string DenominatorStatus { get; set; } = "OK";
  14. }
  15. /// <summary>
  16. /// KPI 计算运行时分发器:解析当前生效配置 → LEGACY_CODE(用旧代码值) / CONFIG_SQL(只读执行 SQL) / LEGACY_TVF(用旧代码值)。
  17. /// 无生效配置 → 回落 LEGACY_CODE(向后兼容)。CONFIG_SQL 失败**不 fallback**、不写值、返回 FAILED(保留上一成功值)。
  18. /// 每次执行落 ado_smart_ops_kpi_calc_run_log(含引擎/版本/SQL 指纹/分子分母/耗时/状态)。
  19. /// </summary>
  20. public sealed class KpiCalcDispatcher : ITransient
  21. {
  22. private readonly ISqlSugarClient _db;
  23. private readonly AdoSmartOpsKpiCalcConfigService _configService;
  24. private readonly KpiSqlReadOnlyExecutor _executor;
  25. public KpiCalcDispatcher(ISqlSugarClient db, AdoSmartOpsKpiCalcConfigService configService, KpiSqlReadOnlyExecutor executor)
  26. {
  27. _db = db;
  28. _configService = configService;
  29. _executor = executor;
  30. }
  31. public async Task<KpiCalcDispatchResult> DispatchAsync(
  32. string metricCode, string moduleCode, long tenantId, long factoryId,
  33. DateTime bizDate, DateTime periodStart, DateTime periodEnd, string sourceZtid,
  34. string batchId, string triggerType,
  35. decimal? legacyValue, string legacyDenomStatus, CancellationToken ct)
  36. {
  37. var startedAt = DateTime.Now;
  38. var config = await _configService.GetActiveAsync(tenantId, metricCode);
  39. var engine = config?.CalcEngineType ?? "LEGACY_CODE";
  40. var result = new KpiCalcDispatchResult { EngineType = engine };
  41. // 无配置 / LEGACY_CODE / LEGACY_TVF → 用旧代码算出的值(旧 Build 已计算并传入)
  42. if (config == null || engine == "LEGACY_CODE" || engine == "LEGACY_TVF")
  43. {
  44. result.MetricValue = legacyValue;
  45. result.DenominatorStatus = legacyDenomStatus;
  46. result.ShouldUpsert = true;
  47. await WriteRunLogAsync(config, metricCode, moduleCode, tenantId, batchId, triggerType, engine, bizDate,
  48. startedAt, legacyValue == null ? "NO_DATA" : "SUCCESS", legacyValue, null, null, 0, null, null, null);
  49. return result;
  50. }
  51. // CONFIG_SQL:只读执行已发布 SQL
  52. var pars = new KpiSqlRunParams
  53. {
  54. TenantId = tenantId, FactoryId = factoryId, ModuleCode = moduleCode, MetricCode = metricCode,
  55. BizDate = bizDate, PeriodStart = periodStart, PeriodEnd = periodEnd, SourceZtid = sourceZtid,
  56. };
  57. var exec = await _executor.ExecuteAsync(config.DataSourceCode, config.SqlScript ?? "", config.TimeoutSeconds, false, pars, ct);
  58. switch (exec.Status)
  59. {
  60. case "SUCCESS":
  61. result.MetricValue = exec.MetricValue;
  62. result.DenominatorStatus = "OK";
  63. result.ShouldUpsert = true;
  64. break;
  65. case "NO_DATA":
  66. result.MetricValue = null;
  67. result.DenominatorStatus = "NO_DATA";
  68. result.ShouldUpsert = true; // 遵循现有契约:NO_DATA 写 null,不伪造 0
  69. break;
  70. default: // FAILED
  71. result.MetricValue = null;
  72. result.DenominatorStatus = "FAILED";
  73. result.ShouldUpsert = false; // 不 fallback、不写错误新值、保留上一成功值
  74. break;
  75. }
  76. await WriteRunLogAsync(config, metricCode, moduleCode, tenantId, batchId, triggerType, engine, bizDate,
  77. startedAt, exec.Status, exec.MetricValue, exec.NumeratorValue, exec.DenominatorValue,
  78. exec.RowCount, exec.ErrorCode, exec.ErrorMessage, exec.SqlHash);
  79. return result;
  80. }
  81. private async Task WriteRunLogAsync(
  82. AdoSmartOpsKpiCalcConfig? config, string metricCode, string moduleCode, long tenantId,
  83. string batchId, string triggerType, string engine, DateTime bizDate, DateTime startedAt,
  84. string status, decimal? metricValue, decimal? numerator, decimal? denominator, int rowCount,
  85. string? errorCode, string? errorMessage, string? sqlHash)
  86. {
  87. try
  88. {
  89. var finishedAt = DateTime.Now;
  90. var log = new AdoSmartOpsKpiCalcRunLog
  91. {
  92. TenantId = tenantId,
  93. BatchId = batchId,
  94. ModuleCode = moduleCode,
  95. MetricCode = metricCode,
  96. ConfigId = config?.Id,
  97. VersionNo = config?.VersionNo,
  98. EngineType = engine,
  99. DataSourceCode = config?.DataSourceCode,
  100. BizDate = bizDate.Date,
  101. StartedAt = startedAt,
  102. FinishedAt = finishedAt,
  103. DurationMs = (long)(finishedAt - startedAt).TotalMilliseconds,
  104. Status = status,
  105. RowCount = rowCount,
  106. MetricValue = metricValue,
  107. NumeratorValue = numerator,
  108. DenominatorValue = denominator,
  109. ErrorCode = errorCode,
  110. ErrorMessage = errorMessage,
  111. SqlHash = sqlHash,
  112. ParameterSnapshot = $"{{\"tenant_id\":{tenantId},\"metric_code\":\"{metricCode}\",\"biz_date\":\"{bizDate:yyyy-MM-dd}\"}}",
  113. TriggerType = triggerType,
  114. CreateTime = finishedAt,
  115. };
  116. // 同批次同指标唯一:存在则更新、否则插入(避免撞唯一键 uk_kpi_calc_run_metric_batch)
  117. var existing = await _db.Queryable<AdoSmartOpsKpiCalcRunLog>().ClearFilter<ITenantIdFilter>()
  118. .Where(x => x.MetricCode == metricCode && x.BatchId == batchId).FirstAsync();
  119. if (existing != null)
  120. {
  121. log.Id = existing.Id;
  122. await _db.Updateable(log).ExecuteCommandAsync();
  123. }
  124. else
  125. {
  126. await _db.Insertable(log).ExecuteCommandAsync();
  127. }
  128. }
  129. catch
  130. {
  131. // 运行日志失败不影响主计算链路
  132. }
  133. }
  134. }