|
|
@@ -11,8 +11,9 @@ namespace Admin.NET.Plugin.AiDOP.DataPlatform;
|
|
|
/// (种子把 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)/ S5_L1_004(Rep_总账_存货_V3 TVF)仍走 legacy 直连 T8,
|
|
|
-/// 其对应实体 T8_KC_DD_LIST_CLLIST / T8_CJ_BG_HEAD_REP / T8_REP_ZONGZHANG_CUNHUO_V3 本批不贴源。
|
|
|
+/// 说明: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 不贴源。
|
|
|
/// </summary>
|
|
|
public sealed class T8BaseInboundMdpSyncService : ITransient
|
|
|
{
|
|
|
@@ -28,6 +29,9 @@ public sealed class T8BaseInboundMdpSyncService : ITransient
|
|
|
("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)
|
|
|
@@ -89,6 +93,8 @@ public sealed class T8BaseInboundMdpSyncService : ITransient
|
|
|
"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}")
|
|
|
};
|
|
|
|
|
|
@@ -117,15 +123,15 @@ WHERE s.process_status='PENDING' AND s.source_biz_key IS NOT NULL AND s.source_b
|
|
|
|
|
|
private static string BuildTzHeadTransform() => $@"
|
|
|
INSERT INTO mdp_std_t8_kc_tz_head
|
|
|
- (tenant_id, source_system, src_id, ztid, lbs, hzyn, zfyn, shyn, shtime, date0,
|
|
|
+ (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")}, {JInt("hzyn")}, {JInt("zfyn")}, {JInt("shyn")},
|
|
|
+ {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),
|
|
|
+ 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),
|
|
|
@@ -203,7 +209,35 @@ ON DUPLICATE KEY UPDATE
|
|
|
source_row_id=VALUES(source_row_id), sync_batch_id=VALUES(sync_batch_id),
|
|
|
sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP";
|
|
|
|
|
|
- /// <summary>建 6 stg(通用信封契约)+ 6 std(KPI 实际引用的 typed 列)。幂等 IF NOT EXISTS。</summary>
|
|
|
+ // 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";
|
|
|
+
|
|
|
+ /// <summary>建 8 stg(通用信封契约)+ 8 std(KPI 实际引用的 typed 列)。幂等 IF NOT EXISTS。</summary>
|
|
|
private async Task EnsureTablesAsync()
|
|
|
{
|
|
|
foreach (var (_, stgTable, _) in Entities)
|
|
|
@@ -211,6 +245,7 @@ ON DUPLICATE KEY UPDATE
|
|
|
|
|
|
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",
|
|
|
@@ -231,6 +266,25 @@ ON DUPLICATE KEY UPDATE
|
|
|
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");
|
|
|
+ }
|
|
|
+
|
|
|
+ /// <summary>幂等补列:目标表缺该列时 ADD COLUMN;已有则跳过。用于向既有 std 表安全追加字段。</summary>
|
|
|
+ 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) => $@"
|