using Admin.NET.Plugin.AiDOP.DataPlatform.Executors; using SqlSugar; namespace Admin.NET.Plugin.AiDOP.DataPlatform; /// /// T8 基表双模式入站(S5/S6/S7 KPI 共用贴源+标准层)。 /// 照抄 S6_REPORT 范式:执行器 → mdp_stg_t8_* → transform → mdp_std_t8_*。 /// 源:mdp_source=T8_V5_SQLSERVER(MdpSourceScopeFactory 复用 Database.json 的 t8_v5 只读连接)。 /// 实体:mdp_entity T8_KC_TZ_HEAD / T8_KC_TZ_LIST / T8_KC_DD_HEAD / T8_KC_DD_LIST / T8_KC_ZJ_LIST / T8_SYS_PELIST /// (种子把 target_table_name 指向对应 mdp_stg_t8_*、sync_mode=FULL、biz_key_expr=Id)。 /// 只读源、只写 mdp;不写 T8;FULL 全量(源表极小、re-upsert 刷新可带回字段/作废标志更新)。 /// 下游 S5/S6/S7 KPI 已改读 mdp_std_t8_*;故 KPI 跑批前须先跑一次本入站。 /// 说明:S5_L1_002 已中台化——其依赖的 Cj_Bg_Head_Rep(报工头,取min kgdate) 与 kc_dd_list_cllist(订单物料明细) /// 随本服务贴源到 mdp_std_t8_cj_bg_head_rep / mdp_std_t8_kc_dd_list_cllist(biz_key=Id,实测唯一)。 /// 仅 S5_L1_004(Rep_总账_存货_V3 TVF,报表聚合结果无逐行主键)仍走 legacy 直连 T8,其实体 T8_REP_ZONGZHANG_CUNHUO_V3 不贴源。 /// public sealed class T8BaseInboundMdpSyncService : ITransient { private readonly ISqlSugarClient _db; private readonly MdpSourcePullDispatcher _pullDispatcher; // 入站实体码(对应 6 张 T8 基表)→ 目标 std 表;stg 表名由 mdp_entity.target_table_name 决定。 private static readonly (string EntityCode, string StgTable, string StdTable)[] Entities = { ("T8_KC_TZ_HEAD", "mdp_stg_t8_kc_tz_head", "mdp_std_t8_kc_tz_head"), ("T8_KC_TZ_LIST", "mdp_stg_t8_kc_tz_list", "mdp_std_t8_kc_tz_list"), ("T8_KC_DD_HEAD", "mdp_stg_t8_kc_dd_head", "mdp_std_t8_kc_dd_head"), ("T8_KC_DD_LIST", "mdp_stg_t8_kc_dd_list", "mdp_std_t8_kc_dd_list"), ("T8_KC_ZJ_LIST", "mdp_stg_t8_kc_zj_list", "mdp_std_t8_kc_zj_list"), ("T8_SYS_PELIST", "mdp_stg_t8_sys_pelist", "mdp_std_t8_sys_pelist"), // S5_L1_002 中台化新增:报工头(取 min kgdate 开工时间) 与 订单物料明细(分母行数)。 ("T8_CJ_BG_HEAD_REP", "mdp_stg_t8_cj_bg_head_rep", "mdp_std_t8_cj_bg_head_rep"), ("T8_KC_DD_LIST_CLLIST", "mdp_stg_t8_kc_dd_list_cllist", "mdp_std_t8_kc_dd_list_cllist"), }; public T8BaseInboundMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher) { _db = db; _pullDispatcher = pullDispatcher; } public async Task RunInboundAsync( long tenantId = 0, bool fullRefresh = true, string? entityCode = null, CancellationToken cancellationToken = default) { cancellationToken.ThrowIfCancellationRequested(); await EnsureTablesAsync(); var now = DateTime.Now; var batchId = $"T8_BASE_IN_{now:yyyyMMddHHmmss}"; var result = new T8BaseInboundResult { BatchId = batchId }; var targets = string.IsNullOrWhiteSpace(entityCode) ? Entities : Entities.Where(e => string.Equals(e.EntityCode, entityCode.Trim(), StringComparison.OrdinalIgnoreCase)).ToArray(); foreach (var (code, stgTable, stdTable) in targets) { cancellationToken.ThrowIfCancellationRequested(); var pullCtx = new MdpPullContext { TenantId = tenantId, FullRefresh = fullRefresh, TaskCode = "T8_BASE_MDP_INBOUND", BatchId = $"{batchId}_{code}" }; var pull = await _pullDispatcher.PullByEntityCodeAsync(code, pullCtx, cancellationToken); var transformBatch = $"{pullCtx.BatchId}_STD"; var stdRows = await TransformStandardAsync(code, stgTable, stdTable, tenantId, transformBatch, now); result.Tables.Add(new T8BaseInboundTableResult { EntityCode = code, StgRows = pull.RowsWritten, StdRows = stdRows }); } return result; } /// 把 PENDING 贴源行按各表 typed 契约投影到 std(JSON key 区分大小写:源标识列为 Id/id,其余小写)。 private async Task TransformStandardAsync( string entityCode, string stgTable, string stdTable, long tenantId, string batchId, DateTime now) { var sql = entityCode switch { "T8_KC_TZ_HEAD" => BuildTzHeadTransform(), "T8_KC_TZ_LIST" => BuildTzListTransform(), "T8_KC_DD_HEAD" => BuildDdHeadTransform(), "T8_KC_DD_LIST" => BuildDdListTransform(), "T8_KC_ZJ_LIST" => BuildZjListTransform(), "T8_SYS_PELIST" => BuildPelistTransform(), "T8_CJ_BG_HEAD_REP" => BuildCjBgHeadRepTransform(), "T8_KC_DD_LIST_CLLIST" => BuildDdListCllistTransform(), _ => throw new InvalidOperationException($"未支持的 T8 入站实体:{entityCode}") }; var affected = await _db.Ado.ExecuteCommandAsync(sql, new SugarParameter("@tid", tenantId), new SugarParameter("@batch", batchId), new SugarParameter("@now", now)); await _db.Ado.ExecuteCommandAsync( $"UPDATE {stgTable} SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING'"); return affected; } // JSON 取值片段:数字/整数经 NULLIF 去 'null';datetime 把 ISO 'T' 换空格再 STR_TO_DATE(尾部小数位被忽略)。 private static string JStr(string k) => $"NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.{k}')),'null')"; private static string JNum(string k) => $"CAST({JStr(k)} AS DECIMAL(18,6))"; private static string JInt(string k) => $"CAST({JStr(k)} AS SIGNED)"; // ISO 时间戳(System.Text.Json 输出如 2026-05-08T13:45:16.92):'T'→空格、去亚秒再 STR_TO_DATE, // 否则 MySQL 严格模式对残留 '.92' 报 Truncated incorrect datetime value。KPI 只用日级,丢亚秒无影响。 private static string JDate(string k) => $"STR_TO_DATE(SUBSTRING_INDEX(REPLACE({JStr(k)},'T',' '),'.',1),'%Y-%m-%d %H:%i:%s')"; private static string SrcId => $"CAST(COALESCE({JStr("Id")},{JStr("id")}) AS SIGNED)"; private static string StgHead(string table) => $@" FROM {table} s WHERE s.process_status='PENDING' AND s.source_biz_key IS NOT NULL AND s.source_biz_key<>''"; private static string BuildTzHeadTransform() => $@" INSERT INTO mdp_std_t8_kc_tz_head (tenant_id, source_system, src_id, ztid, lbs, lynoid, hzyn, zfyn, shyn, shtime, date0, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'), {SrcId}, {JStr("ztid")}, {JStr("lbs")}, {JStr("lynoid")}, {JInt("hzyn")}, {JInt("zfyn")}, {JInt("shyn")}, {JDate("shtime")}, {JDate("date0")}, IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now {StgHead("mdp_stg_t8_kc_tz_head")} ON DUPLICATE KEY UPDATE src_id=VALUES(src_id), ztid=VALUES(ztid), lbs=VALUES(lbs), lynoid=VALUES(lynoid), hzyn=VALUES(hzyn), zfyn=VALUES(zfyn), shyn=VALUES(shyn), shtime=VALUES(shtime), date0=VALUES(date0), source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP"; private static string BuildTzListTransform() => $@" INSERT INTO mdp_std_t8_kc_tz_list (tenant_id, source_system, src_id, idid, code, lynoid, slzx, sl, gdyn, gdtime, rwnoid, jhdate, addtime, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'), {SrcId}, {JInt("idid")}, {JStr("code")}, {JStr("lynoid")}, {JNum("slzx")}, {JNum("sl")}, {JInt("gdyn")}, {JDate("gdtime")}, {JStr("rwnoid")}, {JDate("jhdate")}, {JDate("addtime")}, IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now {StgHead("mdp_stg_t8_kc_tz_list")} ON DUPLICATE KEY UPDATE src_id=VALUES(src_id), idid=VALUES(idid), code=VALUES(code), lynoid=VALUES(lynoid), slzx=VALUES(slzx), sl=VALUES(sl), gdyn=VALUES(gdyn), gdtime=VALUES(gdtime), rwnoid=VALUES(rwnoid), jhdate=VALUES(jhdate), addtime=VALUES(addtime), source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP"; private static string BuildDdHeadTransform() => $@" INSERT INTO mdp_std_t8_kc_dd_head (tenant_id, source_system, src_id, ztid, lbs, zf, shyn, noid, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'), {SrcId}, {JStr("ztid")}, {JStr("lbs")}, {JInt("zf")}, {JInt("shyn")}, {JStr("noid")}, IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now {StgHead("mdp_stg_t8_kc_dd_head")} ON DUPLICATE KEY UPDATE src_id=VALUES(src_id), ztid=VALUES(ztid), lbs=VALUES(lbs), zf=VALUES(zf), shyn=VALUES(shyn), noid=VALUES(noid), source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP"; private static string BuildDdListTransform() => $@" INSERT INTO mdp_std_t8_kc_dd_list (tenant_id, source_system, src_id, idid, rwnoid, code, sl, slzx, jhdate, gdyn, gdtime, addtime, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'), {SrcId}, {JInt("idid")}, {JStr("rwnoid")}, {JStr("code")}, {JNum("sl")}, {JNum("slzx")}, {JDate("jhdate")}, {JInt("gdyn")}, {JDate("gdtime")}, {JDate("addtime")}, IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now {StgHead("mdp_stg_t8_kc_dd_list")} ON DUPLICATE KEY UPDATE src_id=VALUES(src_id), idid=VALUES(idid), rwnoid=VALUES(rwnoid), code=VALUES(code), sl=VALUES(sl), slzx=VALUES(slzx), jhdate=VALUES(jhdate), gdyn=VALUES(gdyn), gdtime=VALUES(gdtime), addtime=VALUES(addtime), source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP"; private static string BuildZjListTransform() => $@" INSERT INTO mdp_std_t8_kc_zj_list (tenant_id, source_system, src_id, ztid, lyid, zjyn, shdate, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'), {SrcId}, {JStr("ztid")}, {JInt("lyid")}, {JInt("zjyn")}, {JDate("shdate")}, IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now {StgHead("mdp_stg_t8_kc_zj_list")} ON DUPLICATE KEY UPDATE src_id=VALUES(src_id), ztid=VALUES(ztid), lyid=VALUES(lyid), zjyn=VALUES(zjyn), shdate=VALUES(shdate), source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP"; private static string BuildPelistTransform() => $@" INSERT INTO mdp_std_t8_sys_pelist (tenant_id, source_system, src_id, ztid, zzzt, gw, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'), {SrcId}, {JStr("ztid")}, {JStr("zzzt")}, {JStr("gw")}, IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now {StgHead("mdp_stg_t8_sys_pelist")} ON DUPLICATE KEY UPDATE src_id=VALUES(src_id), ztid=VALUES(ztid), zzzt=VALUES(zzzt), gw=VALUES(gw), source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP"; // S5_L1_002 分子:报工头,按 noid 取 min(kgdate)=开工时间(下游 KPI 侧 GROUP BY noid 重算 MIN,保留逐行 kgdate)。 private static string BuildCjBgHeadRepTransform() => $@" INSERT INTO mdp_std_t8_cj_bg_head_rep (tenant_id, source_system, src_id, noid, kgdate, ztid, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'), {SrcId}, {JStr("noid")}, {JDate("kgdate")}, {JStr("ztid")}, IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now {StgHead("mdp_stg_t8_cj_bg_head_rep")} ON DUPLICATE KEY UPDATE src_id=VALUES(src_id), noid=VALUES(noid), kgdate=VALUES(kgdate), ztid=VALUES(ztid), source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP"; // S5_L1_002 分母:订单物料明细,count 行数 per 工单(关联 kc_dd_head.src_id = cllist.idid)。 private static string BuildDdListCllistTransform() => $@" INSERT INTO mdp_std_t8_kc_dd_list_cllist (tenant_id, source_system, src_id, idid, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(s.tenant_id,@tid), IFNULL(NULLIF(s.source_system,''),'T8'), {SrcId}, {JInt("idid")}, IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now {StgHead("mdp_stg_t8_kc_dd_list_cllist")} ON DUPLICATE KEY UPDATE src_id=VALUES(src_id), idid=VALUES(idid), source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP"; /// 建 8 stg(通用信封契约)+ 8 std(KPI 实际引用的 typed 列)。幂等 IF NOT EXISTS。 private async Task EnsureTablesAsync() { foreach (var (_, stgTable, _) in Entities) await _db.Ado.ExecuteCommandAsync(StgDdl(stgTable)); await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_tz_head", "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, lbs varchar(100) DEFAULT NULL, " + "lynoid varchar(400) DEFAULT NULL, " + "hzyn int DEFAULT NULL, zfyn int DEFAULT NULL, shyn int DEFAULT NULL, " + "shtime datetime DEFAULT NULL, date0 date DEFAULT NULL")); await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_tz_list", "src_id bigint DEFAULT NULL, idid bigint DEFAULT NULL, code varchar(400) DEFAULT NULL, " + "lynoid varchar(400) DEFAULT NULL, slzx decimal(18,6) DEFAULT NULL, sl decimal(18,6) DEFAULT NULL, " + "gdyn bigint DEFAULT NULL, gdtime datetime DEFAULT NULL, rwnoid varchar(400) DEFAULT NULL, " + "jhdate date DEFAULT NULL, addtime datetime DEFAULT NULL")); await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_dd_head", "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, lbs varchar(100) DEFAULT NULL, " + "zf int DEFAULT NULL, shyn int DEFAULT NULL, noid varchar(400) DEFAULT NULL")); await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_dd_list", "src_id bigint DEFAULT NULL, idid bigint DEFAULT NULL, rwnoid varchar(400) DEFAULT NULL, " + "code varchar(400) DEFAULT NULL, sl decimal(18,6) DEFAULT NULL, slzx decimal(18,6) DEFAULT NULL, " + "jhdate date DEFAULT NULL, gdyn int DEFAULT NULL, gdtime datetime DEFAULT NULL, addtime datetime DEFAULT NULL")); await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_zj_list", "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, lyid bigint DEFAULT NULL, " + "zjyn int DEFAULT NULL, shdate datetime DEFAULT NULL")); await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_sys_pelist", "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, zzzt varchar(100) DEFAULT NULL, " + "gw varchar(400) DEFAULT NULL")); // S5_L1_002 中台化新增 2 张 std(stg 由上方 Entities 循环建)。 await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_cj_bg_head_rep", "src_id bigint DEFAULT NULL, noid varchar(400) DEFAULT NULL, kgdate date DEFAULT NULL, " + "ztid varchar(100) DEFAULT NULL")); await _db.Ado.ExecuteCommandAsync(StdDdl("mdp_std_t8_kc_dd_list_cllist", "src_id bigint DEFAULT NULL, idid bigint DEFAULT NULL")); // 既有 mdp_std_t8_kc_tz_head(CREATE IF NOT EXISTS 不会补列)幂等补 lynoid(S5_L1_002 分子按工单号分组用)。 await EnsureColumnAsync("mdp_std_t8_kc_tz_head", "lynoid", "varchar(400) DEFAULT NULL"); } /// 幂等补列:目标表缺该列时 ADD COLUMN;已有则跳过。用于向既有 std 表安全追加字段。 private async Task EnsureColumnAsync(string table, string column, string columnDef) { var exists = await _db.Ado.GetIntAsync( "SELECT COUNT(*) FROM information_schema.COLUMNS " + "WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME=@c", new SugarParameter("@t", table), new SugarParameter("@c", column)); if (exists == 0) await _db.Ado.ExecuteCommandAsync($"ALTER TABLE {table} ADD COLUMN {column} {columnDef}"); } private static string StgDdl(string table) => $@" CREATE TABLE IF NOT EXISTS {table} ( id bigint NOT NULL AUTO_INCREMENT, tenant_id bigint NOT NULL DEFAULT 0, source_system varchar(50) DEFAULT NULL, source_table varchar(200) DEFAULT NULL, source_row_id varchar(200) DEFAULT NULL, source_biz_key varchar(300) DEFAULT NULL, raw_data json DEFAULT NULL, sync_batch_id varchar(100) DEFAULT NULL, create_time datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, sync_time datetime DEFAULT CURRENT_TIMESTAMP, process_status varchar(20) NOT NULL DEFAULT 'PENDING', process_message varchar(500) DEFAULT NULL, update_time datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (id), UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key), KEY idx_batch (sync_batch_id) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='T8 基表贴源层(双模式入站)'"; private static string StdDdl(string table, string bizCols) => $@" CREATE TABLE IF NOT EXISTS {table} ( id bigint NOT NULL AUTO_INCREMENT, tenant_id bigint NOT NULL DEFAULT 0, source_system varchar(50) NOT NULL DEFAULT 'T8', {bizCols}, source_row_id varchar(200) NOT NULL, source_biz_key varchar(300) NOT NULL, sync_batch_id varchar(100) NOT NULL, sync_time datetime NOT NULL, update_time datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (id), UNIQUE KEY uk_{table} (tenant_id, source_system, source_biz_key), KEY idx_{table}_batch (sync_batch_id) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='T8 基表标准层(双模式入站)'"; } /// T8 基表入站结果。 public sealed class T8BaseInboundResult { public string BatchId { get; set; } = string.Empty; public List Tables { get; } = new(); } public sealed class T8BaseInboundTableResult { public string EntityCode { get; set; } = string.Empty; public int StgRows { get; set; } public int StdRows { get; set; } }