T8BaseInboundMdpSyncService.cs 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  2. using Admin.NET.Plugin.AiDOP.Infrastructure;
  3. using SqlSugar;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  5. /// <summary>
  6. /// T8 基表双模式入站(S5/S6/S7 KPI 共用贴源+标准层)。
  7. /// 照抄 S6_REPORT 范式:执行器 → mdp_stg_t8_* → transform → mdp_std_t8_*。
  8. /// 源:mdp_source=T8_V5_SQLSERVER(MdpSourceScopeFactory 复用 Database.json 的 t8_v5 只读连接)。
  9. /// 实体: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
  10. /// (种子把 target_table_name 指向对应 mdp_stg_t8_*、sync_mode=FULL、biz_key_expr=Id)。
  11. /// 只读源、只写 mdp;不写 T8;FULL 全量(源表极小、re-upsert 刷新可带回字段/作废标志更新)。
  12. /// 下游 S5/S6/S7 KPI 已改读 mdp_std_t8_*;故 KPI 跑批前须先跑一次本入站。
  13. /// 说明:S5_L1_002 已中台化——其依赖的 Cj_Bg_Head_Rep(报工头,取min kgdate) 与 kc_dd_list_cllist(订单物料明细)
  14. /// 随本服务贴源到 mdp_std_t8_cj_bg_head_rep / mdp_std_t8_kc_dd_list_cllist(biz_key=Id,实测唯一)。
  15. /// 仅 S5_L1_004(Rep_总账_存货_V3 TVF,报表聚合结果无逐行主键)仍走 legacy 直连 T8,其实体 T8_REP_ZONGZHANG_CUNHUO_V3 不贴源。
  16. /// </summary>
  17. public sealed partial class T8BaseInboundMdpSyncService : ITransient
  18. {
  19. /// <summary>T8 当前唯一账套;源表无 tenant_id,须经 <see cref="AidopSourceTenantMap"/> 解析。</summary>
  20. private const string SourceZtid = "pbxfxp";
  21. private readonly ISqlSugarClient _db;
  22. private readonly MdpSourcePullDispatcher _pullDispatcher;
  23. private readonly MdpNeutralSourceGate _neutralGate;
  24. private readonly EmployeePositionMapService _positionMap;
  25. private readonly MdpSourceScopeFactory _sourceScope;
  26. private static readonly string STenant = MdpJsonSql.TenantFromStg("s", "@tid");
  27. // 入站实体码(对应 6 张 T8 基表)→ 目标 std 表;stg 表名由 mdp_entity.target_table_name 决定。
  28. private static readonly (string EntityCode, string StgTable, string StdTable)[] Entities =
  29. {
  30. ("T8_KC_TZ_HEAD", "mdp_stg_t8_kc_tz_head", "mdp_std_t8_kc_tz_head"),
  31. ("T8_KC_TZ_LIST", "mdp_stg_t8_kc_tz_list", "mdp_std_t8_kc_tz_list"),
  32. ("T8_KC_DD_HEAD", "mdp_stg_t8_kc_dd_head", "mdp_std_t8_kc_dd_head"),
  33. ("T8_KC_DD_LIST", "mdp_stg_t8_kc_dd_list", "mdp_std_t8_kc_dd_list"),
  34. ("T8_KC_ZJ_LIST", "mdp_stg_t8_kc_zj_list", "mdp_std_t8_kc_zj_list"),
  35. ("T8_SYS_PELIST", "mdp_stg_t8_sys_pelist", "mdp_std_t8_sys_pelist"),
  36. // S5_L1_002 中台化新增:报工头(取 min kgdate 开工时间) 与 订单物料明细(分母行数)。
  37. ("T8_CJ_BG_HEAD_REP", "mdp_stg_t8_cj_bg_head_rep", "mdp_std_t8_cj_bg_head_rep"),
  38. ("T8_KC_DD_LIST_CLLIST", "mdp_stg_t8_kc_dd_list_cllist", "mdp_std_t8_kc_dd_list_cllist"),
  39. };
  40. public T8BaseInboundMdpSyncService(
  41. ISqlSugarClient db,
  42. MdpSourcePullDispatcher pullDispatcher,
  43. MdpNeutralSourceGate neutralGate,
  44. EmployeePositionMapService positionMap,
  45. MdpSourceScopeFactory sourceScope)
  46. {
  47. _db = db;
  48. _pullDispatcher = pullDispatcher;
  49. _neutralGate = neutralGate;
  50. _positionMap = positionMap;
  51. _sourceScope = sourceScope;
  52. }
  53. public async Task<T8BaseInboundResult> RunInboundAsync(
  54. long tenantId = 0,
  55. bool fullRefresh = true,
  56. string? entityCode = null,
  57. CancellationToken cancellationToken = default,
  58. string? sourceZtid = null)
  59. {
  60. cancellationToken.ThrowIfCancellationRequested();
  61. await EnsureTablesAsync();
  62. // 账套以登记表 source_scope 为准;未传入时才用适配器默认账套。
  63. var ztid = string.IsNullOrWhiteSpace(sourceZtid) ? SourceZtid : sourceZtid.Trim();
  64. var resolvedTenantId = AidopSourceTenantMap.ResolveTenantId(ztid, tenantId);
  65. var now = DateTime.Now;
  66. var batchId = $"T8_BASE_IN_{now:yyyyMMddHHmmss}";
  67. var result = new T8BaseInboundResult { BatchId = batchId };
  68. var targets = string.IsNullOrWhiteSpace(entityCode)
  69. ? Entities
  70. : Entities.Where(e => string.Equals(e.EntityCode, entityCode.Trim(), StringComparison.OrdinalIgnoreCase)).ToArray();
  71. foreach (var (code, stgTable, stdTable) in targets)
  72. {
  73. cancellationToken.ThrowIfCancellationRequested();
  74. var pullCtx = new MdpPullContext
  75. {
  76. TenantId = resolvedTenantId,
  77. FullRefresh = fullRefresh,
  78. TaskCode = "T8_BASE_MDP_INBOUND",
  79. BatchId = $"{batchId}_{code}"
  80. };
  81. var pull = await _pullDispatcher.PullByEntityCodeAsync(code, pullCtx, cancellationToken);
  82. var transformBatch = $"{pullCtx.BatchId}_STD";
  83. var stdRows = await TransformStandardAsync(code, stgTable, stdTable, resolvedTenantId, transformBatch, now);
  84. result.Tables.Add(new T8BaseInboundTableResult
  85. {
  86. EntityCode = code,
  87. StgRows = pull.RowsWritten,
  88. StdRows = stdRows
  89. });
  90. }
  91. result.NeutralProjection = await ProjectNeutralAsync(resolvedTenantId, batchId, now, cancellationToken);
  92. return result;
  93. }
  94. /// <summary>把 PENDING 贴源行按各表 typed 契约投影到 std(JSON key 区分大小写:源标识列为 Id/id,其余小写)。</summary>
  95. private async Task<int> TransformStandardAsync(
  96. string entityCode, string stgTable, string stdTable, long tenantId, string batchId, DateTime now)
  97. {
  98. var sql = entityCode switch
  99. {
  100. "T8_KC_TZ_HEAD" => BuildTzHeadTransform(),
  101. "T8_KC_TZ_LIST" => BuildTzListTransform(),
  102. "T8_KC_DD_HEAD" => BuildDdHeadTransform(),
  103. "T8_KC_DD_LIST" => BuildDdListTransform(),
  104. "T8_KC_ZJ_LIST" => BuildZjListTransform(),
  105. "T8_SYS_PELIST" => BuildPelistTransform(),
  106. "T8_CJ_BG_HEAD_REP" => BuildCjBgHeadRepTransform(),
  107. "T8_KC_DD_LIST_CLLIST" => BuildDdListCllistTransform(),
  108. _ => throw new InvalidOperationException($"未支持的 T8 入站实体:{entityCode}")
  109. };
  110. var affected = await MdpSchemaAligner.ExecuteAsync(_db, sql,
  111. new SugarParameter("@tid", tenantId),
  112. new SugarParameter("@batch", batchId),
  113. new SugarParameter("@now", now));
  114. await MdpSchemaAligner.ExecuteAsync(_db,
  115. tenantId > 0
  116. ? $"UPDATE {stgTable} SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING' AND tenant_id = @tid"
  117. : $"UPDATE {stgTable} SET process_status='DONE', update_time=NOW() WHERE process_status='PENDING'",
  118. tenantId > 0 ? new SugarParameter("@tid", tenantId) : null!);
  119. return affected;
  120. }
  121. // JSON 取值片段:数字/整数经 NULLIF 去 'null';datetime 把 ISO 'T' 换空格再 STR_TO_DATE(尾部小数位被忽略)。
  122. private static string JStr(string k) => $"NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.{k}')),'null')";
  123. private static string JNum(string k) => $"CAST({JStr(k)} AS DECIMAL(18,6))";
  124. private static string JInt(string k) => $"CAST({JStr(k)} AS SIGNED)";
  125. // ISO 时间戳(System.Text.Json 输出如 2026-05-08T13:45:16.92):'T'→空格、去亚秒再 STR_TO_DATE,
  126. // 否则 MySQL 严格模式对残留 '.92' 报 Truncated incorrect datetime value。KPI 只用日级,丢亚秒无影响。
  127. private static string JDate(string k) => $"STR_TO_DATE(SUBSTRING_INDEX(REPLACE({JStr(k)},'T',' '),'.',1),'%Y-%m-%d %H:%i:%s')";
  128. private static string SrcId => $"CAST(COALESCE({JStr("Id")},{JStr("id")}) AS SIGNED)";
  129. private static string StgHead(string table) => $@"
  130. FROM {table} s
  131. WHERE s.process_status='PENDING' AND s.source_biz_key IS NOT NULL AND s.source_biz_key<>'' AND {MdpJsonSql.TenantGuard(STenant)}";
  132. private static string BuildTzHeadTransform() => $@"
  133. INSERT INTO mdp_std_t8_kc_tz_head
  134. (tenant_id, source_system, src_id, ztid, lbs, lynoid, hzyn, zfyn, shyn, shtime, date0, zbrname,
  135. source_row_id, source_biz_key, sync_batch_id, sync_time)
  136. SELECT {STenant}, IFNULL(NULLIF(s.source_system,''),'T8'),
  137. {SrcId}, {JStr("ztid")}, {JStr("lbs")}, {JStr("lynoid")}, {JInt("hzyn")}, {JInt("zfyn")}, {JInt("shyn")},
  138. {JDate("shtime")}, {JDate("date0")}, {JStr("zbrname")},
  139. IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
  140. {StgHead("mdp_stg_t8_kc_tz_head")}
  141. ON DUPLICATE KEY UPDATE
  142. src_id=VALUES(src_id), ztid=VALUES(ztid), lbs=VALUES(lbs), lynoid=VALUES(lynoid),
  143. hzyn=VALUES(hzyn), zfyn=VALUES(zfyn), shyn=VALUES(shyn),
  144. shtime=VALUES(shtime), date0=VALUES(date0), zbrname=VALUES(zbrname),
  145. source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
  146. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
  147. private static string BuildTzListTransform() => $@"
  148. INSERT INTO mdp_std_t8_kc_tz_list
  149. (tenant_id, source_system, src_id, idid, code, lynoid, slzx, sl, gdyn, gdtime, rwnoid, jhdate, addtime,
  150. hw, ckcode, pcnoid, kcsl,
  151. source_row_id, source_biz_key, sync_batch_id, sync_time)
  152. SELECT {STenant}, IFNULL(NULLIF(s.source_system,''),'T8'),
  153. {SrcId}, {JInt("idid")}, {JStr("code")}, {JStr("lynoid")}, {JNum("slzx")}, {JNum("sl")},
  154. {JInt("gdyn")}, {JDate("gdtime")}, {JStr("rwnoid")}, {JDate("jhdate")}, {JDate("addtime")},
  155. {JStr("hw")}, {JStr("ckcode")}, {JStr("pcnoid")}, {JNum("KcSl")},
  156. IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
  157. {StgHead("mdp_stg_t8_kc_tz_list")}
  158. ON DUPLICATE KEY UPDATE
  159. src_id=VALUES(src_id), idid=VALUES(idid), code=VALUES(code), lynoid=VALUES(lynoid),
  160. slzx=VALUES(slzx), sl=VALUES(sl), gdyn=VALUES(gdyn), gdtime=VALUES(gdtime),
  161. rwnoid=VALUES(rwnoid), jhdate=VALUES(jhdate), addtime=VALUES(addtime),
  162. hw=VALUES(hw), ckcode=VALUES(ckcode), pcnoid=VALUES(pcnoid), kcsl=VALUES(kcsl),
  163. source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
  164. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
  165. private static string BuildDdHeadTransform() => $@"
  166. INSERT INTO mdp_std_t8_kc_dd_head
  167. (tenant_id, source_system, src_id, ztid, lbs, zf, shyn, noid,
  168. source_row_id, source_biz_key, sync_batch_id, sync_time)
  169. SELECT {STenant}, IFNULL(NULLIF(s.source_system,''),'T8'),
  170. {SrcId}, {JStr("ztid")}, {JStr("lbs")}, {JInt("zf")}, {JInt("shyn")}, {JStr("noid")},
  171. IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
  172. {StgHead("mdp_stg_t8_kc_dd_head")}
  173. ON DUPLICATE KEY UPDATE
  174. src_id=VALUES(src_id), ztid=VALUES(ztid), lbs=VALUES(lbs),
  175. zf=VALUES(zf), shyn=VALUES(shyn), noid=VALUES(noid),
  176. source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
  177. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
  178. private static string BuildDdListTransform() => $@"
  179. INSERT INTO mdp_std_t8_kc_dd_list
  180. (tenant_id, source_system, src_id, idid, rwnoid, code, sl, slzx, jhdate, gdyn, gdtime, addtime,
  181. source_row_id, source_biz_key, sync_batch_id, sync_time)
  182. SELECT {STenant}, IFNULL(NULLIF(s.source_system,''),'T8'),
  183. {SrcId}, {JInt("idid")}, {JStr("rwnoid")}, {JStr("code")}, {JNum("sl")}, {JNum("slzx")},
  184. {JDate("jhdate")}, {JInt("gdyn")}, {JDate("gdtime")}, {JDate("addtime")},
  185. IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
  186. {StgHead("mdp_stg_t8_kc_dd_list")}
  187. ON DUPLICATE KEY UPDATE
  188. src_id=VALUES(src_id), idid=VALUES(idid), rwnoid=VALUES(rwnoid), code=VALUES(code),
  189. sl=VALUES(sl), slzx=VALUES(slzx), jhdate=VALUES(jhdate), gdyn=VALUES(gdyn), gdtime=VALUES(gdtime),
  190. addtime=VALUES(addtime),
  191. source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
  192. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
  193. private static string BuildZjListTransform() => $@"
  194. INSERT INTO mdp_std_t8_kc_zj_list
  195. (tenant_id, source_system, src_id, ztid, lyid, zjyn, shdate,
  196. source_row_id, source_biz_key, sync_batch_id, sync_time)
  197. SELECT {STenant}, IFNULL(NULLIF(s.source_system,''),'T8'),
  198. {SrcId}, {JStr("ztid")}, {JInt("lyid")}, {JInt("zjyn")}, {JDate("shdate")},
  199. IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
  200. {StgHead("mdp_stg_t8_kc_zj_list")}
  201. ON DUPLICATE KEY UPDATE
  202. src_id=VALUES(src_id), ztid=VALUES(ztid), lyid=VALUES(lyid), zjyn=VALUES(zjyn), shdate=VALUES(shdate),
  203. source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
  204. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
  205. private static string BuildPelistTransform() => $@"
  206. INSERT INTO mdp_std_t8_sys_pelist
  207. (tenant_id, source_system, src_id, ztid, zzzt, gw,
  208. source_row_id, source_biz_key, sync_batch_id, sync_time)
  209. SELECT {STenant}, IFNULL(NULLIF(s.source_system,''),'T8'),
  210. {SrcId}, {JStr("ztid")}, {JStr("zzzt")}, {JStr("gw")},
  211. IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
  212. {StgHead("mdp_stg_t8_sys_pelist")}
  213. ON DUPLICATE KEY UPDATE
  214. src_id=VALUES(src_id), ztid=VALUES(ztid), zzzt=VALUES(zzzt), gw=VALUES(gw),
  215. source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
  216. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
  217. // S5_L1_002 分子:报工头,按 noid 取 min(kgdate)=开工时间(下游 KPI 侧 GROUP BY noid 重算 MIN,保留逐行 kgdate)。
  218. private static string BuildCjBgHeadRepTransform() => $@"
  219. INSERT INTO mdp_std_t8_cj_bg_head_rep
  220. (tenant_id, source_system, src_id, noid, kgdate, ztid,
  221. source_row_id, source_biz_key, sync_batch_id, sync_time)
  222. SELECT {STenant}, IFNULL(NULLIF(s.source_system,''),'T8'),
  223. {SrcId}, {JStr("noid")}, {JDate("kgdate")}, {JStr("ztid")},
  224. IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
  225. {StgHead("mdp_stg_t8_cj_bg_head_rep")}
  226. ON DUPLICATE KEY UPDATE
  227. src_id=VALUES(src_id), noid=VALUES(noid), kgdate=VALUES(kgdate), ztid=VALUES(ztid),
  228. source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
  229. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
  230. // S5_L1_002 分母:订单物料明细,count 行数 per 工单(关联 kc_dd_head.src_id = cllist.idid)。
  231. private static string BuildDdListCllistTransform() => $@"
  232. INSERT INTO mdp_std_t8_kc_dd_list_cllist
  233. (tenant_id, source_system, src_id, idid, code, sl,
  234. source_row_id, source_biz_key, sync_batch_id, sync_time)
  235. SELECT {STenant}, IFNULL(NULLIF(s.source_system,''),'T8'),
  236. {SrcId}, {JInt("idid")}, {JStr("code")}, {JNum("sl")},
  237. IFNULL(NULLIF(s.source_row_id,''),s.source_biz_key), s.source_biz_key, @batch, @now
  238. {StgHead("mdp_stg_t8_kc_dd_list_cllist")}
  239. ON DUPLICATE KEY UPDATE
  240. src_id=VALUES(src_id), idid=VALUES(idid), code=VALUES(code), sl=VALUES(sl),
  241. source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
  242. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
  243. /// <summary>建 8 stg(通用信封契约)+ 8 std(KPI 实际引用的 typed 列)。幂等 IF NOT EXISTS。</summary>
  244. private async Task EnsureTablesAsync()
  245. {
  246. foreach (var (_, stgTable, _) in Entities)
  247. await MdpSchemaAligner.ExecuteAsync(_db, StgDdl(stgTable));
  248. await MdpSchemaAligner.ExecuteAsync(_db, StdDdl("mdp_std_t8_kc_tz_head",
  249. "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, lbs varchar(100) DEFAULT NULL, " +
  250. "lynoid varchar(400) DEFAULT NULL, " +
  251. "hzyn int DEFAULT NULL, zfyn int DEFAULT NULL, shyn int DEFAULT NULL, " +
  252. "shtime datetime DEFAULT NULL, date0 date DEFAULT NULL, zbrname varchar(100) DEFAULT NULL"));
  253. await MdpSchemaAligner.ExecuteAsync(_db, StdDdl("mdp_std_t8_kc_tz_list",
  254. "src_id bigint DEFAULT NULL, idid bigint DEFAULT NULL, code varchar(400) DEFAULT NULL, " +
  255. "lynoid varchar(400) DEFAULT NULL, slzx decimal(18,6) DEFAULT NULL, sl decimal(18,6) DEFAULT NULL, " +
  256. "gdyn bigint DEFAULT NULL, gdtime datetime DEFAULT NULL, rwnoid varchar(400) DEFAULT NULL, " +
  257. "jhdate date DEFAULT NULL, addtime datetime DEFAULT NULL, " +
  258. "hw varchar(80) DEFAULT NULL, ckcode varchar(80) DEFAULT NULL, pcnoid varchar(80) DEFAULT NULL, kcsl decimal(18,6) DEFAULT NULL"));
  259. await MdpSchemaAligner.ExecuteAsync(_db, StdDdl("mdp_std_t8_kc_dd_head",
  260. "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, lbs varchar(100) DEFAULT NULL, " +
  261. "zf int DEFAULT NULL, shyn int DEFAULT NULL, noid varchar(400) DEFAULT NULL"));
  262. await MdpSchemaAligner.ExecuteAsync(_db, StdDdl("mdp_std_t8_kc_dd_list",
  263. "src_id bigint DEFAULT NULL, idid bigint DEFAULT NULL, rwnoid varchar(400) DEFAULT NULL, " +
  264. "code varchar(400) DEFAULT NULL, sl decimal(18,6) DEFAULT NULL, slzx decimal(18,6) DEFAULT NULL, " +
  265. "jhdate date DEFAULT NULL, gdyn int DEFAULT NULL, gdtime datetime DEFAULT NULL, addtime datetime DEFAULT NULL"));
  266. await MdpSchemaAligner.ExecuteAsync(_db, StdDdl("mdp_std_t8_kc_zj_list",
  267. "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, lyid bigint DEFAULT NULL, " +
  268. "zjyn int DEFAULT NULL, shdate datetime DEFAULT NULL"));
  269. await MdpSchemaAligner.ExecuteAsync(_db, StdDdl("mdp_std_t8_sys_pelist",
  270. "src_id bigint DEFAULT NULL, ztid varchar(100) DEFAULT NULL, zzzt varchar(100) DEFAULT NULL, " +
  271. "gw varchar(400) DEFAULT NULL"));
  272. // S5_L1_002 中台化新增 2 张 std(stg 由上方 Entities 循环建)。
  273. await MdpSchemaAligner.ExecuteAsync(_db, StdDdl("mdp_std_t8_cj_bg_head_rep",
  274. "src_id bigint DEFAULT NULL, noid varchar(400) DEFAULT NULL, kgdate date DEFAULT NULL, " +
  275. "ztid varchar(100) DEFAULT NULL"));
  276. await MdpSchemaAligner.ExecuteAsync(_db, StdDdl("mdp_std_t8_kc_dd_list_cllist",
  277. "src_id bigint DEFAULT NULL, idid bigint DEFAULT NULL, code varchar(400) DEFAULT NULL, sl decimal(18,6) DEFAULT NULL"));
  278. await EnsureColumnAsync("mdp_std_t8_kc_dd_list_cllist", "code", "varchar(400) DEFAULT NULL");
  279. await EnsureColumnAsync("mdp_std_t8_kc_dd_list_cllist", "sl", "decimal(18,6) DEFAULT NULL");
  280. // 既有 mdp_std_t8_kc_tz_head(CREATE IF NOT EXISTS 不会补列)幂等补 lynoid(S5_L1_002 分子按工单号分组用)。
  281. await EnsureColumnAsync("mdp_std_t8_kc_tz_head", "lynoid", "varchar(400) DEFAULT NULL");
  282. await EnsureColumnAsync("mdp_std_t8_kc_tz_head", "zbrname", "varchar(100) DEFAULT NULL");
  283. await EnsureColumnAsync("mdp_std_t8_kc_tz_list", "hw", "varchar(80) DEFAULT NULL");
  284. await EnsureColumnAsync("mdp_std_t8_kc_tz_list", "ckcode", "varchar(80) DEFAULT NULL");
  285. await EnsureColumnAsync("mdp_std_t8_kc_tz_list", "pcnoid", "varchar(80) DEFAULT NULL");
  286. await EnsureColumnAsync("mdp_std_t8_kc_tz_list", "kcsl", "decimal(18,6) DEFAULT NULL");
  287. }
  288. /// <summary>幂等补列:目标表缺该列时 ADD COLUMN;已有则跳过。用于向既有 std 表安全追加字段。</summary>
  289. private async Task EnsureColumnAsync(string table, string column, string columnDef)
  290. {
  291. var exists = await _db.Ado.GetIntAsync(
  292. "SELECT COUNT(*) FROM information_schema.COLUMNS " +
  293. "WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME=@c",
  294. new SugarParameter("@t", table), new SugarParameter("@c", column));
  295. if (exists == 0)
  296. await MdpSchemaAligner.ExecuteAsync(_db, $"ALTER TABLE {table} ADD COLUMN {column} {columnDef}");
  297. }
  298. private static string StgDdl(string table) => $@"
  299. CREATE TABLE IF NOT EXISTS {table} (
  300. id bigint NOT NULL AUTO_INCREMENT,
  301. tenant_id bigint NOT NULL DEFAULT 0,
  302. source_system varchar(50) DEFAULT NULL,
  303. source_table varchar(200) DEFAULT NULL,
  304. source_row_id varchar(200) DEFAULT NULL,
  305. source_biz_key varchar(300) DEFAULT NULL,
  306. raw_data json DEFAULT NULL,
  307. sync_batch_id varchar(100) DEFAULT NULL,
  308. create_time datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
  309. sync_time datetime DEFAULT CURRENT_TIMESTAMP,
  310. process_status varchar(20) NOT NULL DEFAULT 'PENDING',
  311. process_message varchar(500) DEFAULT NULL,
  312. update_time datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  313. PRIMARY KEY (id),
  314. UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key),
  315. KEY idx_batch (sync_batch_id)
  316. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='T8 基表贴源层(双模式入站)'";
  317. private static string StdDdl(string table, string bizCols) => $@"
  318. CREATE TABLE IF NOT EXISTS {table} (
  319. id bigint NOT NULL AUTO_INCREMENT,
  320. tenant_id bigint NOT NULL DEFAULT 0,
  321. source_system varchar(50) NOT NULL DEFAULT 'T8',
  322. {bizCols},
  323. source_row_id varchar(200) NOT NULL,
  324. source_biz_key varchar(300) NOT NULL,
  325. sync_batch_id varchar(100) NOT NULL,
  326. sync_time datetime NOT NULL,
  327. update_time datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  328. PRIMARY KEY (id),
  329. UNIQUE KEY uk_{table} (tenant_id, source_system, source_biz_key),
  330. KEY idx_{table}_batch (sync_batch_id)
  331. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='T8 基表标准层(双模式入站)'";
  332. }
  333. /// <summary>T8 基表入站结果。</summary>
  334. public sealed class T8BaseInboundResult
  335. {
  336. public string BatchId { get; set; } = string.Empty;
  337. public List<T8BaseInboundTableResult> Tables { get; } = new();
  338. /// <summary>中立层投影结果。未就绪或单步失败时写明原因,入站本身仍返回。</summary>
  339. public string NeutralProjection { get; set; } = "";
  340. }
  341. public sealed class T8BaseInboundTableResult
  342. {
  343. public string EntityCode { get; set; } = string.Empty;
  344. public int StgRows { get; set; }
  345. public int StdRows { get; set; }
  346. }