S8MdpAuthorityHealthResolver.cs 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360
  1. using Admin.NET.Plugin.AiDOP.Service.S8.Rules.DataAccess;
  2. using Microsoft.Extensions.Logging;
  3. using SqlSugar;
  4. using System.Text.RegularExpressions;
  5. namespace Admin.NET.Plugin.AiDOP.Service.S8.Rules.Health;
  6. /// <summary>
  7. /// <see cref="IS8AuthorityHealthResolver"/> 的数据中台适配实现:
  8. /// 从 <c>mdp_transform_run_log</c> 采集生产运行事实,必要时再校验快照表的发布不变量,
  9. /// 然后交给纯函数 <see cref="S8AuthorityHealthEvaluator"/> 判定。
  10. ///
  11. /// <para>本类是「临时适配」:中台一旦提供数据集级 freshness 接口,
  12. /// 换掉本类即可,端口与调用方不动。</para>
  13. /// </summary>
  14. public class S8MdpAuthorityHealthResolver : IS8AuthorityHealthResolver, ITransient
  15. {
  16. /// <summary>合法标识符(表名 / 列名)。声明值虽由开发者书写,仍不直接拼进 SQL 而先做白名单校验。</summary>
  17. private static readonly Regex IdentifierPattern = new("^[A-Za-z_][A-Za-z0-9_]*$", RegexOptions.Compiled);
  18. private readonly ISqlSugarClient _db;
  19. private readonly IS8DatasetCatalog _datasetCatalog;
  20. private readonly ILogger<S8MdpAuthorityHealthResolver> _logger;
  21. public S8MdpAuthorityHealthResolver(
  22. ISqlSugarClient db,
  23. IS8DatasetCatalog datasetCatalog,
  24. ILogger<S8MdpAuthorityHealthResolver> logger)
  25. {
  26. _db = db;
  27. _datasetCatalog = datasetCatalog;
  28. _logger = logger;
  29. }
  30. public async Task<S8AuthorityHealthResult> ResolveAsync(long tenantId, string datasetCode)
  31. {
  32. var observedAt = DateTime.Now;
  33. try
  34. {
  35. var spec = _datasetCatalog.Find(datasetCode)?.AuthoritySpec;
  36. // 未声明 AuthoritySpec / 声明不完整 → 交给纯函数判 UNKNOWN,不在这里编造默认值。
  37. if (spec is null || string.IsNullOrWhiteSpace(spec.ProducerJobCode))
  38. {
  39. return S8AuthorityHealthEvaluator.Evaluate(new S8AuthorityObservation
  40. {
  41. DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt, SpecFound = false
  42. });
  43. }
  44. // 生产者不可信时不必查库 —— 结论与运行日志无关。
  45. if (!spec.ProducerTrusted)
  46. {
  47. return S8AuthorityHealthEvaluator.Evaluate(new S8AuthorityObservation
  48. {
  49. DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt,
  50. SpecFound = true, ProducerTrusted = false,
  51. UntrustedReasonCode = spec.UntrustedReasonCode,
  52. StaleWindow = spec.StaleWindow
  53. });
  54. }
  55. var lookbackFrom = observedAt - spec.LookbackWindow;
  56. var lastSuccess = await QueryLastSuccessAsync(
  57. tenantId, spec.ProducerJobCode, lookbackFrom, spec.RequiresPublicationEvidence);
  58. // 需要发布证据却一条都没有:跑批可能成功过,但没人能指认「当前快照」是哪一批。
  59. // 这与「压根没成功过」是两件事,故给独立原因码;结论同为 UNKNOWN、拦恢复。
  60. // 判定放在这里而不是纯函数里,是因为「有没有 publish」是一个查询事实,
  61. // 而 S8AuthorityHealthEvaluator 的契约是不读行、不读库。
  62. if (spec.RequiresPublicationEvidence && lastSuccess is null)
  63. {
  64. var anySuccess = await QueryLastSuccessAsync(
  65. tenantId, spec.ProducerJobCode, lookbackFrom, requirePublication: false);
  66. if (anySuccess is not null)
  67. {
  68. return new S8AuthorityHealthResult
  69. {
  70. DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt,
  71. State = S8AuthorityHealthState.Unknown,
  72. ReasonCode = S8AuthorityHealthReason.NoPublication,
  73. Reason = "回溯窗口内有成功运行,但没有一次带发布证据,无法指认当前快照"
  74. };
  75. }
  76. }
  77. // 锚点:有成功则以其结束时刻为界,否则以回溯窗口起点为界。
  78. var anchor = lastSuccess?.EndTime ?? lookbackFrom;
  79. var newerFailure = await QueryNewerTerminalFailureAsync(tenantId, spec.ProducerJobCode, anchor);
  80. var newerInFlight = await QueryNewerInFlightAsync(tenantId, spec.ProducerJobCode, anchor);
  81. var (invariant, detail) = await ResolveSnapshotInvariantAsync(tenantId, spec, lastSuccess);
  82. return S8AuthorityHealthEvaluator.Evaluate(new S8AuthorityObservation
  83. {
  84. DatasetCode = datasetCode,
  85. TenantId = tenantId,
  86. ObservedAt = observedAt,
  87. SpecFound = true,
  88. ProducerTrusted = true,
  89. LastSuccessEndAt = lastSuccess?.EndTime,
  90. HasNewerTerminalFailure = newerFailure is not null,
  91. NewerTerminalFailureStatus = newerFailure?.Status,
  92. OldestNewerInFlightStartAt = newerInFlight,
  93. SnapshotInvariant = invariant,
  94. SnapshotDetail = detail,
  95. StaleWindow = spec.StaleWindow
  96. });
  97. }
  98. catch (Exception ex)
  99. {
  100. // fail-safe:探测自身失败一律判 UNKNOWN、拦截恢复。绝不因为查不动库就放行。
  101. _logger.LogWarning(ex,
  102. "authority_health_probe_failed tenant={Tenant} dataset={Dataset}", tenantId, datasetCode);
  103. return new S8AuthorityHealthResult
  104. {
  105. DatasetCode = datasetCode, TenantId = tenantId, ObservedAt = observedAt,
  106. State = S8AuthorityHealthState.Unknown,
  107. ReasonCode = S8AuthorityHealthReason.ResolverFailed,
  108. Reason = $"Authority 健康探测失败:{ex.GetType().Name}"
  109. };
  110. }
  111. }
  112. /// <summary>
  113. /// 最近一次<b>终态成功</b>的生产运行。
  114. ///
  115. /// <para><b>三处刻意为之,改动前先读完</b>:</para>
  116. /// <list type="number">
  117. /// <item>租户条件是严格相等,<b>绝不能</b>写成 <c>(tenant_id = @TenantId OR tenant_id = 0)</c>。
  118. /// 仓内 <c>MdpMonitorService.BuildMdpRunLogTenantWhere</c> 就是后者,本类不得复用:
  119. /// 该表确实存在 <c>tenant_id = 0</c> 的平台行(实测 S1 有 111 条、
  120. /// 且 <c>S5_PURCHASE_RECEIPT_MDP_SYNC</c> 有 325 条平台级 SUCCESS),
  121. /// 一条这样的行会让所有租户同时判成健康。</item>
  122. /// <item>必须显式过滤终态。该表有 298 条永不回收的 RUNNING 行(S1 占 42 条,最老 570 小时),
  123. /// 裸 <c>ORDER BY start_time DESC LIMIT 1</c> 会取到孤儿行,其 <c>end_time</c> 为 NULL,
  124. /// 后续新鲜度运算全部失效。</item>
  125. /// <item>按 <c>start_time</c> 排序(走 <c>idx_job_start</c> 的反向索引扫描,无 filesort),
  126. /// 但新鲜度用取到那行的 <c>end_time</c> 计算 —— 那才是 Authority 真正变成当前态的时刻。</item>
  127. /// </list>
  128. /// </summary>
  129. /// <param name="requirePublication">
  130. /// true 时只认<b>带发布证据</b>的成功运行。用于 <c>RequiresPublicationEvidence</c> 的数据集:
  131. /// 增量跑批(<c>RunInboundAsync</c>)同样以 SUCCESS 收尾,但它只处理变更行、不发布当前态,
  132. /// 因此不带 publish。若不加这层过滤,一次增量成功就会把「最近一次成功」的锚点从上一轮
  133. /// 完整发布挪走,新鲜度与失败锚点都会算错。
  134. /// </param>
  135. private async Task<RunRow?> QueryLastSuccessAsync(
  136. long tenantId, string jobCode, DateTime lookbackFrom, bool requirePublication)
  137. {
  138. var sql =
  139. $"""
  140. SELECT r.batch_id AS BatchId, r.status AS Status, r.start_time AS StartTime, r.end_time AS EndTime,
  141. JSON_UNQUOTE(JSON_EXTRACT(r.summary_json, '$.publish.batchId')) AS PublishedBatchId,
  142. JSON_EXTRACT(r.summary_json, '$.publish.currentRows') AS PublishedCurrentRows
  143. FROM mdp_transform_run_log r
  144. WHERE r.job_code = @JobCode
  145. AND r.tenant_id = @TenantId
  146. AND r.status = 'SUCCESS'
  147. AND r.end_time IS NOT NULL
  148. AND r.start_time >= @LookbackFrom
  149. {(requirePublication
  150. ? "AND JSON_EXTRACT(r.summary_json, '$.publish.batchId') IS NOT NULL"
  151. : string.Empty)}
  152. ORDER BY r.start_time DESC, r.id DESC
  153. LIMIT 1
  154. """;
  155. var rows = await _db.Ado.SqlQueryAsync<RunRow>(sql, new SugarParameter[]
  156. {
  157. new("@JobCode", jobCode), new("@TenantId", tenantId), new("@LookbackFrom", lookbackFrom)
  158. });
  159. return rows.FirstOrDefault();
  160. }
  161. /// <summary>
  162. /// 锚点之后是否出现终态非成功的运行。
  163. /// <para>用 <c>status &lt;&gt; 'SUCCESS'</c> 而非 <c>IN ('FAILED', ...)</c>:
  164. /// 将来若新增未知的终态状态,它会被归为「失败」而不是被静默当成健康 —— 失败方向必须保守。</para>
  165. /// </summary>
  166. private async Task<RunRow?> QueryNewerTerminalFailureAsync(long tenantId, string jobCode, DateTime anchor)
  167. {
  168. const string sql =
  169. """
  170. SELECT r.batch_id AS BatchId, r.status AS Status, r.start_time AS StartTime, r.end_time AS EndTime
  171. FROM mdp_transform_run_log r
  172. WHERE r.job_code = @JobCode
  173. AND r.tenant_id = @TenantId
  174. AND r.end_time IS NOT NULL
  175. AND r.status <> 'SUCCESS'
  176. AND r.start_time > @Anchor
  177. ORDER BY r.start_time DESC, r.id DESC
  178. LIMIT 1
  179. """;
  180. var rows = await _db.Ado.SqlQueryAsync<RunRow>(sql, new SugarParameter[]
  181. {
  182. new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor)
  183. });
  184. return rows.FirstOrDefault();
  185. }
  186. /// <summary>
  187. /// 锚点之后是否存在尚未终态的运行,取最早那条的开始时刻。
  188. /// <para><c>end_time IS NULL</c> 与 <c>status = 'RUNNING'</c> 在该表上实测完全等价(298/298 双向成立),
  189. /// 用前者是因为它对将来新增的非终态状态同样成立。</para>
  190. /// </summary>
  191. private async Task<DateTime?> QueryNewerInFlightAsync(long tenantId, string jobCode, DateTime anchor)
  192. {
  193. const string sql =
  194. """
  195. SELECT MIN(r.start_time) AS StartTime
  196. FROM mdp_transform_run_log r
  197. WHERE r.job_code = @JobCode
  198. AND r.tenant_id = @TenantId
  199. AND r.end_time IS NULL
  200. AND r.start_time > @Anchor
  201. """;
  202. var rows = await _db.Ado.SqlQueryAsync<RunRow>(sql, new SugarParameter[]
  203. {
  204. new("@JobCode", jobCode), new("@TenantId", tenantId), new("@Anchor", anchor)
  205. });
  206. return rows.FirstOrDefault()?.StartTime;
  207. }
  208. /// <summary>
  209. /// 快照型 Authority 的发布不变量。
  210. ///
  211. /// <para>当前批次数 = 0 时返回 <see cref="S8SnapshotInvariant.Unconfirmed"/> 而非「健康的空」:
  212. /// 生产侧的发布语句在本轮零行时命中 0 行、在库里不留任何 batch marker,
  213. /// 于是「源侧本轮确实没有数据」与「发布压根没跑成」字面等同。
  214. /// 这一点必须保守 —— 猜成健康会让真实预警被整体判成已恢复。
  215. /// 消除它需要生产侧补一个发布标记,属下一批的工作。</para>
  216. ///
  217. /// <para>作用域只按租户。快照的真实作用域是 (租户, 工厂),但运行日志表没有 factory 列,
  218. /// 无法把两侧对齐;若某租户将来出现多工厂,本查询会看到多个当前批次并判 VIOLATED ——
  219. /// 是拦截而非放行,失败方向安全。</para>
  220. /// </summary>
  221. private async Task<(string Invariant, string? Detail)> ResolveSnapshotInvariantAsync(
  222. long tenantId, S8AuthoritySpec spec, RunRow? lastSuccess)
  223. {
  224. var lastSuccessBatchId = lastSuccess?.BatchId;
  225. if (!string.Equals(spec.AuthorityKind, S8AuthorityKind.PublishedSnapshot, StringComparison.Ordinal))
  226. return (S8SnapshotInvariant.NotApplicable, null);
  227. if (string.IsNullOrWhiteSpace(spec.SnapshotTable))
  228. return (S8SnapshotInvariant.Violated, "AuthorityKind 声明为 PUBLISHED_SNAPSHOT 但未声明 SnapshotTable");
  229. var table = RequireIdentifier(spec.SnapshotTable, nameof(spec.SnapshotTable));
  230. var batchCol = RequireIdentifier(spec.SnapshotBatchColumn, nameof(spec.SnapshotBatchColumn));
  231. var flagCol = RequireIdentifier(spec.SnapshotCurrentFlagColumn, nameof(spec.SnapshotCurrentFlagColumn));
  232. var batchSql =
  233. $"""
  234. SELECT COUNT(DISTINCT s.`{batchCol}`) AS CurrentBatches, MIN(s.`{batchCol}`) AS BatchId
  235. FROM `{table}` s
  236. WHERE s.tenant_id = @TenantId AND s.`{flagCol}` = 1
  237. """;
  238. var batchRows = await _db.Ado.SqlQueryAsync<SnapshotRow>(
  239. batchSql, new SugarParameter[] { new("@TenantId", tenantId) });
  240. var batch = batchRows.FirstOrDefault();
  241. var currentBatches = batch?.CurrentBatches ?? 0;
  242. if (currentBatches == 0)
  243. {
  244. // 零个当前批次有两种截然不同的成因,DWD 表本身分不出来:
  245. // ① 源侧本轮确实没有数据 —— 发布语句命中 0 行,成功发布了一个空快照(合法)
  246. // ② 发布压根没跑成 —— 上一批次被退休了,新批次没被翻牌(故障)
  247. // 生产侧的发布证据是唯一能区分二者的东西,且它与 status='SUCCESS' 写在同一条 UPDATE 里,
  248. // 不存在「成功了但没有证据」的中间态。
  249. var publishedBatchId = lastSuccess?.PublishedBatchId;
  250. var publishedRows = lastSuccess?.PublishedCurrentRows;
  251. if (!string.IsNullOrWhiteSpace(publishedBatchId)
  252. && string.Equals(publishedBatchId, lastSuccessBatchId, StringComparison.Ordinal)
  253. && publishedRows == 0)
  254. {
  255. return (S8SnapshotInvariant.Satisfied, null);
  256. }
  257. // 没有发布证据(历史行,或生产侧尚未升级)→ 维持保守判定,与升级前行为完全一致。
  258. var why = string.IsNullOrWhiteSpace(publishedBatchId)
  259. ? "最近一次成功运行没有留下发布证据"
  260. : $"发布证据指向批次 {publishedBatchId}(当前行数 {publishedRows?.ToString() ?? "未知"})," +
  261. $"与最近一次成功运行的批次 {lastSuccessBatchId} 不一致";
  262. return (S8SnapshotInvariant.Unconfirmed,
  263. $"{table} 在租户 {tenantId} 下当前批次数为 0,且{why}" +
  264. "(合法空快照与发布失败在库里字面等同,不猜)");
  265. }
  266. if (currentBatches > 1)
  267. return (S8SnapshotInvariant.Violated,
  268. $"{table} 在租户 {tenantId} 下同时存在 {currentBatches} 个当前批次");
  269. if (!string.IsNullOrWhiteSpace(lastSuccessBatchId)
  270. && !string.Equals(batch!.BatchId, lastSuccessBatchId, StringComparison.Ordinal))
  271. {
  272. return (S8SnapshotInvariant.Violated,
  273. $"{table} 的当前批次 {batch.BatchId} 与最近一次成功运行的批次 {lastSuccessBatchId} 不一致");
  274. }
  275. // 必填列校验用绝对计数,空快照上天然真空成立。
  276. foreach (var column in spec.SnapshotRequiredColumns)
  277. {
  278. var col = RequireIdentifier(column, nameof(spec.SnapshotRequiredColumns));
  279. var nullSql =
  280. $"""
  281. SELECT COUNT(*) AS NullCount
  282. FROM `{table}` s
  283. WHERE s.tenant_id = @TenantId AND s.`{flagCol}` = 1 AND s.`{col}` IS NULL
  284. """;
  285. var nullRows = await _db.Ado.SqlQueryAsync<SnapshotRow>(
  286. nullSql, new SugarParameter[] { new("@TenantId", tenantId) });
  287. var nullCount = nullRows.FirstOrDefault()?.NullCount ?? 0;
  288. if (nullCount > 0)
  289. return (S8SnapshotInvariant.Violated,
  290. $"{table} 当前快照中 {col} 为空的行数 = {nullCount}");
  291. }
  292. return (S8SnapshotInvariant.Satisfied, null);
  293. }
  294. private static string RequireIdentifier(string value, string field)
  295. {
  296. if (!IdentifierPattern.IsMatch(value))
  297. throw new InvalidOperationException($"AuthoritySpec.{field} 不是合法标识符:{value}");
  298. return value;
  299. }
  300. /// <summary>运行日志投影行。属性名与 SELECT 别名一一对应。</summary>
  301. private sealed class RunRow
  302. {
  303. public string? BatchId { get; set; }
  304. public string? Status { get; set; }
  305. public DateTime? StartTime { get; set; }
  306. public DateTime? EndTime { get; set; }
  307. /// <summary>生产侧写入的发布证据;历史行没有该字段,取出为 null。</summary>
  308. public string? PublishedBatchId { get; set; }
  309. /// <summary>发布后的当前行数。0 是合法值(成功发布了一个空快照)。</summary>
  310. public int? PublishedCurrentRows { get; set; }
  311. }
  312. /// <summary>快照校验投影行。</summary>
  313. private sealed class SnapshotRow
  314. {
  315. public int CurrentBatches { get; set; }
  316. public string? BatchId { get; set; }
  317. public int NullCount { get; set; }
  318. }
  319. }