| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145 |
- using Microsoft.Extensions.Logging;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform;
- /// <summary>
- /// 删除某租户某标准对象里来源与登记不一致的中立行。不删贴源,不删目录之外的表。
- /// </summary>
- public sealed class MdpNeutralSourceCleanup : ITransient
- {
- private readonly ISqlSugarClient _db;
- private readonly ILogger _logger;
- public MdpNeutralSourceCleanup(ISqlSugarClient db, ILoggerFactory loggerFactory)
- {
- _db = db;
- _logger = loggerFactory.CreateLogger(nameof(MdpNeutralSourceCleanup));
- }
- public async Task<int> PurgeModuleAsync(long tenantId, string moduleCode, CancellationToken cancellationToken = default)
- {
- var total = 0;
- foreach (var def in MdpStdObjectCatalog.All.Where(d => d.Module == moduleCode))
- {
- cancellationToken.ThrowIfCancellationRequested();
- total += await PurgeObjectAsync(tenantId, def.Code, cancellationToken);
- }
- return total;
- }
- public async Task<int> PurgeObjectAsync(long tenantId, string stdObject, CancellationToken cancellationToken = default)
- {
- var def = MdpStdObjectCatalog.Find(stdObject);
- if (def == null || tenantId <= 0)
- return 0;
- var total = 0;
- foreach (var table in def.Tables)
- {
- cancellationToken.ThrowIfCancellationRequested();
- if (!await TableReadyAsync(table))
- continue;
- var docFilter = string.IsNullOrEmpty(def.DocType) || !await ColumnExistsAsync(table, "doc_type")
- ? ""
- : " AND t.doc_type=@doc";
- if (!await RegisteredSourceHasRowsAsync(tenantId, def.Code, table, docFilter, cancellationToken))
- continue;
- var sql = $"""
- DELETE t FROM `{table}` t
- INNER JOIN mdp_tenant_std_source s
- ON s.tenant_id=t.tenant_id AND s.std_object=@obj
- WHERE t.tenant_id=@tenant
- AND t.source_system<>s.source_system
- {docFilter}
- """;
- var n = await _db.Ado.ExecuteCommandAsync(sql, new { tenant = tenantId, obj = def.Code, doc = def.DocType });
- if (n > 0)
- _logger.LogWarning(
- "中立层来源清理 tenant={Tenant} object={Object} table={Table} deleted={Deleted}",
- tenantId, def.Code, table, n);
- total += n;
- }
- return total;
- }
- /// <summary>
- /// 登记来源在该表一行都没有时,不执行清理。
- /// <para>
- /// 清理是「按登记收口」,不是「清空对象」。若登记来源尚未产出任何行,删掉非登记来源就等于
- /// 把该对象整片清空:指标随即算出 NO_DATA 且全程不报错,看板一夜之间没数(S6/S7 就是这么掉的)。
- /// 这种情况下真正的问题是登记与数据不一致,须人工核对登记,不能靠重算去「修正」。
- /// 故此处保留存量并按 <c>PURGE_SKIPPED_EMPTY_SOURCE</c> 告警,宁可留旧数据也不静默清空。
- /// </para>
- /// </summary>
- private async Task<bool> RegisteredSourceHasRowsAsync(
- long tenantId, string stdObject, string table, string docFilter, CancellationToken cancellationToken)
- {
- cancellationToken.ThrowIfCancellationRequested();
- var rows = await _db.Ado.SqlQueryAsync<SourceRowCount>(
- $"""
- SELECT MAX(s.source_system) AS RegisteredSource,
- SUM(t.source_system=s.source_system) AS RegisteredRows,
- SUM(t.source_system<>s.source_system) AS ForeignRows
- FROM `{table}` t
- INNER JOIN mdp_tenant_std_source s
- ON s.tenant_id=t.tenant_id AND s.std_object=@obj
- WHERE t.tenant_id=@tenant
- {docFilter}
- """,
- new { tenant = tenantId, obj = stdObject, doc = MdpStdObjectCatalog.Find(stdObject)?.DocType });
- var stat = rows.FirstOrDefault();
- if (stat == null || stat.ForeignRows <= 0)
- return true;
- if (stat.RegisteredRows > 0)
- return true;
- _logger.LogError(
- "中立层来源清理已跳过:登记来源无数据 tenant={Tenant} object={Object} table={Table} 登记来源={Source} 待清理行={Foreign}",
- tenantId, stdObject, table, stat.RegisteredSource, stat.ForeignRows);
- await _db.Ado.ExecuteCommandAsync(
- """
- INSERT INTO mdp_source_gate_log
- (tenant_id, std_object, source_system, gate_reason, row_count, sample_keys, sync_batch_id)
- VALUES (@tenant, @obj, @src, 'PURGE_SKIPPED_EMPTY_SOURCE', @rows, @sample, '')
- """,
- new
- {
- tenant = tenantId,
- obj = stdObject,
- src = stat.RegisteredSource ?? "",
- rows = stat.ForeignRows,
- sample = table
- });
- return false;
- }
- private sealed class SourceRowCount
- {
- public string? RegisteredSource { get; set; }
- public long RegisteredRows { get; set; }
- public long ForeignRows { get; set; }
- }
- private async Task<bool> TableReadyAsync(string table)
- {
- var n = await _db.Ado.SqlQueryAsync<int>(
- """
- SELECT COUNT(*) FROM information_schema.COLUMNS
- WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@table AND COLUMN_NAME='source_system'
- """,
- new { table });
- return n.FirstOrDefault() > 0;
- }
- private async Task<bool> ColumnExistsAsync(string table, string column)
- {
- var n = await _db.Ado.SqlQueryAsync<int>(
- """
- SELECT COUNT(*) FROM information_schema.COLUMNS
- WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@table AND COLUMN_NAME=@column
- """,
- new { table, column });
- return n.FirstOrDefault() > 0;
- }
- }
|