MdpNeutralSourceCleanup.cs 5.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145
  1. using Microsoft.Extensions.Logging;
  2. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  3. /// <summary>
  4. /// 删除某租户某标准对象里来源与登记不一致的中立行。不删贴源,不删目录之外的表。
  5. /// </summary>
  6. public sealed class MdpNeutralSourceCleanup : ITransient
  7. {
  8. private readonly ISqlSugarClient _db;
  9. private readonly ILogger _logger;
  10. public MdpNeutralSourceCleanup(ISqlSugarClient db, ILoggerFactory loggerFactory)
  11. {
  12. _db = db;
  13. _logger = loggerFactory.CreateLogger(nameof(MdpNeutralSourceCleanup));
  14. }
  15. public async Task<int> PurgeModuleAsync(long tenantId, string moduleCode, CancellationToken cancellationToken = default)
  16. {
  17. var total = 0;
  18. foreach (var def in MdpStdObjectCatalog.All.Where(d => d.Module == moduleCode))
  19. {
  20. cancellationToken.ThrowIfCancellationRequested();
  21. total += await PurgeObjectAsync(tenantId, def.Code, cancellationToken);
  22. }
  23. return total;
  24. }
  25. public async Task<int> PurgeObjectAsync(long tenantId, string stdObject, CancellationToken cancellationToken = default)
  26. {
  27. var def = MdpStdObjectCatalog.Find(stdObject);
  28. if (def == null || tenantId <= 0)
  29. return 0;
  30. var total = 0;
  31. foreach (var table in def.Tables)
  32. {
  33. cancellationToken.ThrowIfCancellationRequested();
  34. if (!await TableReadyAsync(table))
  35. continue;
  36. var docFilter = string.IsNullOrEmpty(def.DocType) || !await ColumnExistsAsync(table, "doc_type")
  37. ? ""
  38. : " AND t.doc_type=@doc";
  39. if (!await RegisteredSourceHasRowsAsync(tenantId, def.Code, table, docFilter, cancellationToken))
  40. continue;
  41. var sql = $"""
  42. DELETE t FROM `{table}` t
  43. INNER JOIN mdp_tenant_std_source s
  44. ON s.tenant_id=t.tenant_id AND s.std_object=@obj
  45. WHERE t.tenant_id=@tenant
  46. AND t.source_system<>s.source_system
  47. {docFilter}
  48. """;
  49. var n = await _db.Ado.ExecuteCommandAsync(sql, new { tenant = tenantId, obj = def.Code, doc = def.DocType });
  50. if (n > 0)
  51. _logger.LogWarning(
  52. "中立层来源清理 tenant={Tenant} object={Object} table={Table} deleted={Deleted}",
  53. tenantId, def.Code, table, n);
  54. total += n;
  55. }
  56. return total;
  57. }
  58. /// <summary>
  59. /// 登记来源在该表一行都没有时,不执行清理。
  60. /// <para>
  61. /// 清理是「按登记收口」,不是「清空对象」。若登记来源尚未产出任何行,删掉非登记来源就等于
  62. /// 把该对象整片清空:指标随即算出 NO_DATA 且全程不报错,看板一夜之间没数(S6/S7 就是这么掉的)。
  63. /// 这种情况下真正的问题是登记与数据不一致,须人工核对登记,不能靠重算去「修正」。
  64. /// 故此处保留存量并按 <c>PURGE_SKIPPED_EMPTY_SOURCE</c> 告警,宁可留旧数据也不静默清空。
  65. /// </para>
  66. /// </summary>
  67. private async Task<bool> RegisteredSourceHasRowsAsync(
  68. long tenantId, string stdObject, string table, string docFilter, CancellationToken cancellationToken)
  69. {
  70. cancellationToken.ThrowIfCancellationRequested();
  71. var rows = await _db.Ado.SqlQueryAsync<SourceRowCount>(
  72. $"""
  73. SELECT MAX(s.source_system) AS RegisteredSource,
  74. SUM(t.source_system=s.source_system) AS RegisteredRows,
  75. SUM(t.source_system<>s.source_system) AS ForeignRows
  76. FROM `{table}` t
  77. INNER JOIN mdp_tenant_std_source s
  78. ON s.tenant_id=t.tenant_id AND s.std_object=@obj
  79. WHERE t.tenant_id=@tenant
  80. {docFilter}
  81. """,
  82. new { tenant = tenantId, obj = stdObject, doc = MdpStdObjectCatalog.Find(stdObject)?.DocType });
  83. var stat = rows.FirstOrDefault();
  84. if (stat == null || stat.ForeignRows <= 0)
  85. return true;
  86. if (stat.RegisteredRows > 0)
  87. return true;
  88. _logger.LogError(
  89. "中立层来源清理已跳过:登记来源无数据 tenant={Tenant} object={Object} table={Table} 登记来源={Source} 待清理行={Foreign}",
  90. tenantId, stdObject, table, stat.RegisteredSource, stat.ForeignRows);
  91. await _db.Ado.ExecuteCommandAsync(
  92. """
  93. INSERT INTO mdp_source_gate_log
  94. (tenant_id, std_object, source_system, gate_reason, row_count, sample_keys, sync_batch_id)
  95. VALUES (@tenant, @obj, @src, 'PURGE_SKIPPED_EMPTY_SOURCE', @rows, @sample, '')
  96. """,
  97. new
  98. {
  99. tenant = tenantId,
  100. obj = stdObject,
  101. src = stat.RegisteredSource ?? "",
  102. rows = stat.ForeignRows,
  103. sample = table
  104. });
  105. return false;
  106. }
  107. private sealed class SourceRowCount
  108. {
  109. public string? RegisteredSource { get; set; }
  110. public long RegisteredRows { get; set; }
  111. public long ForeignRows { get; set; }
  112. }
  113. private async Task<bool> TableReadyAsync(string table)
  114. {
  115. var n = await _db.Ado.SqlQueryAsync<int>(
  116. """
  117. SELECT COUNT(*) FROM information_schema.COLUMNS
  118. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@table AND COLUMN_NAME='source_system'
  119. """,
  120. new { table });
  121. return n.FirstOrDefault() > 0;
  122. }
  123. private async Task<bool> ColumnExistsAsync(string table, string column)
  124. {
  125. var n = await _db.Ado.SqlQueryAsync<int>(
  126. """
  127. SELECT COUNT(*) FROM information_schema.COLUMNS
  128. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@table AND COLUMN_NAME=@column
  129. """,
  130. new { table, column });
  131. return n.FirstOrDefault() > 0;
  132. }
  133. }