NeutralRequiredColumnsMonitor.cs 3.4 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879
  1. using Microsoft.Extensions.Logging;
  2. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  3. /// <summary>
  4. /// 运行期核对中立层必填列是否整列为空。
  5. /// <para>
  6. /// <see cref="NeutralRequiredColumns"/> 是「投影 SQL 有没有写这一列」的静态守卫,守不住
  7. /// 「写了但取到的全是 NULL」和「这批存量是加列之前物化的」两种情况。而整列为空不会报错,
  8. /// 只会让依赖它的 KPI 一起算出 NO_DATA —— 2026-09 S5/S7 全租户掉数就是这么发生的。
  9. /// 故物化之后必须实测一次,把整列为空变成一条可查询的告警。
  10. /// </para>
  11. /// </summary>
  12. public sealed class NeutralRequiredColumnsMonitor : ITransient
  13. {
  14. private readonly ISqlSugarClient _db;
  15. private readonly ILogger _logger;
  16. public NeutralRequiredColumnsMonitor(ISqlSugarClient db, ILoggerFactory loggerFactory)
  17. {
  18. _db = db;
  19. _logger = loggerFactory.CreateLogger(nameof(NeutralRequiredColumnsMonitor));
  20. }
  21. /// <summary>返回该租户该表里「有行但整列为空」的必填列名。顺带落 mdp_source_gate_log 与错误日志。</summary>
  22. public async Task<IReadOnlyList<string>> AssertAsync(
  23. long tenantId, string stdObject, string table, string? batchId, CancellationToken cancellationToken = default)
  24. {
  25. if (tenantId <= 0) return [];
  26. var required = NeutralRequiredColumns.For(stdObject).Select(c => c.Name).ToList();
  27. if (required.Count == 0) return [];
  28. cancellationToken.ThrowIfCancellationRequested();
  29. var present = await _db.Ado.SqlQueryAsync<string>(
  30. """
  31. SELECT COLUMN_NAME FROM information_schema.COLUMNS
  32. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@table
  33. """,
  34. new { table });
  35. var columns = required
  36. .Where(c => present.Contains(c, StringComparer.OrdinalIgnoreCase))
  37. .ToList();
  38. if (columns.Count == 0) return [];
  39. // 列名取自 NeutralRequiredColumns 常量,不来自外部输入,故可直接拼进 SQL。
  40. var counters = string.Join(", ", columns.Select(c => $"SUM(`{c}` IS NOT NULL) AS `{c}`"));
  41. var stats = await _db.Ado.GetDataTableAsync(
  42. $"SELECT COUNT(*) AS `__total`, {counters} FROM `{table}` WHERE tenant_id=@tenant",
  43. new { tenant = tenantId });
  44. if (stats.Rows.Count == 0) return [];
  45. var row = stats.Rows[0];
  46. if (Convert.ToInt64(row["__total"]) <= 0) return [];
  47. var empty = columns
  48. .Where(c => row[c] == DBNull.Value || Convert.ToInt64(row[c]) == 0)
  49. .ToList();
  50. if (empty.Count == 0) return empty;
  51. _logger.LogError(
  52. "中立层必填列整列为空 tenant={Tenant} object={Object} table={Table} columns={Columns}",
  53. tenantId, stdObject, table, string.Join(",", empty));
  54. await _db.Ado.ExecuteCommandAsync(
  55. """
  56. INSERT INTO mdp_source_gate_log
  57. (tenant_id, std_object, source_system, gate_reason, row_count, sample_keys, sync_batch_id)
  58. VALUES (@tenant, @obj, '', 'REQUIRED_COLUMN_EMPTY', @rows, @sample, @batch)
  59. """,
  60. new
  61. {
  62. tenant = tenantId,
  63. obj = stdObject,
  64. rows = empty.Count,
  65. sample = string.Join(",", empty),
  66. batch = batchId ?? ""
  67. });
  68. return empty;
  69. }
  70. }