using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.DataPlatform; /// /// 删除某租户某标准对象里来源与登记不一致的中立行。不删贴源,不删目录之外的表。 /// 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 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 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; } /// /// 登记来源在该表一行都没有时,不执行清理。 /// /// 清理是「按登记收口」,不是「清空对象」。若登记来源尚未产出任何行,删掉非登记来源就等于 /// 把该对象整片清空:指标随即算出 NO_DATA 且全程不报错,看板一夜之间没数(S6/S7 就是这么掉的)。 /// 这种情况下真正的问题是登记与数据不一致,须人工核对登记,不能靠重算去「修正」。 /// 故此处保留存量并按 PURGE_SKIPPED_EMPTY_SOURCE 告警,宁可留旧数据也不静默清空。 /// /// private async Task RegisteredSourceHasRowsAsync( long tenantId, string stdObject, string table, string docFilter, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); var rows = await _db.Ado.SqlQueryAsync( $""" 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 TableReadyAsync(string table) { var n = await _db.Ado.SqlQueryAsync( """ 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 ColumnExistsAsync(string table, string column) { var n = await _db.Ado.SqlQueryAsync( """ 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; } }