S1MdpSyncTransformService.cs 122 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.S1Refresh;
  3. using Admin.NET.Plugin.AiDOP.SmartOps;
  4. namespace Admin.NET.Plugin.AiDOP.Order;
  5. /// <summary>
  6. /// S1 首批 MDP 同步和标准化转换服务。
  7. /// Phase2:源→stg 经统一执行器;std/DWD/KPI 仍读 stg。
  8. /// </summary>
  9. public class S1MdpSyncTransformService : ITransient
  10. {
  11. private const string JobCode = "S1_MDP_SYNC_TRANSFORM";
  12. private readonly ISqlSugarClient _db;
  13. private readonly SmartOpsKpiAtomicBuildService _atomicBuild;
  14. private readonly MdpModuleStagingPuller _stagingPuller;
  15. private readonly IKpiTargetResolver _kpiTargetResolver;
  16. private readonly IS1MdpFullRunLock _fullRunLock;
  17. public S1MdpSyncTransformService(
  18. ISqlSugarClient db,
  19. SmartOpsKpiAtomicBuildService atomicBuild,
  20. MdpModuleStagingPuller stagingPuller,
  21. IKpiTargetResolver kpiTargetResolver,
  22. IS1MdpFullRunLock fullRunLock)
  23. {
  24. _db = db;
  25. _atomicBuild = atomicBuild;
  26. _stagingPuller = stagingPuller;
  27. _kpiTargetResolver = kpiTargetResolver;
  28. _fullRunLock = fullRunLock;
  29. }
  30. public async Task<S1MdpSyncTransformResult> RunFullAsync(
  31. S1MdpRunScope scope,
  32. CancellationToken cancellationToken = default,
  33. string triggerType = "AUTO",
  34. long? rebuildJobId = null,
  35. Func<S1MdpProgressUpdate, Task>? reportProgress = null)
  36. {
  37. scope = S1MdpRunScope.Create(scope.TenantId, scope.FactoryId);
  38. cancellationToken.ThrowIfCancellationRequested();
  39. var holderId = $"{triggerType}:{scope.ScopeKey}:{Environment.MachineName}:{Guid.NewGuid():N}";
  40. var lease = await _fullRunLock.TryAcquireAsync(scope.LockName, holderId, rebuildJobId, cancellationToken);
  41. if (lease == null)
  42. throw new S1MdpAlreadyRunningException();
  43. await using (lease)
  44. {
  45. using var heartbeatCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
  46. var heartbeat = KeepLeaseAliveAsync(lease, heartbeatCts.Token);
  47. var now = DateTime.Now;
  48. var batchId = $"S1_MDP_FULL_{scope.TenantId}_{scope.FactoryId}_{now:yyyyMMddHHmmss}";
  49. var runLogId = await InsertTransformRunLogAsync(scope, batchId, now, triggerType);
  50. var result = new S1MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  51. try
  52. {
  53. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Preparing, 0, 5, "准备运行环境"));
  54. await EnsureS1RuntimeObjectsAsync();
  55. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Staging, 1, 10, "正在拉取源数据"));
  56. result.StageRows = await SyncStagingAsync(scope, batchId, now, cancellationToken);
  57. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Staging, 1, 35, "拉取源数据完成", result.StageRows, S1MdpRebuildStage.Staging));
  58. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Standard, 2, 35, "正在标准化数据"));
  59. result.StandardRows = await TransformStandardAsync(scope, batchId, now, cancellationToken);
  60. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Standard, 2, 50, "标准化数据完成", result.StandardRows, S1MdpRebuildStage.Standard));
  61. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Dwd, 3, 50, "正在生成 DWD 明细"));
  62. result.DwdRows = await BuildDwdAsync(scope, batchId, now, result, cancellationToken);
  63. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Dwd, 3, 65, "生成 DWD 明细完成", result.DwdRows, S1MdpRebuildStage.Dwd));
  64. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Kpi, 4, 65, "正在重算 KPI"));
  65. result.KpiRows = await BuildS1KpiValuesAsync(scope, batchId, now, cancellationToken);
  66. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Kpi, 4, 80, "重算 KPI 完成", result.KpiRows, S1MdpRebuildStage.Kpi));
  67. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Atomic, 5, 82, "正在重建原子数据"));
  68. result.AtomicRows = await _atomicBuild.BuildOrderDeliveryDomainForAllDatesAsync(
  69. scope.TenantId, scope.FactoryId, batchId, cancellationToken);
  70. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Atomic, 5, 98, "重建原子数据完成", result.AtomicRows, S1MdpRebuildStage.Atomic));
  71. await ReportAsync(reportProgress, new S1MdpProgressUpdate(S1MdpRebuildStage.Finalizing, 5, 99, "正在收尾"));
  72. await MarkTransformRunSuccessAsync(runLogId, now, result);
  73. return result;
  74. }
  75. catch (Exception ex)
  76. {
  77. var summary = ex is OperationCanceledException
  78. ? "HTTP 请求已取消或服务中断"
  79. : ex.Message;
  80. await MarkTransformRunFailedAsync(runLogId, now, summary);
  81. throw;
  82. }
  83. finally
  84. {
  85. heartbeatCts.Cancel();
  86. try { await heartbeat; } catch (OperationCanceledException) { }
  87. }
  88. }
  89. }
  90. private static async Task ReportAsync(Func<S1MdpProgressUpdate, Task>? report, S1MdpProgressUpdate update)
  91. {
  92. if (report == null) return;
  93. await report(update);
  94. }
  95. /// <summary>贴源已由 FILE 导入写入后,只跑标准层/DWD/KPI,不再 Pull。</summary>
  96. public async Task<S1MdpSyncTransformResult> RunPostStagingAsync(
  97. S1MdpRunScope scope,
  98. CancellationToken cancellationToken = default,
  99. string triggerType = "FILE_IMPORT")
  100. {
  101. scope = S1MdpRunScope.Create(scope.TenantId, scope.FactoryId);
  102. cancellationToken.ThrowIfCancellationRequested();
  103. var now = DateTime.Now;
  104. var batchId = $"S1_MDP_FILE_{scope.TenantId}_{scope.FactoryId}_{now:yyyyMMddHHmmss}";
  105. var runLogId = await InsertTransformRunLogAsync(scope, batchId, now, triggerType);
  106. var result = new S1MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  107. try
  108. {
  109. await EnsureS1RuntimeObjectsAsync();
  110. result.StandardRows = await TransformStandardAsync(scope, batchId, now, cancellationToken);
  111. result.DwdRows = await BuildDwdAsync(scope, batchId, now, result, cancellationToken);
  112. result.KpiRows = await BuildS1KpiValuesAsync(scope, batchId, now, cancellationToken);
  113. result.AtomicRows = await _atomicBuild.BuildOrderDeliveryDomainForAllDatesAsync(
  114. scope.TenantId, scope.FactoryId, batchId, cancellationToken);
  115. await MarkTransformRunSuccessAsync(runLogId, now, result);
  116. return result;
  117. }
  118. catch (Exception ex)
  119. {
  120. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  121. throw;
  122. }
  123. }
  124. private async Task EnsureS1RuntimeObjectsAsync()
  125. {
  126. await _db.Ado.ExecuteCommandAsync(
  127. """
  128. CREATE TABLE IF NOT EXISTS dwd_requirement_examine_detail (
  129. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  130. tenant_id BIGINT NOT NULL,
  131. factory_id BIGINT NULL,
  132. stat_date DATE NOT NULL,
  133. row_id BIGINT NOT NULL,
  134. parent_row_id BIGINT NULL,
  135. examine_id BIGINT NULL,
  136. order_entry_id BIGINT NULL,
  137. bill_no VARCHAR(100) NULL,
  138. morder_no VARCHAR(100) NULL,
  139. num VARCHAR(100) NULL,
  140. item_number VARCHAR(100) NULL,
  141. item_name VARCHAR(200) NULL,
  142. bom_number VARCHAR(100) NULL,
  143. model VARCHAR(200) NULL,
  144. kitting_time DATE NULL,
  145. item_type VARCHAR(50) NULL,
  146. erp_cls_name VARCHAR(80) NULL,
  147. qty DECIMAL(18,4) NULL,
  148. wastage DECIMAL(18,4) NULL,
  149. -- ── 齐套三量的标准名(B3 决定:不另起别名列)──────────────────────────────
  150. -- need_count / use_qty / lack_qty 就是本表「需求量 / 已满足量 / 缺口量」的标准列名,
  151. -- 不新增 required_qty / satisfied_qty / shortage_qty 别名列:
  152. -- 别名列 = 同一个数字有两个写入点,迟早对不上;而重算一套更会凭空造出与源系统不一致的口径。
  153. -- 恒等式 lack_qty + use_qty = need_count 由源系统保证,
  154. -- 实测(aidopdev,只读)1019060/1019060 行全部成立,最大绝对偏差 0.0000,三列均无 NULL。
  155. need_count DECIMAL(18,4) NULL COMMENT '标准列·需求量(齐套恒等式:lack_qty + use_qty = need_count)',
  156. sqty DECIMAL(18,4) NULL,
  157. use_qty DECIMAL(18,4) NULL COMMENT '标准列·已满足量(不要新增 satisfied_qty 别名)',
  158. self_lack_qty DECIMAL(18,4) NULL,
  159. lack_qty DECIMAL(18,4) NULL COMMENT '标准列·缺口量(不要新增 shortage_qty 别名)',
  160. mo_qty DECIMAL(18,4) NULL,
  161. make_qty DECIMAL(18,4) NULL,
  162. purchase_qty DECIMAL(18,4) NULL,
  163. purchase_occupy_qty DECIMAL(18,4) NULL,
  164. satisfy_time DATE NULL,
  165. have_ic_subs VARCHAR(10) NULL,
  166. substitute_code VARCHAR(100) NULL,
  167. create_time DATETIME NULL,
  168. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  169. sync_batch_id VARCHAR(100) NOT NULL,
  170. calc_batch_id VARCHAR(100) NOT NULL,
  171. calc_time DATETIME NOT NULL,
  172. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  173. -- ⚠️ 以下三列必须与 UpdateScripts/1.0.516.sql 的 ALTER ... ADD COLUMN 保持「同名同类型同顺序」。
  174. -- 本表没有 SqlSugar 实体、也没有 UpdateScripts 的 CREATE,schema 由「本 DDL + ALTER 脚本」两处共同定义:
  175. -- 新库走本 DDL 一次建全,旧库走 ALTER 追加。ALTER ADD COLUMN 默认追加到表尾,
  176. -- 故这三列在此也必须放在 update_time 之后,否则新库与迁移库的列序会分叉。
  177. bom_level INT NULL COMMENT 'BOM 层级 ← raw_data.$.level;1=成品根,>1=投入料',
  178. material_role VARCHAR(32) NULL COMMENT '物料角色枚举:FINISHED_GOOD_ROOT / INPUT_MATERIAL',
  179. is_current_flag TINYINT NOT NULL DEFAULT 0 COMMENT '1=属于当前快照批次;同一 (tenant, factory) 作用域内恒只有一个批次为 1',
  180. UNIQUE KEY uk_dwd_req_exam_detail (tenant_id, row_id, calc_batch_id),
  181. KEY idx_req_exam_tenant_batch (tenant_id, calc_batch_id),
  182. KEY idx_req_exam_bill (tenant_id, bill_no),
  183. KEY idx_req_exam_item (tenant_id, item_number),
  184. KEY idx_req_exam_current (tenant_id, is_current_flag, morder_no)
  185. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S1需求明细核验DWD';
  186. """);
  187. await _db.Ado.ExecuteCommandAsync(
  188. """
  189. INSERT INTO mdp_entity
  190. (tenant_id, source_id, entity_code, entity_name, entity_type, source_table_name, target_table_name, sync_mode, batch_size, status, remark, create_time, update_time)
  191. SELECT 0, s.id, 'S1_REQUIREMENT_EXAMINE_RESULT', 'S1需求核验结果主表', 'TABLE',
  192. 'b_examine_result', 'mdp_stg_so', 'FULL', 5000, 1, '需求核验结果主表,进入 S1 贴源层', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  193. FROM mdp_source s
  194. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  195. LIMIT 1
  196. ON DUPLICATE KEY UPDATE
  197. source_id=VALUES(source_id), entity_name=VALUES(entity_name), source_table_name=VALUES(source_table_name),
  198. target_table_name=VALUES(target_table_name), sync_mode=VALUES(sync_mode), status=VALUES(status),
  199. remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  200. """);
  201. await _db.Ado.ExecuteCommandAsync(
  202. """
  203. INSERT INTO mdp_entity
  204. (tenant_id, source_id, entity_code, entity_name, entity_type, source_table_name, target_table_name, sync_mode, batch_size, status, remark, create_time, update_time)
  205. SELECT 0, s.id, 'S1_REQUIREMENT_EXAMINE_DETAIL', 'S1需求核验BOM明细', 'TABLE',
  206. 'b_bom_child_examine', 'mdp_stg_so', 'FULL', 5000, 1, '需求核验BOM明细,进入 S1 贴源层和 DWD', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  207. FROM mdp_source s
  208. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  209. LIMIT 1
  210. ON DUPLICATE KEY UPDATE
  211. source_id=VALUES(source_id), entity_name=VALUES(entity_name), source_table_name=VALUES(source_table_name),
  212. target_table_name=VALUES(target_table_name), sync_mode=VALUES(sync_mode), status=VALUES(status),
  213. remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  214. """);
  215. await _db.Ado.ExecuteCommandAsync(
  216. """
  217. INSERT INTO mdp_entity
  218. (tenant_id, source_id, entity_code, entity_name, entity_type, source_table_name, target_table_name, sync_mode, incr_column, batch_size, status, remark, create_time, update_time)
  219. SELECT 0, s.id, v.entity_code, v.entity_name, 'TABLE', v.source_table_name, 'mdp_stg_so', v.sync_mode, v.incr_column, 5000, 1, v.remark, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  220. FROM mdp_source s
  221. JOIN (
  222. SELECT 'S1_PRODUCT_DESIGN' AS entity_code, 'S1产品设计主表' AS entity_name, 'ado_product_design' AS source_table_name, 'INCR' AS sync_mode, 'UpdateTime' AS incr_column, '产品设计主表,进入订单标准层' AS remark
  223. UNION ALL SELECT 'S1_PRODUCT_DESIGN_BOM', 'S1产品设计BOM', 'ado_product_design_bom', 'FULL', NULL, '产品设计BOM,表无时间戳列,只能全量抽取'
  224. UNION ALL SELECT 'S1_PRODUCT_DESIGN_ROUTING', 'S1产品设计工艺路线', 'ado_product_design_routing', 'FULL', NULL, '产品设计工艺路线,表无时间戳列,只能全量抽取'
  225. ) v
  226. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  227. ON DUPLICATE KEY UPDATE
  228. source_id=VALUES(source_id), entity_name=VALUES(entity_name), source_table_name=VALUES(source_table_name),
  229. target_table_name=VALUES(target_table_name), sync_mode=VALUES(sync_mode), incr_column=VALUES(incr_column),
  230. status=VALUES(status), remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  231. """);
  232. }
  233. /// <summary>Phase2/3:源→stg 走执行器(DB/API);遗留 SyncOneEntityAsync 不再调用。</summary>
  234. private async Task<int> SyncStagingAsync(S1MdpRunScope scope, string batchId, DateTime now, CancellationToken cancellationToken)
  235. {
  236. _ = now;
  237. var codes = S1MdpEntityConfig.All.Select(x => x.EntityCode);
  238. return await _stagingPuller.PullEntitiesAsync(
  239. codes, batchId, scope.TenantId, fullRefresh: true, taskCode: "S1_MDP_INBOUND",
  240. cancellationToken, scope.FactoryId, requireMatchingSourceTenant: true);
  241. }
  242. /// <summary>
  243. /// 双模式入站:仅抽 stg(可指定 DB 实体码或 *_API),再跑 std/DWD/KPI。
  244. /// </summary>
  245. public async Task<S1MdpSyncTransformResult> RunInboundAsync(
  246. S1MdpRunScope scope,
  247. IEnumerable<string>? entityCodes = null,
  248. bool fullRefresh = false,
  249. CancellationToken cancellationToken = default)
  250. {
  251. scope = S1MdpRunScope.Create(scope.TenantId, scope.FactoryId);
  252. cancellationToken.ThrowIfCancellationRequested();
  253. var now = DateTime.Now;
  254. var batchId = $"S1_MDP_IN_{scope.TenantId}_{scope.FactoryId}_{now:yyyyMMddHHmmss}";
  255. var runLogId = await InsertTransformRunLogAsync(scope, batchId, now, "INBOUND");
  256. var result = new S1MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  257. try
  258. {
  259. await EnsureS1RuntimeObjectsAsync();
  260. var codes = entityCodes?.ToList() ?? S1MdpEntityConfig.All.Select(x => x.EntityCode).ToList();
  261. result.StageRows = await _stagingPuller.PullEntitiesAsync(
  262. codes, batchId, scope.TenantId, fullRefresh, "S1_MDP_INBOUND", cancellationToken,
  263. scope.FactoryId, requireMatchingSourceTenant: true);
  264. result.StandardRows = await TransformStandardAsync(scope, batchId, now, cancellationToken);
  265. result.DwdRows = await BuildDwdAsync(scope, batchId, now, result, cancellationToken);
  266. result.KpiRows = await BuildS1KpiValuesAsync(scope, batchId, now, cancellationToken);
  267. result.AtomicRows = await _atomicBuild.BuildOrderDeliveryDomainForAllDatesAsync(
  268. scope.TenantId, scope.FactoryId, batchId, cancellationToken);
  269. await MarkTransformRunSuccessAsync(runLogId, now, result);
  270. return result;
  271. }
  272. catch (Exception ex)
  273. {
  274. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  275. throw;
  276. }
  277. }
  278. /// <summary>Phase3:已停用。源→stg 改由执行器路径;保留方法供紧急回滚对照,勿再调用。</summary>
  279. [Obsolete("Phase3: use MdpModuleStagingPuller / executors")]
  280. private async Task<int> SyncOneEntityAsync(S1MdpEntityConfig entity, string batchId, DateTime now, long tenantId)
  281. {
  282. if (tenantId <= 0) throw new InvalidOperationException("S1 同步日志必须指定有效 tenantId");
  283. var entityRow = await _db.Ado.SqlQuerySingleAsync<S1MdpEntityRow>(
  284. "SELECT id AS Id, entity_name AS EntityName FROM mdp_entity WHERE tenant_id=0 AND entity_code=@EntityCode LIMIT 1",
  285. new SugarParameter("@EntityCode", entity.EntityCode));
  286. if (entityRow == null) throw Oops.Oh($"未找到 MDP 实体配置:{entity.EntityCode}");
  287. var columns = await _db.Ado.SqlQueryAsync<S1ColumnRow>(
  288. """
  289. SELECT COLUMN_NAME AS ColumnName
  290. FROM information_schema.COLUMNS
  291. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@TableName
  292. ORDER BY ORDINAL_POSITION
  293. """,
  294. new SugarParameter("@TableName", entity.SourceTable));
  295. if (columns.Count == 0) throw Oops.Oh($"未找到源表:{entity.SourceTable}");
  296. var names = columns.Select(u => u.ColumnName).ToList();
  297. var tenantExpr = BuildOptionalColumnExpr(names, "tenant_id", "0");
  298. var factoryExpr = BuildOptionalColumnExpr(names, "factory_id", "NULL");
  299. var companyExpr = BuildOptionalColumnExpr(names, "company_id", "NULL");
  300. var sourceRowExpr = names.Any(u => string.Equals(u, entity.SourceRowIdExpression, StringComparison.OrdinalIgnoreCase))
  301. ? $"s.`{FindColumn(names, entity.SourceRowIdExpression)}`"
  302. : entity.SourceRowIdExpression;
  303. var rawDataExpr = BuildJsonObjectExpression(names);
  304. var rowsRead = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM `{entity.SourceTable}`");
  305. var logId = await InsertSyncLogAsync(tenantId, entityRow.Id, entityRow.EntityName, batchId, rowsRead);
  306. var started = DateTime.Now;
  307. try
  308. {
  309. var affected = await _db.Ado.ExecuteCommandAsync(
  310. $"""
  311. INSERT INTO `{entity.TargetTable}`
  312. (tenant_id, factory_id, company_id, source_system, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time, process_status, raw_data)
  313. SELECT
  314. {tenantExpr},
  315. {factoryExpr},
  316. {companyExpr},
  317. 'AIDOP',
  318. @SourceTable,
  319. CAST({sourceRowExpr} AS CHAR),
  320. CAST(COALESCE({entity.SourceBizKeyExpression}, CAST({sourceRowExpr} AS CHAR)) AS CHAR),
  321. @BatchId,
  322. @Now,
  323. 'PENDING',
  324. {rawDataExpr}
  325. FROM `{entity.SourceTable}` s
  326. ON DUPLICATE KEY UPDATE
  327. tenant_id=VALUES(tenant_id),
  328. factory_id=VALUES(factory_id),
  329. company_id=VALUES(company_id),
  330. source_row_id=VALUES(source_row_id),
  331. sync_batch_id=VALUES(sync_batch_id),
  332. sync_time=VALUES(sync_time),
  333. process_status=VALUES(process_status),
  334. raw_data=VALUES(raw_data),
  335. update_time=CURRENT_TIMESTAMP
  336. """,
  337. new SugarParameter("@SourceTable", entity.SourceTable),
  338. new SugarParameter("@BatchId", batchId),
  339. new SugarParameter("@Now", now));
  340. await MarkSyncLogSuccessAsync(logId, started, affected);
  341. return rowsRead;
  342. }
  343. catch (Exception ex)
  344. {
  345. await MarkSyncLogFailedAsync(logId, started, ex.Message);
  346. throw;
  347. }
  348. }
  349. private async Task<int> TransformStandardAsync(S1MdpRunScope scope, string batchId, DateTime now, CancellationToken cancellationToken)
  350. {
  351. using var db = _db.CopyNew();
  352. db.Ado.CommandTimeOut = Math.Max(db.Ado.CommandTimeOut, 180);
  353. var total = 0;
  354. foreach (var command in BuildStandardCommands(scope, batchId, now))
  355. {
  356. cancellationToken.ThrowIfCancellationRequested();
  357. total += await db.Ado.ExecuteCommandAsync(command.Sql, command.Parameters);
  358. }
  359. return total;
  360. }
  361. private async Task<int> BuildDwdAsync(
  362. S1MdpRunScope scope, string batchId, DateTime now,
  363. S1MdpSyncTransformResult result, CancellationToken cancellationToken)
  364. {
  365. using var db = _db.CopyNew();
  366. db.Ado.CommandTimeOut = Math.Max(db.Ado.CommandTimeOut, 180);
  367. var total = 0;
  368. foreach (var command in BuildDwdCommands(scope, batchId, now))
  369. {
  370. cancellationToken.ThrowIfCancellationRequested();
  371. total += await db.Ado.ExecuteCommandAsync(command.Sql, command.Parameters);
  372. }
  373. // DWD 阶段的最后一步:把本批次发布为「当前快照」。
  374. // 必须在所有 DWD 写入完成之后,否则读者会看到一个还没灌完的批次被标成 current。
  375. // 不计入 total —— total 语义是「本轮构建的 DWD 行数」,翻牌影响的行数属于发布动作,不是新建行。
  376. cancellationToken.ThrowIfCancellationRequested();
  377. var publishedRows = await PublishCurrentSnapshotAsync(db, scope, batchId);
  378. // 发布证据随运行日志一起落盘(见 MarkTransformRunSuccessAsync)。
  379. // 走到这里就意味着两条 DWD 写入与发布语句都已成功;若中途抛错,
  380. // 调用方的 catch 会写 FAILED,本字段保持 null —— 失败方向是保守的。
  381. result.Publication = new S1MdpSnapshotPublication
  382. {
  383. Table = "dwd_requirement_examine_detail",
  384. BatchId = batchId,
  385. FactoryId = scope.FactoryId,
  386. CurrentRows = publishedRows
  387. };
  388. return total;
  389. }
  390. /// <summary>
  391. /// 把 <paramref name="batchId"/> 原子地发布为该 (tenant, factory) 作用域下的当前快照。
  392. ///
  393. /// 【为什么必须是「一条」语句】
  394. /// BuildDwdAsync 全程没有事务(用的是 _db.CopyNew() 上的裸 ExecuteCommandAsync,没有 BeginTran)。
  395. /// 因此任何「两条语句」的写法都会在两条之间留下一个对并发读者可见的错误窗口:
  396. /// · 先清后置(先把旧批次清 0,再把新批次置 1)→ 中间窗口「零个当前批次」,
  397. /// 读者拿到空结果,会被当成「本工单已齐套、无缺料」——缺料监控直接漏报。
  398. /// · 先置后清(先把新批次置 1,再把旧批次清 0)→ 中间窗口「两个当前批次」,
  399. /// 读者按 is_current_flag=1 取数会同时拿到新旧两批,缺料数量凭空翻倍。
  400. /// · 而且第二条语句一旦失败/进程被杀,上述错误窗口就不再是「瞬时」,而是永久残留。
  401. /// 单条 UPDATE 由 InnoDB 保证语句级原子性,不存在中间可见状态,也没有"第二条没跑成"的失败模式。
  402. ///
  403. /// 【为什么 factory_id 要 COALESCE 归一】
  404. /// 同一次 run(scope.FactoryId=1)产出的行,factory_id 既可能是 1 也可能是 NULL:
  405. /// 入库过滤用的是 COALESCE(NULLIF(d.factory_id,0),1)=@FactoryId,而落库值取的是
  406. /// COALESCE(NULLIF(d.factory_id,0), NULLIF(r.factory_id,0), NULLIF(so.factory_id,0)),三者全空时就是 NULL。
  407. /// 实测(aidopdev,只读):1019060 行里 296189 行 factory_id IS NULL,
  408. /// 且 3439 个批次里有 594 个批次「同时含 factory_id=1 和 factory_id IS NULL」。
  409. /// 若作用域谓词直接写 factory_id=@FactoryId:
  410. /// ① NULL 行永远不被 = 匹配(SQL 三值逻辑,NULL=1 得 NULL 而非 TRUE),旧批次的 NULL 行会一直挂着 is_current_flag=1;
  411. /// ② 新批次也只有一半的行被置 1。
  412. /// 结果就是任务里点名要避免的「两个当前批次」永久共存。
  413. /// 故必须用 COALESCE(NULLIF(factory_id,0),1) 归一 —— 这也是本仓库通用的 factory 作用域写法。
  414. /// </summary>
  415. /// <summary>
  416. /// 发布当前快照,并返回本批次发布后的当前行数。
  417. /// <para>返回 0 是<b>合法结果</b>——源侧本轮没有数据时,成功发布的就是一个空快照。
  418. /// 调用方把这个数字记进运行日志,下游才能把「正常的空」与「没发布成」区分开。</para>
  419. /// </summary>
  420. private static async Task<int> PublishCurrentSnapshotAsync(ISqlSugarClient db, S1MdpRunScope scope, string batchId)
  421. {
  422. await db.Ado.ExecuteCommandAsync(
  423. """
  424. UPDATE dwd_requirement_examine_detail
  425. SET is_current_flag = CASE WHEN calc_batch_id=@BatchId THEN 1 ELSE 0 END
  426. WHERE tenant_id=@TenantId
  427. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  428. AND (calc_batch_id=@BatchId OR is_current_flag=1)
  429. """,
  430. new SugarParameter("@BatchId", batchId),
  431. new SugarParameter("@TenantId", scope.TenantId),
  432. new SugarParameter("@FactoryId", scope.FactoryId));
  433. // 复核发布结果而不是相信影响行数:ODKU 与 CASE 的影响行数都不等于「当前有几行」。
  434. return await db.Ado.GetIntAsync(
  435. """
  436. SELECT COUNT(*) FROM dwd_requirement_examine_detail
  437. WHERE tenant_id=@TenantId
  438. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  439. AND calc_batch_id=@BatchId
  440. AND is_current_flag=1
  441. """,
  442. new SugarParameter("@BatchId", batchId),
  443. new SugarParameter("@TenantId", scope.TenantId),
  444. new SugarParameter("@FactoryId", scope.FactoryId));
  445. }
  446. private async Task<int> BuildS1KpiValuesAsync(S1MdpRunScope scope, string batchId, DateTime now, CancellationToken cancellationToken)
  447. {
  448. var affected = 0;
  449. for (var dayOffset = 13; dayOffset >= 0; dayOffset--)
  450. {
  451. var statDate = now.Date.AddDays(-dayOffset);
  452. var rows = await CalculateS1KpiValuesAsync(scope, batchId, statDate);
  453. foreach (var row in rows)
  454. {
  455. cancellationToken.ThrowIfCancellationRequested();
  456. if (row.TenantId != scope.TenantId || row.FactoryId != scope.FactoryId)
  457. continue;
  458. affected += await UpsertS1KpiValueAsync(row, statDate, now);
  459. }
  460. }
  461. return affected;
  462. }
  463. private async Task<List<S1KpiCalcRow>> CalculateS1KpiValuesAsync(S1MdpRunScope scope, string batchId, DateTime statDate)
  464. {
  465. return await _db.Ado.SqlQueryAsync<S1KpiCalcRow>(
  466. """
  467. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  468. 'S1_L1_001' AS MetricCode,
  469. ROUND(AVG(TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) / 24), 4) AS MetricValue
  470. FROM mdp_std_so
  471. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  472. AND order_date IS NOT NULL
  473. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  474. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) >= order_date
  475. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  476. UNION ALL
  477. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  478. 'S1_L1_002' AS MetricCode,
  479. ROUND(100 * SUM(CASE WHEN TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) <= 72 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  480. FROM mdp_std_so
  481. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  482. AND order_date IS NOT NULL
  483. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  484. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  485. UNION ALL
  486. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  487. 'S1_L1_003' AS MetricCode,
  488. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(planner_no, '')), 1), 4) AS MetricValue
  489. FROM mdp_std_so
  490. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  491. AND order_date IS NOT NULL AND order_date <= @StatDate
  492. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  493. UNION ALL
  494. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  495. 'S1_L1_004' AS MetricCode,
  496. ROUND(AVG(GREATEST(IFNULL(remaining_qty, 0), 0)) / GREATEST(SUM(IFNULL(shipped_qty, 0)), 1) * 30, 4) AS MetricValue
  497. FROM dwd_ship_trans
  498. WHERE calc_batch_id=@BatchId
  499. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  500. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  501. HAVING SUM(IFNULL(shipped_qty, 0)) > 0
  502. UNION ALL
  503. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  504. 'S1_L2_010' AS MetricCode,
  505. ROUND(AVG(TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) / 24), 4) AS MetricValue
  506. FROM mdp_std_so
  507. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  508. AND order_date IS NOT NULL
  509. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  510. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) >= order_date
  511. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  512. UNION ALL
  513. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  514. 'S1_L2_011' AS MetricCode,
  515. ROUND(100 * SUM(CASE WHEN TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) <= 72 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  516. FROM mdp_std_so
  517. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  518. AND order_date IS NOT NULL
  519. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  520. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  521. UNION ALL
  522. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  523. 'S1_L2_012' AS MetricCode,
  524. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(planner_no, '')), 1), 4) AS MetricValue
  525. FROM mdp_std_so
  526. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  527. AND order_date IS NOT NULL AND order_date <= @StatDate
  528. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  529. UNION ALL
  530. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  531. 'S1_L2_013' AS MetricCode,
  532. ROUND(100 * SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  533. FROM dwd_ship_trans
  534. WHERE calc_batch_id=@BatchId
  535. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  536. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  537. UNION ALL
  538. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  539. 'S1_L2_014' AS MetricCode,
  540. ROUND(100 * SUM(CASE WHEN shipped_qty >= planned_ship_qty AND planned_ship_qty > 0 THEN 1 ELSE 0 END)
  541. / GREATEST(SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END), 1), 4) AS MetricValue
  542. FROM dwd_ship_trans
  543. WHERE calc_batch_id=@BatchId
  544. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  545. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  546. HAVING SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) > 0
  547. UNION ALL
  548. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  549. 'S1_L2_015' AS MetricCode,
  550. ROUND(100 * SUM(CASE WHEN linkage_status = 'LINKED' THEN 1 ELSE 0 END)
  551. / GREATEST(SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END), 1), 4) AS MetricValue
  552. FROM dwd_ship_trans
  553. WHERE calc_batch_id=@BatchId
  554. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  555. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  556. HAVING SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) > 0
  557. UNION ALL
  558. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  559. 'S1_L2_001' AS MetricCode,
  560. ROUND(AVG(TIMESTAMPDIFF(HOUR, create_time, update_time) / 24), 4) AS MetricValue
  561. FROM (
  562. SELECT tenant_id,
  563. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f') AS create_time,
  564. COALESCE(
  565. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  566. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  567. ) AS update_time
  568. FROM mdp_stg_so
  569. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  570. AND source_table = 'ado_contract_review'
  571. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-PROC-%'
  572. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-DEL-%'
  573. ) t
  574. WHERE create_time IS NOT NULL AND update_time IS NOT NULL AND update_time >= create_time
  575. GROUP BY tenant_id
  576. UNION ALL
  577. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  578. 'S1_L2_002' AS MetricCode,
  579. ROUND(100 * SUM(CASE WHEN TIMESTAMPDIFF(HOUR, create_time, update_time) <= 72 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  580. FROM (
  581. SELECT tenant_id,
  582. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f') AS create_time,
  583. COALESCE(
  584. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  585. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  586. ) AS update_time
  587. FROM mdp_stg_so
  588. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  589. AND source_table = 'ado_contract_review'
  590. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-PROC-%'
  591. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-DEL-%'
  592. ) t
  593. WHERE create_time IS NOT NULL AND update_time IS NOT NULL AND update_time >= create_time
  594. GROUP BY tenant_id
  595. UNION ALL
  596. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  597. 'S1_L2_003' AS MetricCode,
  598. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(owner_account, '')), 1), 4) AS MetricValue
  599. FROM (
  600. SELECT tenant_id,
  601. COALESCE(
  602. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ResponsibleAccount')),
  603. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))
  604. ) AS owner_account
  605. FROM mdp_stg_so
  606. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  607. AND source_table = 'ado_contract_review'
  608. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-PROC-%'
  609. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-DEL-%'
  610. ) t
  611. GROUP BY tenant_id
  612. UNION ALL
  613. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  614. 'S1_L2_007' AS MetricCode,
  615. ROUND(AVG(TIMESTAMPDIFF(HOUR,create_time,update_time)/24),4) AS MetricValue
  616. FROM (
  617. SELECT tenant_id,
  618. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f') AS create_time,
  619. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f') AS update_time
  620. FROM mdp_stg_so
  621. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  622. AND source_table='ado_contract_review'
  623. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') LIKE 'UAT-PROC-%'
  624. ) t
  625. WHERE create_time IS NOT NULL AND update_time>=create_time
  626. GROUP BY tenant_id
  627. UNION ALL
  628. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  629. 'S1_L2_008' AS MetricCode,
  630. ROUND(100*SUM(CASE WHEN TIMESTAMPDIFF(HOUR,create_time,update_time)<=72 THEN 1 ELSE 0 END)
  631. / NULLIF(COUNT(*),0),4) AS MetricValue
  632. FROM (
  633. SELECT tenant_id,
  634. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f') AS create_time,
  635. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f') AS update_time
  636. FROM mdp_stg_so
  637. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  638. AND source_table='ado_contract_review'
  639. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') LIKE 'UAT-PROC-%'
  640. ) t
  641. WHERE create_time IS NOT NULL AND update_time IS NOT NULL
  642. GROUP BY tenant_id
  643. UNION ALL
  644. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  645. 'S1_L2_009' AS MetricCode,
  646. ROUND(COUNT(*)/NULLIF(COUNT(DISTINCT NULLIF(owner_account,'')),0),4) AS MetricValue
  647. FROM (
  648. SELECT tenant_id,
  649. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ResponsibleAccount')),
  650. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))) AS owner_account
  651. FROM mdp_stg_so
  652. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  653. AND source_table='ado_contract_review'
  654. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') LIKE 'UAT-PROC-%'
  655. ) t
  656. GROUP BY tenant_id
  657. UNION ALL
  658. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  659. 'S1_L2_004' AS MetricCode,
  660. ROUND(AVG(cycle_hours / 24), 4) AS MetricValue
  661. FROM (
  662. SELECT tenant_id,
  663. COALESCE(
  664. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingDesignCycle')) REGEXP '^-?[0-9]+$'
  665. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingDesignCycle')) AS DECIMAL(18,6)) END,
  666. TIMESTAMPDIFF(HOUR,
  667. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualStart')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  668. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualEnd')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  669. )
  670. ) AS cycle_hours
  671. FROM mdp_stg_so
  672. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  673. AND source_table = 'ado_product_design'
  674. ) t
  675. WHERE cycle_hours IS NOT NULL AND cycle_hours >= 0
  676. GROUP BY tenant_id
  677. UNION ALL
  678. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  679. 'S1_L2_005' AS MetricCode,
  680. ROUND(100 * SUM(CASE WHEN actual_end <= plan_end THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  681. FROM (
  682. SELECT tenant_id,
  683. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingPlanEnd')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f') AS plan_end,
  684. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualEnd')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f') AS actual_end
  685. FROM mdp_stg_so
  686. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  687. AND source_table = 'ado_product_design'
  688. ) t
  689. WHERE plan_end IS NOT NULL AND actual_end IS NOT NULL
  690. GROUP BY tenant_id
  691. UNION ALL
  692. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  693. 'S1_L2_006' AS MetricCode,
  694. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(design_lead, '')), 1), 4) AS MetricValue
  695. FROM (
  696. SELECT tenant_id,
  697. COALESCE(
  698. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DesignLeadAccount')),
  699. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DesignLeadName')),
  700. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))
  701. ) AS design_lead
  702. FROM mdp_stg_so
  703. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  704. AND source_table = 'ado_product_design'
  705. ) t
  706. GROUP BY tenant_id
  707. UNION ALL
  708. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  709. CASE stage_no
  710. WHEN 1 THEN 'S1_L3_001'
  711. WHEN 2 THEN 'S1_L3_002'
  712. WHEN 3 THEN 'S1_L3_003'
  713. WHEN 4 THEN 'S1_L3_004'
  714. WHEN 5 THEN 'S1_L3_005'
  715. END AS MetricCode,
  716. ROUND(AVG(cycle_hours), 4) AS MetricValue
  717. FROM (
  718. SELECT tenant_id,
  719. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) REGEXP '^-?[0-9]+$'
  720. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) END AS stage_no,
  721. COALESCE(
  722. TIMESTAMPDIFF(HOUR,
  723. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  724. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  725. ),
  726. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  727. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  728. ) AS cycle_hours
  729. FROM mdp_stg_so
  730. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  731. AND source_table = 'ado_contract_review_flow'
  732. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-PROC-%'
  733. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-DEL-%'
  734. ) t
  735. WHERE stage_no BETWEEN 1 AND 5
  736. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  737. GROUP BY tenant_id, stage_no
  738. UNION ALL
  739. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  740. CASE stage_no
  741. WHEN 1 THEN 'S1_L3_101'
  742. WHEN 2 THEN 'S1_L3_102'
  743. WHEN 3 THEN 'S1_L3_103'
  744. WHEN 4 THEN 'S1_L3_104'
  745. WHEN 5 THEN 'S1_L3_105'
  746. END AS MetricCode,
  747. ROUND(AVG(cycle_hours), 4) AS MetricValue
  748. FROM (
  749. SELECT tenant_id,
  750. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) REGEXP '^-?[0-9]+$'
  751. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) END AS stage_no,
  752. COALESCE(
  753. TIMESTAMPDIFF(HOUR,
  754. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  755. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  756. ),
  757. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  758. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  759. ) AS cycle_hours
  760. FROM mdp_stg_so
  761. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  762. AND source_table = 'ado_contract_review_flow'
  763. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-PROC-%'
  764. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-DEL-%'
  765. ) t
  766. WHERE stage_no BETWEEN 1 AND 5
  767. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  768. GROUP BY tenant_id, stage_no
  769. UNION ALL
  770. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  771. metric_code AS MetricCode,
  772. ROUND(AVG(cycle_hours), 4) AS MetricValue
  773. FROM (
  774. SELECT tenant_id,
  775. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')), ''), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  776. WHEN 'LAW' THEN 'S1_L4_001'
  777. WHEN 'PRE_SALES' THEN 'S1_L4_002'
  778. WHEN 'MPS' THEN 'S1_L4_003'
  779. WHEN 'TEST' THEN 'S1_L4_004'
  780. WHEN '法律事务部' THEN 'S1_L4_001'
  781. WHEN '技术售前组' THEN 'S1_L4_002'
  782. WHEN '综合主计划' THEN 'S1_L4_003'
  783. WHEN '试验站' THEN 'S1_L4_004'
  784. END AS metric_code,
  785. COALESCE(
  786. TIMESTAMPDIFF(HOUR,
  787. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  788. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  789. ),
  790. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  791. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  792. ) AS cycle_hours
  793. FROM mdp_stg_so
  794. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  795. AND source_table = 'ado_contract_review_flow'
  796. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) = '1'
  797. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-PROC-%'
  798. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-DEL-%'
  799. ) t
  800. WHERE metric_code IS NOT NULL
  801. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  802. GROUP BY tenant_id, metric_code
  803. UNION ALL
  804. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  805. metric_code AS MetricCode,
  806. ROUND(AVG(cycle_hours), 4) AS MetricValue
  807. FROM (
  808. SELECT tenant_id,
  809. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')), ''), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  810. WHEN 'LAW' THEN 'S1_L4_101'
  811. WHEN 'PRE_SALES' THEN 'S1_L4_102'
  812. WHEN 'MPS' THEN 'S1_L4_103'
  813. WHEN 'TEST' THEN 'S1_L4_104'
  814. WHEN '法律事务部' THEN 'S1_L4_101'
  815. WHEN '技术售前组' THEN 'S1_L4_102'
  816. WHEN '综合主计划' THEN 'S1_L4_103'
  817. WHEN '试验站' THEN 'S1_L4_104'
  818. END AS metric_code,
  819. COALESCE(
  820. TIMESTAMPDIFF(HOUR,
  821. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  822. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  823. ),
  824. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  825. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  826. ) AS cycle_hours
  827. FROM mdp_stg_so
  828. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  829. AND source_table = 'ado_contract_review_flow'
  830. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) = '1'
  831. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-PROC-%'
  832. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-DEL-%'
  833. ) t
  834. WHERE metric_code IS NOT NULL
  835. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  836. GROUP BY tenant_id, metric_code
  837. UNION ALL
  838. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  839. CONCAT('S1_L3_',branch_code,LPAD(stage_no,2,'0')) AS MetricCode,
  840. ROUND(AVG(cycle_hours),4) AS MetricValue
  841. FROM (
  842. SELECT tenant_id,
  843. CASE
  844. WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-PROC-%' THEN '3'
  845. WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-DEL-%' THEN '4'
  846. END AS branch_code,
  847. CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) AS stage_no,
  848. TIMESTAMPDIFF(
  849. HOUR,
  850. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f'),
  851. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f')
  852. ) AS cycle_hours
  853. FROM mdp_stg_so
  854. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  855. AND source_table='ado_contract_review_flow'
  856. AND (COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-PROC-%'
  857. OR COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-DEL-%')
  858. ) t
  859. WHERE branch_code IS NOT NULL AND stage_no BETWEEN 1 AND 5
  860. AND cycle_hours IS NOT NULL AND cycle_hours>=0
  861. GROUP BY tenant_id,branch_code,stage_no
  862. UNION ALL
  863. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  864. CONCAT('S1_L4_',branch_code,dept_code) AS MetricCode,
  865. ROUND(AVG(cycle_hours),4) AS MetricValue
  866. FROM (
  867. SELECT tenant_id,
  868. CASE
  869. WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-PROC-%' THEN '3'
  870. WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-DEL-%' THEN '4'
  871. END AS branch_code,
  872. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')),''),
  873. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  874. WHEN 'LAW' THEN '01' WHEN '法律事务部' THEN '01'
  875. WHEN 'PRE_SALES' THEN '02' WHEN '技术售前组' THEN '02'
  876. WHEN 'MPS' THEN '03' WHEN '综合主计划' THEN '03'
  877. WHEN 'TEST' THEN '04' WHEN '试验站' THEN '04'
  878. END AS dept_code,
  879. TIMESTAMPDIFF(
  880. HOUR,
  881. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f'),
  882. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f')
  883. ) AS cycle_hours
  884. FROM mdp_stg_so
  885. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  886. AND source_table='ado_contract_review_flow'
  887. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo'))='1'
  888. AND (COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-PROC-%'
  889. OR COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-DEL-%')
  890. ) t
  891. WHERE branch_code IS NOT NULL AND dept_code IS NOT NULL
  892. AND cycle_hours IS NOT NULL AND cycle_hours>=0
  893. GROUP BY tenant_id,branch_code,dept_code
  894. """,
  895. new SugarParameter("@BatchId", batchId),
  896. new SugarParameter("@StatDate", statDate),
  897. new SugarParameter("@TenantId", scope.TenantId),
  898. new SugarParameter("@FactoryId", scope.FactoryId));
  899. }
  900. private async Task<int> UpsertS1KpiValueAsync(S1KpiCalcRow row, DateTime statDate, DateTime now)
  901. {
  902. var meta = await _db.Ado.SqlQuerySingleAsync<S1KpiMetaRow>(
  903. """
  904. SELECT MetricLevel, Direction, YellowThreshold, RedThreshold
  905. FROM ado_smart_ops_kpi_master
  906. WHERE TenantId=@TenantId AND ModuleCode='S1' AND MetricCode=@MetricCode AND IsEnabled=1
  907. LIMIT 1
  908. """,
  909. new SugarParameter("@TenantId", row.TenantId),
  910. new SugarParameter("@MetricCode", row.MetricCode));
  911. if (meta == null || row.MetricValue == null)
  912. return 0;
  913. var table = ResolveKpiValueTable(meta.MetricLevel);
  914. var current = await _db.Ado.SqlQuerySingleAsync<S1KpiValueRow>(
  915. $"""
  916. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  917. FROM {table}
  918. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  919. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  920. ORDER BY id
  921. LIMIT 1
  922. """,
  923. new SugarParameter("@TenantId", row.TenantId),
  924. new SugarParameter("@FactoryId", row.FactoryId),
  925. new SugarParameter("@MetricCode", row.MetricCode),
  926. new SugarParameter("@BizDate", statDate));
  927. var prior = await _db.Ado.SqlQuerySingleAsync<S1KpiValueRow>(
  928. $"""
  929. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  930. FROM {table}
  931. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  932. AND metric_code=@MetricCode AND biz_date<@BizDate AND is_deleted=0
  933. ORDER BY biz_date DESC, id DESC
  934. LIMIT 1
  935. """,
  936. new SugarParameter("@TenantId", row.TenantId),
  937. new SugarParameter("@FactoryId", row.FactoryId),
  938. new SugarParameter("@MetricCode", row.MetricCode),
  939. new SugarParameter("@BizDate", statDate));
  940. var actual = Math.Round(row.MetricValue.Value, 4);
  941. var snap = await _kpiTargetResolver.ResolveAsync(row.TenantId, row.FactoryId, row.MetricCode, "S1", statDate);
  942. var target = snap.TargetValue;
  943. var status = AidopS4KpiMerge.AchievementLevel(actual, target, meta.Direction, meta.YellowThreshold, meta.RedThreshold);
  944. var trend = ResolveTrendFlag(actual, prior?.MetricValue);
  945. if (current != null)
  946. {
  947. return await _db.Ado.ExecuteCommandAsync(
  948. $"""
  949. UPDATE {table}
  950. SET metric_value=@MetricValue, target_value=@TargetValue, status_color=@StatusColor, trend_flag=@TrendFlag,
  951. target_config_id=@TargetConfigId, target_source=@TargetSource, target_resolved_at=@TargetResolvedAt,
  952. is_active=1, status='ACTIVE', calc_time=@CalcTime, update_time=@CalcTime
  953. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  954. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  955. """,
  956. new SugarParameter("@MetricValue", actual),
  957. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  958. new SugarParameter("@StatusColor", status),
  959. new SugarParameter("@TrendFlag", trend),
  960. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  961. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  962. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt),
  963. new SugarParameter("@CalcTime", now),
  964. new SugarParameter("@TenantId", row.TenantId),
  965. new SugarParameter("@FactoryId", row.FactoryId),
  966. new SugarParameter("@MetricCode", row.MetricCode),
  967. new SugarParameter("@BizDate", statDate));
  968. }
  969. var nextId = Yitter.IdGenerator.YitIdHelper.NextId();
  970. return await _db.Ado.ExecuteCommandAsync(
  971. $"""
  972. INSERT INTO {table}
  973. (id, tenant_id, factory_id, status, biz_date, create_time, update_time, is_deleted, is_active,
  974. module_code, metric_code, metric_value, target_value, status_color, trend_flag, calc_time,
  975. target_config_id, target_source, target_resolved_at)
  976. VALUES
  977. (@Id, @TenantId, @FactoryId, 'ACTIVE', @BizDate, @CalcTime, @CalcTime, 0, 1,
  978. 'S1', @MetricCode, @MetricValue, @TargetValue, @StatusColor, @TrendFlag, @CalcTime,
  979. @TargetConfigId, @TargetSource, @TargetResolvedAt)
  980. """,
  981. new SugarParameter("@Id", nextId),
  982. new SugarParameter("@TenantId", row.TenantId),
  983. new SugarParameter("@FactoryId", row.FactoryId),
  984. new SugarParameter("@BizDate", statDate),
  985. new SugarParameter("@CalcTime", now),
  986. new SugarParameter("@MetricCode", row.MetricCode),
  987. new SugarParameter("@MetricValue", actual),
  988. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  989. new SugarParameter("@StatusColor", status),
  990. new SugarParameter("@TrendFlag", trend),
  991. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  992. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  993. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt));
  994. }
  995. private IEnumerable<S1MdpSqlCommand> BuildStandardCommands(S1MdpRunScope scope, string batchId, DateTime now)
  996. {
  997. yield return Cmd(
  998. """
  999. INSERT INTO mdp_std_so
  1000. (tenant_id, factory_id, company_id, source_system, order_id, order_entry_id, order_no, order_line, order_type,
  1001. customer_id, customer_no, customer_name, customer_order_no, country, item_code, item_name, item_spec,
  1002. map_number, map_name, bom_number, unit, order_qty, delivered_notice_qty, delivered_qty, price, tax_price,
  1003. amount, total_amount, order_date, customer_request_date, plan_delivery_date, promised_delivery_date,
  1004. capacity_date, material_ready_date, planner_no, planner_name, order_status, review_status, review_stage,
  1005. flow_state, progress, urgent, closed, deleted_flag, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  1006. SELECT
  1007. COALESCE(NULLIF(e.tenant_id, 0), NULLIF(h.tenant_id, 0),
  1008. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1009. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) AS SIGNED) END),
  1010. COALESCE(NULLIF(NULLIF(e.factory_id, 0), e.tenant_id),
  1011. NULLIF(NULLIF(h.factory_id, 0), h.tenant_id)),
  1012. NULLIF(e.company_id, 0),
  1013. 'AIDOP',
  1014. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.Id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.Id')) AS SIGNED) END,
  1015. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.Id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.Id')) AS SIGNED) END,
  1016. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no')), e.source_biz_key),
  1017. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.entry_seq')), JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.Id'))) AS CHAR),
  1018. CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.order_type')) AS CHAR),
  1019. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_id')) AS SIGNED) END,
  1020. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_no')),
  1021. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_name')),
  1022. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.custom_order_bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_from'))),
  1023. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.country')),
  1024. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.item_number')),
  1025. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.item_name')),
  1026. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.specification')),
  1027. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.map_number')),
  1028. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.map_name')),
  1029. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bom_number')),
  1030. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.unit')),
  1031. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.qty')) AS DECIMAL(18,6)) END,
  1032. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.deliver_notice_count')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.deliver_notice_count')) AS DECIMAL(18,6)) END,
  1033. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.deliver_count')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.deliver_count')) AS DECIMAL(18,6)) END,
  1034. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.price')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.price')) AS DECIMAL(18,6)) END,
  1035. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tax_price')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tax_price')) AS DECIMAL(18,6)) END,
  1036. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.amount')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.amount')) AS DECIMAL(18,6)) END,
  1037. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.total_amount')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.total_amount')) AS DECIMAL(18,6)) END,
  1038. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.date')), 'null'), ''),
  1039. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.rdate')), 'null'), ''),
  1040. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.plan_date')), 'null'), ''),
  1041. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.date')), 'null'), ''),
  1042. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.sys_capacity_date')), 'null'), ''),
  1043. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.sys_material_date')), 'null'), ''),
  1044. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.planner_no')),
  1045. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.planner_name')),
  1046. CASE WHEN JSON_EXTRACT(h.raw_data,'$.closed') IN (1, true) THEN 'CLOSED' ELSE 'OPEN' END,
  1047. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.FlowStatus')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.flowstate'))),
  1048. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.CurrentDept')), JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.CurrentStage'))),
  1049. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.flowstate')),
  1050. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.progress')),
  1051. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(COALESCE(e.raw_data, h.raw_data),'$.urgent')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(COALESCE(e.raw_data, h.raw_data),'$.urgent')) AS SIGNED) END,
  1052. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.closed')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.closed')) AS SIGNED) END,
  1053. COALESCE(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(COALESCE(e.raw_data, h.raw_data),'$.IsDeleted')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(COALESCE(e.raw_data, h.raw_data),'$.IsDeleted')) AS SIGNED) END, 0),
  1054. e.source_table,
  1055. e.source_row_id,
  1056. e.source_biz_key,
  1057. @BatchId,
  1058. @Now
  1059. FROM mdp_stg_so e
  1060. LEFT JOIN mdp_stg_so h ON h.source_table='crm_seorder'
  1061. AND h.tenant_id = e.tenant_id
  1062. AND h.sync_batch_id = (SELECT lb.sync_batch_id FROM mdp_stg_so lb
  1063. WHERE lb.source_table='crm_seorder' AND lb.tenant_id=@TenantId
  1064. ORDER BY lb.sync_time DESC, lb.id DESC LIMIT 1)
  1065. AND JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.Id')) = JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.seorder_id'))
  1066. LEFT JOIN mdp_stg_so r ON r.source_table='ado_contract_review'
  1067. AND r.tenant_id = e.tenant_id
  1068. AND JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.BillNo')) = COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.contract_no')), JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no')))
  1069. WHERE e.source_table='crm_seorderentry'
  1070. AND e.tenant_id=@TenantId
  1071. AND COALESCE(NULLIF(e.factory_id, 0), 1)=@FactoryId
  1072. -- ── 只取「最新一次 sync_batch_id」,使 mdp_std_so 成为源的真实镜像 ────────────────
  1073. -- 贴源层是纯 upsert、从不淘汰:源侧硬删除的订单行会永远留在 mdp_stg_so 里,
  1074. -- 而本 INSERT 原先不带任何批次条件、每轮重读整张贴源表,于是把这些孤儿一并
  1075. -- 物化进标准层。实测租户 797403760988229:贴源 414 行,源侧实际只剩 28 行,
  1076. -- 其余 386 行中 361 行是源侧已删除的孤儿、25 行是 biz_key 分隔符从 ':' 改成 '#'
  1077. -- 之后被弃用的旧格式重复行(1.0.259 改的 biz_key_expr)。
  1078. --
  1079. -- 这些孤儿不是只影响 S8:mdp_std_so 同时驱动 dwd_ship_trans(每条孤儿凭空
  1080. -- 造出一条早已过期、必然判 DELAYED 的发运行)与 S1/S7/S9 的多项 KPI。
  1081. --
  1082. -- 该等价关系(不在最新批次 = 源侧已删除)成立的前提是该 entity 每轮都全量重读。
  1083. -- 本批已把 S1_SEORDER / S1_SEORDER_ENTRY 的 sync_mode 由 INCR 改为 FULL
  1084. -- (见 UpdateScripts/1.0.521.sql):MdpSyncWindowResolver 对 FULL 实体会强制
  1085. -- ctx.FullRefresh=true,从而在**任何**调用路径上都跳过增量水位(含
  1086. -- RunInboundAsync 与 MdpHotWatchService 这两条 fullRefresh=false 的路径)。
  1087. -- ⚠️ 若哪天把它们改回 INCREMENTAL,本过滤会把「未变更的存量行」误当孤儿丢掉,
  1088. -- 届时必须回到贴源层整体替换的正解,不能只改这里。
  1089. AND e.sync_batch_id = (SELECT lb.sync_batch_id FROM mdp_stg_so lb
  1090. WHERE lb.source_table='crm_seorderentry' AND lb.tenant_id=@TenantId
  1091. ORDER BY lb.sync_time DESC, lb.id DESC LIMIT 1)
  1092. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no'))), '') <> ''
  1093. AND COALESCE(NULLIF(e.tenant_id, 0), NULLIF(h.tenant_id, 0),
  1094. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1095. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1096. ON DUPLICATE KEY UPDATE
  1097. factory_id=VALUES(factory_id), customer_no=VALUES(customer_no), customer_name=VALUES(customer_name), item_code=VALUES(item_code),
  1098. item_name=VALUES(item_name), item_spec=VALUES(item_spec), order_qty=VALUES(order_qty),
  1099. delivered_notice_qty=VALUES(delivered_notice_qty), delivered_qty=VALUES(delivered_qty),
  1100. plan_delivery_date=VALUES(plan_delivery_date), promised_delivery_date=VALUES(promised_delivery_date),
  1101. capacity_date=VALUES(capacity_date), material_ready_date=VALUES(material_ready_date),
  1102. order_status=VALUES(order_status), review_status=VALUES(review_status), review_stage=VALUES(review_stage),
  1103. flow_state=VALUES(flow_state), progress=VALUES(progress), closed=VALUES(closed), deleted_flag=VALUES(deleted_flag),
  1104. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  1105. """, scope, batchId, now);
  1106. yield return Cmd(
  1107. """
  1108. INSERT INTO mdp_std_ship_trans
  1109. (tenant_id, factory_id, company_id, source_system, trans_type, plan_id, plan_no, plan_line, order_id, order_entry_id,
  1110. order_no, order_line, customer_no, customer_name, country, item_code, item_name, item_spec, qty, plan_qty,
  1111. weight, volume, order_date, plan_ship_date, shipping_site, shipping_address, consignee, telephone,
  1112. status, confirm_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  1113. SELECT
  1114. COALESCE(NULLIF(d.tenant_id, 0),
  1115. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1116. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1117. NULLIF(d.factory_id, 0),
  1118. NULLIF(d.company_id, 0),
  1119. 'AIDOP',
  1120. 'SHIP_PLAN',
  1121. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id')) AS SIGNED) END,
  1122. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.LotSerial')), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id')) AS CHAR)),
  1123. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS CHAR),
  1124. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.seorder_id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.seorder_id')) AS SIGNED) END,
  1125. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sentry_id')) AS SIGNED) END,
  1126. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr'))),
  1127. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sentry_id')) AS CHAR),
  1128. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CustomNo')),
  1129. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CustomName')),
  1130. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Country')),
  1131. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')),
  1132. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemName')),
  1133. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Specification')),
  1134. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Qty')) AS DECIMAL(18,6)) END,
  1135. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Qty')) AS DECIMAL(18,6)) END,
  1136. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Weight')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Weight')) AS DECIMAL(18,6)) END,
  1137. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Volume')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Volume')) AS DECIMAL(18,6)) END,
  1138. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdDate')), 'null'), ''),
  1139. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingDate')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CreateTime'))), 'null'), ''),
  1140. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingSite')),
  1141. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingAddress')),
  1142. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Consignee')),
  1143. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Telephone')),
  1144. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Status')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Status'))),
  1145. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.IsConfirm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.IsConfirm'))) AS CHAR),
  1146. d.source_table,
  1147. d.source_row_id,
  1148. d.source_biz_key,
  1149. @BatchId,
  1150. @Now
  1151. FROM mdp_stg_ship_trans d
  1152. LEFT JOIN mdp_stg_ship_trans m ON m.source_table='ShippingPlan'
  1153. AND m.tenant_id = d.tenant_id
  1154. AND JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.RecID')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id'))
  1155. WHERE d.source_table='ShippingPlanDetail'
  1156. AND d.tenant_id=@TenantId
  1157. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1158. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr'))), '') <> ''
  1159. AND COALESCE(NULLIF(d.tenant_id, 0),
  1160. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1161. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1162. ON DUPLICATE KEY UPDATE
  1163. factory_id=VALUES(factory_id), plan_no=VALUES(plan_no), order_no=VALUES(order_no), customer_no=VALUES(customer_no), customer_name=VALUES(customer_name),
  1164. item_code=VALUES(item_code), item_name=VALUES(item_name), qty=VALUES(qty), plan_qty=VALUES(plan_qty),
  1165. plan_ship_date=VALUES(plan_ship_date), shipping_site=VALUES(shipping_site), shipping_address=VALUES(shipping_address),
  1166. status=VALUES(status), confirm_status=VALUES(confirm_status), sync_batch_id=VALUES(sync_batch_id),
  1167. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  1168. """, scope, batchId, now);
  1169. yield return Cmd(
  1170. """
  1171. INSERT INTO mdp_std_ship_trans
  1172. (tenant_id, factory_id, source_system, trans_type, shipper_rec_id, shipper_no, shipper_line, order_no, order_line,
  1173. customer_no, item_code, item_name, qty_to_ship, picking_qty, real_qty, gross_weight, net_weight, volume,
  1174. plan_ship_date, actual_ship_date, site, status, confirm_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  1175. SELECT
  1176. COALESCE(NULLIF(d.tenant_id, 0),
  1177. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1178. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1179. NULLIF(d.factory_id, 0),
  1180. 'AIDOP',
  1181. 'ASN_SHIPPER',
  1182. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ASNBOLShipperRecID')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ASNBOLShipperRecID')) AS SIGNED) END,
  1183. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Id')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Id'))),
  1184. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS CHAR),
  1185. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.OrdNbr'))),
  1186. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdLine')) AS CHAR),
  1187. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.SoldTo')),
  1188. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ContainerItem')),
  1189. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Descr')),
  1190. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyToShip')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyToShip')) AS DECIMAL(18,6)) END,
  1191. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.PickingQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.PickingQty')) AS DECIMAL(18,6)) END,
  1192. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RealQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RealQty')) AS DECIMAL(18,6)) END,
  1193. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.GrossWeight')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.GrossWeight')) AS DECIMAL(18,6)) END,
  1194. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.NetWeight')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.NetWeight')) AS DECIMAL(18,6)) END,
  1195. CASE WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Volume')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Volume'))) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Volume')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Volume'))) AS DECIMAL(18,6)) END,
  1196. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ShipDate')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShipDate'))), 'null'), ''),
  1197. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ShipDate')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShipDate'))), 'null'), ''),
  1198. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Site')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Site'))),
  1199. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Status')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Status'))),
  1200. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.IsConfirm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.IsConfirm'))) AS CHAR),
  1201. d.source_table,
  1202. d.source_row_id,
  1203. d.source_biz_key,
  1204. @BatchId,
  1205. @Now
  1206. FROM mdp_stg_ship_trans d
  1207. LEFT JOIN mdp_stg_ship_trans m ON m.source_table='ASNBOLShipperMaster'
  1208. AND m.tenant_id = d.tenant_id
  1209. AND JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.RecID')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ASNBOLShipperRecID'))
  1210. WHERE d.source_table='ASNBOLShipperDetail'
  1211. AND d.tenant_id=@TenantId
  1212. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1213. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Id')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Id'))), '') <> ''
  1214. AND COALESCE(NULLIF(d.tenant_id, 0),
  1215. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1216. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1217. ON DUPLICATE KEY UPDATE
  1218. factory_id=VALUES(factory_id), shipper_no=VALUES(shipper_no), order_no=VALUES(order_no), order_line=VALUES(order_line),
  1219. customer_no=VALUES(customer_no), item_code=VALUES(item_code), item_name=VALUES(item_name),
  1220. qty_to_ship=VALUES(qty_to_ship), picking_qty=VALUES(picking_qty), real_qty=VALUES(real_qty),
  1221. actual_ship_date=VALUES(actual_ship_date), site=VALUES(site), status=VALUES(status),
  1222. confirm_status=VALUES(confirm_status), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  1223. update_time=CURRENT_TIMESTAMP
  1224. """, scope, batchId, now);
  1225. yield return Cmd(
  1226. """
  1227. INSERT INTO mdp_std_ship_trans
  1228. (tenant_id, factory_id, source_system, trans_type, order_no, customer_no, item_code, item_name, qty,
  1229. plan_ship_date, status, linkage_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  1230. SELECT
  1231. COALESCE(NULLIF(d.tenant_id, 0),
  1232. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1233. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1234. NULLIF(d.factory_id, 0),
  1235. 'AIDOP',
  1236. 'LINKAGE_PLAN',
  1237. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')),
  1238. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.custom_no')),
  1239. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_number')),
  1240. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Descr')),
  1241. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) AS DECIMAL(18,6)) END,
  1242. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sys_capacity_date')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.fystarttime'))), 'null'), ''),
  1243. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.type')),
  1244. CASE
  1245. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.isuse')) = '1' THEN 'LINKED'
  1246. ELSE 'INACTIVE'
  1247. END,
  1248. d.source_table,
  1249. d.source_row_id,
  1250. d.source_biz_key,
  1251. @BatchId,
  1252. @Now
  1253. FROM mdp_stg_ship_trans d
  1254. WHERE d.source_table='LinkagePlan'
  1255. AND d.tenant_id=@TenantId
  1256. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1257. AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), '') <> ''
  1258. AND COALESCE(NULLIF(d.tenant_id, 0),
  1259. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1260. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1261. ON DUPLICATE KEY UPDATE
  1262. factory_id=VALUES(factory_id), order_no=VALUES(order_no), customer_no=VALUES(customer_no), item_code=VALUES(item_code),
  1263. item_name=VALUES(item_name), qty=VALUES(qty), plan_ship_date=VALUES(plan_ship_date),
  1264. status=VALUES(status), linkage_status=VALUES(linkage_status), sync_batch_id=VALUES(sync_batch_id),
  1265. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  1266. """, scope, batchId, now);
  1267. }
  1268. private IEnumerable<S1MdpSqlCommand> BuildDwdCommands(S1MdpRunScope scope, string batchId, DateTime now)
  1269. {
  1270. yield return Cmd(
  1271. """
  1272. INSERT INTO dwd_requirement_examine_detail
  1273. (tenant_id, factory_id, stat_date, row_id, parent_row_id, examine_id, order_entry_id, bill_no, morder_no,
  1274. num, item_number, item_name, bom_number, model, kitting_time, item_type, erp_cls_name, qty, wastage,
  1275. need_count, sqty, use_qty, self_lack_qty, lack_qty, mo_qty, make_qty, purchase_qty, purchase_occupy_qty,
  1276. satisfy_time, have_ic_subs, substitute_code, create_time, source_system, sync_batch_id, calc_batch_id, calc_time,
  1277. bom_level, material_role, is_current_flag)
  1278. SELECT
  1279. COALESCE(NULLIF(d.tenant_id, 0), NULLIF(r.tenant_id, 0), NULLIF(so.tenant_id, 0),
  1280. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1281. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1282. COALESCE(NULLIF(d.factory_id, 0), NULLIF(r.factory_id, 0), NULLIF(so.factory_id, 0)),
  1283. @StatDate,
  1284. CAST(d.source_row_id AS SIGNED),
  1285. CASE
  1286. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.num')) = '1' THEN NULL
  1287. ELSE p.parent_row_id
  1288. END,
  1289. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id')) REGEXP '^-?[0-9]+$'
  1290. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id')) AS SIGNED) END,
  1291. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$'
  1292. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) AS SIGNED) END,
  1293. COALESCE(so.order_no, JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.bill_no'))),
  1294. JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.morder_no')),
  1295. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.num')) AS CHAR),
  1296. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_number')),
  1297. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_name')),
  1298. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bom_number')),
  1299. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.model')),
  1300. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.kitting_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  1301. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.type')) = '1' THEN '替代件' ELSE '标准件' END,
  1302. COALESCE(
  1303. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls_name')),
  1304. CASE JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls'))
  1305. WHEN '0' THEN '配置类'
  1306. WHEN '1' THEN '自制'
  1307. WHEN '2' THEN '委外加工'
  1308. WHEN '3' THEN '外购'
  1309. WHEN '4' THEN '虚拟件'
  1310. END
  1311. ),
  1312. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) AS DECIMAL(18,4)) END,
  1313. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.wastage')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.wastage')) AS DECIMAL(18,4)) END,
  1314. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.needCount')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.needCount')) AS DECIMAL(18,4)) END,
  1315. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sqty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sqty')) AS DECIMAL(18,4)) END,
  1316. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.use_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.use_qty')) AS DECIMAL(18,4)) END,
  1317. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.self_lack_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.self_lack_qty')) AS DECIMAL(18,4)) END,
  1318. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.lack_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.lack_qty')) AS DECIMAL(18,4)) END,
  1319. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.mo_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.mo_qty')) AS DECIMAL(18,4)) END,
  1320. CASE
  1321. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) = '0' AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls')) = '1'
  1322. THEN CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.lack_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.lack_qty')) AS DECIMAL(18,4)) END
  1323. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  1324. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) AS DECIMAL(18,4))
  1325. END,
  1326. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.purchase_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.purchase_qty')) AS DECIMAL(18,4)) END,
  1327. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.purchase_occupy_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.purchase_occupy_qty')) AS DECIMAL(18,4)) END,
  1328. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.satisfy_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  1329. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.haveicsubs')) = '1' THEN '是' ELSE '否' END,
  1330. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.substitute_code')),
  1331. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.create_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  1332. 'AIDOP',
  1333. d.sync_batch_id,
  1334. @BatchId,
  1335. @Now,
  1336. -- bom_level:仅当 level 是纯数字时才物化,脏值留 NULL(不猜层级)。
  1337. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.level')) REGEXP '^[0-9]+$'
  1338. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.level')) AS SIGNED) END,
  1339. -- material_role:把「层级」翻译成「角色」,让下游(含 S8)永远不必知道 level 的存在。
  1340. -- 失败方向必须保守:只有 level 精确等于 '1' 才判成品根;NULL / 脏值 / 空串一律落 INPUT_MATERIAL,
  1341. -- 绝不产生 NULL、绝不把行丢掉——漏判一个投入料会漏报缺料,误判一个成品根只会少算一行。
  1342. -- (SQL 三值逻辑:level 为 NULL 时 `= '1'` 得 NULL,CASE 自然走 ELSE,符合上述保守方向。)
  1343. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.level')) = '1'
  1344. THEN 'FINISHED_GOOD_ROOT' ELSE 'INPUT_MATERIAL' END,
  1345. -- is_current_flag 一律先写 0,由本阶段最后一步的原子发布语句统一翻牌(见 PublishCurrentSnapshotAsync)。
  1346. 0
  1347. FROM mdp_stg_so d
  1348. INNER JOIN mdp_stg_so r ON r.source_table='b_examine_result'
  1349. AND r.tenant_id = d.tenant_id
  1350. AND r.source_row_id = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id'))
  1351. LEFT JOIN mdp_std_so so ON so.tenant_id = COALESCE(r.tenant_id, d.tenant_id)
  1352. AND COALESCE(NULLIF(so.factory_id, 0), 1)=@FactoryId
  1353. AND so.order_entry_id = CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$'
  1354. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) AS SIGNED) END
  1355. LEFT JOIN (
  1356. SELECT examine_id, MIN(source_row_id) AS parent_row_id
  1357. FROM (
  1358. SELECT JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.examine_id')) AS examine_id, CAST(source_row_id AS SIGNED) AS source_row_id
  1359. FROM mdp_stg_so
  1360. WHERE source_table='b_bom_child_examine'
  1361. AND tenant_id=@TenantId
  1362. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1363. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.num')) = '1'
  1364. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.is_use')) IN ('1', 'true', 'True', 'base64:type16:AQ==')
  1365. AND sync_batch_id = (SELECT lb.sync_batch_id FROM mdp_stg_so lb
  1366. WHERE lb.source_table='b_bom_child_examine' AND lb.tenant_id=@TenantId
  1367. ORDER BY lb.sync_time DESC, lb.id DESC LIMIT 1)
  1368. ) x
  1369. GROUP BY examine_id
  1370. ) p ON p.examine_id = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id'))
  1371. WHERE d.source_table='b_bom_child_examine'
  1372. AND d.tenant_id=@TenantId
  1373. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1374. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.is_use')) IN ('1', 'true', 'True', 'base64:type16:AQ==')
  1375. -- ── B5 表头软删过滤 ────────────────────────────────────────────────────────────
  1376. -- b_examine_result 是软删表:核验单被作废后仍留在源里,只把 IsDeleted 置真。
  1377. -- 之前这里完全没有版本过滤,等于把已作废核验单的缺料当成现行缺料输出。
  1378. -- 实测(aidopdev,只读):贴源 5969 行表头里 5851 行 IsDeleted 为真,仅 118 行存活
  1379. -- ——即旧 DWD 有约 98% 建立在已作废表头上。
  1380. -- bit(1) 列经贴源后有两种编码并存('1'/'0' 与 base64:type16:AQ==/AA==),两种都必须认。
  1381. -- 这里用「存活白名单」而不是「已删黑名单」,与上面 is_use 的写法保持同一方向;
  1382. -- 实测该列在贴源层恰好只有这 4 种取值、无 NULL、无其它编码,故白名单当前不会误杀任何行。
  1383. AND JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.IsDeleted')) IN ('0', 'false', 'False', 'base64:type16:AA==')
  1384. -- ── B6 FULL=镜像 的 DWD 侧权宜实现(STOPGAP)────────────────────────────────────
  1385. -- 正解应当是贴源层按 (tenant_id, factory_id, source_table) 整体替换,但贴源写入由
  1386. -- DataPlatform/Executors/MdpStagingWriter.cs + MdpDbPullExecutor.cs 通用实现,
  1387. -- 且 mdp_stg_so 的唯一键 uk_source_key=(source_system, source_table, source_biz_key)
  1388. -- 连 tenant_id 都不含 —— 改成整体替换必须同时改那两个文件和该唯一键,均不在本次可改文件范围内。
  1389. -- 故在 DWD 侧兜底:只取每个 (tenant_id, source_table) 的「最新一次 sync_batch_id」。
  1390. -- 该等价关系成立的前提是 sync_mode='FULL'(实测两个 entity 均为 FULL):
  1391. -- FULL 每轮重读整张源表,仍存在的行其 sync_batch_id 必被刷新到最新批次;
  1392. -- 留着旧 sync_batch_id 的行 = 上一轮全量里已经查不到 = 源侧已删除的孤儿。
  1393. -- 实测孤儿量很大:贴源 b_bom_child_examine 共 32105 行,最新批次只有 15272 行。
  1394. -- ⚠️ 若哪天该 entity 被改成 INCREMENTAL,本过滤会错杀未变更的存量行,届时必须回到贴源层正解。
  1395. AND d.sync_batch_id = (SELECT lb.sync_batch_id FROM mdp_stg_so lb
  1396. WHERE lb.source_table='b_bom_child_examine' AND lb.tenant_id=@TenantId
  1397. ORDER BY lb.sync_time DESC, lb.id DESC LIMIT 1)
  1398. AND r.sync_batch_id = (SELECT lb.sync_batch_id FROM mdp_stg_so lb
  1399. WHERE lb.source_table='b_examine_result' AND lb.tenant_id=@TenantId
  1400. ORDER BY lb.sync_time DESC, lb.id DESC LIMIT 1)
  1401. AND COALESCE(NULLIF(d.tenant_id, 0), NULLIF(r.tenant_id, 0), NULLIF(so.tenant_id, 0),
  1402. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1403. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1404. ON DUPLICATE KEY UPDATE
  1405. parent_row_id=VALUES(parent_row_id), order_entry_id=VALUES(order_entry_id), bill_no=VALUES(bill_no),
  1406. morder_no=VALUES(morder_no), num=VALUES(num), item_number=VALUES(item_number), item_name=VALUES(item_name),
  1407. bom_number=VALUES(bom_number), model=VALUES(model), kitting_time=VALUES(kitting_time),
  1408. item_type=VALUES(item_type), erp_cls_name=VALUES(erp_cls_name), qty=VALUES(qty), wastage=VALUES(wastage),
  1409. need_count=VALUES(need_count), sqty=VALUES(sqty), use_qty=VALUES(use_qty), self_lack_qty=VALUES(self_lack_qty),
  1410. lack_qty=VALUES(lack_qty), mo_qty=VALUES(mo_qty), make_qty=VALUES(make_qty), purchase_qty=VALUES(purchase_qty),
  1411. purchase_occupy_qty=VALUES(purchase_occupy_qty), satisfy_time=VALUES(satisfy_time), have_ic_subs=VALUES(have_ic_subs),
  1412. substitute_code=VALUES(substitute_code), create_time=VALUES(create_time), sync_batch_id=VALUES(sync_batch_id),
  1413. calc_time=VALUES(calc_time), update_time=CURRENT_TIMESTAMP,
  1414. bom_level=VALUES(bom_level), material_role=VALUES(material_role)
  1415. -- 故意不写 is_current_flag:该列的唯一写者是本阶段最后一步的原子发布语句。
  1416. -- 若在这里也回写,同批次重跑会把已发布状态打回 0,制造瞬时「无当前快照」窗口。
  1417. """, scope, batchId, now);
  1418. yield return Cmd(
  1419. """
  1420. INSERT INTO dwd_ship_trans
  1421. (tenant_id, factory_id, company_id, stat_date, order_id, order_entry_id, order_no, order_line, customer_no,
  1422. customer_name, country, item_code, item_name, item_spec, order_qty, planned_ship_qty, shipped_qty,
  1423. remaining_qty, order_date, customer_request_date, plan_delivery_date, promised_delivery_date,
  1424. plan_ship_date, actual_ship_date, review_status, order_status, delivery_status, linkage_status, risk_level,
  1425. source_system, source_table, source_row_id, source_biz_key, sync_batch_id, calc_batch_id, calc_time)
  1426. SELECT
  1427. so.tenant_id,
  1428. so.factory_id,
  1429. so.company_id,
  1430. @StatDate,
  1431. so.order_id,
  1432. so.order_entry_id,
  1433. so.order_no,
  1434. IFNULL(so.order_line, ''),
  1435. so.customer_no,
  1436. so.customer_name,
  1437. so.country,
  1438. IFNULL(so.item_code, ''),
  1439. so.item_name,
  1440. so.item_spec,
  1441. IFNULL(so.order_qty, 0),
  1442. IFNULL(p.plan_qty, 0),
  1443. IFNULL(a.real_qty, 0),
  1444. GREATEST(IFNULL(so.order_qty, 0) - IFNULL(a.real_qty, 0), 0),
  1445. so.order_date,
  1446. so.customer_request_date,
  1447. so.plan_delivery_date,
  1448. so.promised_delivery_date,
  1449. p.plan_ship_date,
  1450. a.actual_ship_date,
  1451. so.review_status,
  1452. so.order_status,
  1453. CASE
  1454. WHEN IFNULL(so.order_qty, 0) > 0 AND IFNULL(a.real_qty, 0) >= IFNULL(so.order_qty, 0) THEN 'COMPLETED'
  1455. WHEN COALESCE(p.plan_ship_date, so.promised_delivery_date, so.plan_delivery_date) < @Now THEN 'DELAYED'
  1456. WHEN IFNULL(p.plan_qty, 0) > 0 THEN 'PLANNED'
  1457. ELSE 'OPEN'
  1458. END,
  1459. l.linkage_status,
  1460. CASE
  1461. WHEN COALESCE(p.plan_ship_date, so.promised_delivery_date, so.plan_delivery_date) < @Now
  1462. AND IFNULL(a.real_qty, 0) < IFNULL(so.order_qty, 0) THEN 'HIGH'
  1463. WHEN IFNULL(a.real_qty, 0) < IFNULL(so.order_qty, 0) THEN 'MEDIUM'
  1464. ELSE 'LOW'
  1465. END,
  1466. 'AIDOP',
  1467. so.source_table,
  1468. so.source_row_id,
  1469. so.source_biz_key,
  1470. so.sync_batch_id,
  1471. @BatchId,
  1472. @Now
  1473. FROM mdp_std_so so
  1474. LEFT JOIN (
  1475. SELECT tenant_id, order_no, order_entry_id, IFNULL(order_line, '') AS order_line, IFNULL(item_code, '') AS item_code,
  1476. SUM(IFNULL(plan_qty, IFNULL(qty, 0))) AS plan_qty,
  1477. MIN(plan_ship_date) AS plan_ship_date
  1478. FROM mdp_std_ship_trans
  1479. WHERE trans_type='SHIP_PLAN'
  1480. AND tenant_id=@TenantId
  1481. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1482. GROUP BY tenant_id, order_no, order_entry_id, IFNULL(order_line, ''), IFNULL(item_code, '')
  1483. ) p ON so.tenant_id=p.tenant_id
  1484. AND so.order_no=p.order_no
  1485. AND IFNULL(so.item_code, '')=p.item_code
  1486. AND (
  1487. (p.order_entry_id IS NOT NULL AND so.order_entry_id=p.order_entry_id)
  1488. OR (p.order_entry_id IS NULL AND IFNULL(so.order_line, '')=p.order_line)
  1489. )
  1490. LEFT JOIN (
  1491. SELECT tenant_id, order_no, IFNULL(order_line, '') AS order_line, IFNULL(item_code, '') AS item_code,
  1492. SUM(IFNULL(real_qty, IFNULL(qty_to_ship, 0))) AS real_qty,
  1493. MAX(actual_ship_date) AS actual_ship_date
  1494. FROM mdp_std_ship_trans
  1495. WHERE trans_type='ASN_SHIPPER'
  1496. AND tenant_id=@TenantId
  1497. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1498. GROUP BY tenant_id, order_no, IFNULL(order_line, ''), IFNULL(item_code, '')
  1499. ) a ON so.tenant_id=a.tenant_id AND so.order_no=a.order_no AND IFNULL(so.order_line, '')=a.order_line AND IFNULL(so.item_code, '')=a.item_code
  1500. LEFT JOIN (
  1501. SELECT tenant_id, order_no, item_code, MAX(linkage_status) AS linkage_status
  1502. FROM mdp_std_ship_trans
  1503. WHERE IFNULL(linkage_status, '') <> ''
  1504. AND tenant_id=@TenantId
  1505. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1506. GROUP BY tenant_id, order_no, item_code
  1507. ) l ON so.tenant_id=l.tenant_id AND so.order_no=l.order_no AND IFNULL(so.item_code, '')=IFNULL(l.item_code, '')
  1508. WHERE so.tenant_id=@TenantId
  1509. AND COALESCE(NULLIF(so.factory_id, 0), 1)=@FactoryId
  1510. AND IFNULL(so.order_no, '') <> ''
  1511. ON DUPLICATE KEY UPDATE
  1512. factory_id=VALUES(factory_id), customer_no=VALUES(customer_no), customer_name=VALUES(customer_name), item_name=VALUES(item_name),
  1513. order_qty=VALUES(order_qty), planned_ship_qty=VALUES(planned_ship_qty), shipped_qty=VALUES(shipped_qty),
  1514. remaining_qty=VALUES(remaining_qty), plan_ship_date=VALUES(plan_ship_date), actual_ship_date=VALUES(actual_ship_date),
  1515. review_status=VALUES(review_status), order_status=VALUES(order_status), delivery_status=VALUES(delivery_status),
  1516. linkage_status=VALUES(linkage_status), risk_level=VALUES(risk_level), sync_batch_id=VALUES(sync_batch_id),
  1517. calc_batch_id=VALUES(calc_batch_id), calc_time=VALUES(calc_time), update_time=CURRENT_TIMESTAMP
  1518. """, scope, batchId, now);
  1519. }
  1520. private async Task<long> InsertSyncLogAsync(long tenantId, long entityId, string entityName, string batchId, int rowsRead)
  1521. {
  1522. await _db.Ado.ExecuteCommandAsync(
  1523. """
  1524. INSERT INTO mdp_sync_log
  1525. (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type, sync_start, rows_read, status)
  1526. VALUES (@TenantId, @EntityId, 'AIDOPDEV_MYSQL', @EntityName, @BatchId, 'FULL', 'AUTO', NOW(), @RowsRead, 'RUNNING')
  1527. """,
  1528. new SugarParameter("@TenantId", tenantId),
  1529. new SugarParameter("@EntityId", entityId),
  1530. new SugarParameter("@EntityName", entityName),
  1531. new SugarParameter("@BatchId", batchId),
  1532. new SugarParameter("@RowsRead", rowsRead));
  1533. return await _db.Ado.GetLongAsync(
  1534. "SELECT id FROM mdp_sync_log WHERE sync_batch_id=@BatchId AND entity_id=@EntityId ORDER BY id DESC LIMIT 1",
  1535. new List<SugarParameter>
  1536. {
  1537. new("@BatchId", batchId),
  1538. new("@EntityId", entityId)
  1539. });
  1540. }
  1541. private async Task MarkSyncLogSuccessAsync(long logId, DateTime started, int affected)
  1542. {
  1543. await _db.Ado.ExecuteCommandAsync(
  1544. """
  1545. UPDATE mdp_sync_log
  1546. SET sync_end=NOW(), duration_ms=@DurationMs, rows_insert=@RowsInsert, rows_update=0, rows_skip=0, rows_error=0, status='SUCCESS'
  1547. WHERE id=@Id
  1548. """,
  1549. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  1550. new SugarParameter("@RowsInsert", affected),
  1551. new SugarParameter("@Id", logId));
  1552. }
  1553. private async Task MarkSyncLogFailedAsync(long logId, DateTime started, string message)
  1554. {
  1555. try
  1556. {
  1557. await _db.Ado.ExecuteCommandAsync(
  1558. """
  1559. UPDATE mdp_sync_log
  1560. SET sync_end=NOW(), duration_ms=@DurationMs, rows_error=1, status='FAILED', error_msg=@ErrorMsg
  1561. WHERE id=@Id
  1562. """,
  1563. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  1564. new SugarParameter("@ErrorMsg", Truncate(message, 1000)),
  1565. new SugarParameter("@Id", logId));
  1566. }
  1567. catch (Exception ex)
  1568. {
  1569. // 写库自身失败兜底:避免再抛掩盖原异常;遗留 RUNNING 行可由运维手动清理
  1570. Console.Error.WriteLine($"[S1MdpSyncTransform] MarkSyncLogFailed write failed (syncLogId={logId}): {ex.Message}");
  1571. }
  1572. }
  1573. private async Task<long> InsertTransformRunLogAsync(S1MdpRunScope scope, string batchId, DateTime startedAt, string triggerType)
  1574. {
  1575. await _db.Ado.ExecuteCommandAsync(
  1576. """
  1577. INSERT INTO mdp_transform_run_log
  1578. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  1579. VALUES (@TenantId, @JobCode, 'S1 MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  1580. """,
  1581. new SugarParameter("@TenantId", scope.TenantId),
  1582. new SugarParameter("@JobCode", JobCode),
  1583. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  1584. new SugarParameter("@BatchId", batchId),
  1585. new SugarParameter("@StartTime", startedAt));
  1586. return await _db.Ado.GetLongAsync(
  1587. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  1588. new List<SugarParameter> { new("@BatchId", batchId) });
  1589. }
  1590. private async Task MarkTransformRunSuccessAsync(long runLogId, DateTime startedAt, S1MdpSyncTransformResult result)
  1591. {
  1592. var finishedAt = DateTime.Now;
  1593. await _db.Ado.ExecuteCommandAsync(
  1594. """
  1595. UPDATE mdp_transform_run_log
  1596. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  1597. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=@DwdRows,
  1598. summary_json=@SummaryJson, update_time=CURRENT_TIMESTAMP
  1599. WHERE id=@Id
  1600. """,
  1601. new SugarParameter("@EndTime", finishedAt),
  1602. new SugarParameter("@DurationMs", ToDurationMs(finishedAt - startedAt)),
  1603. new SugarParameter("@StageRows", result.StageRows),
  1604. new SugarParameter("@StandardRows", result.StandardRows),
  1605. new SugarParameter("@DwdRows", result.DwdRows),
  1606. new SugarParameter("@SummaryJson", BuildRunSummaryJson(result)),
  1607. new SugarParameter("@Id", runLogId));
  1608. }
  1609. private async Task MarkTransformRunFailedAsync(long runLogId, DateTime startedAt, string message)
  1610. {
  1611. try
  1612. {
  1613. var finishedAt = DateTime.Now;
  1614. var db = _db.CopyNew();
  1615. await db.Ado.ExecuteCommandAsync(
  1616. """
  1617. UPDATE mdp_transform_run_log
  1618. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  1619. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  1620. WHERE id=@Id
  1621. """,
  1622. new SugarParameter("@EndTime", finishedAt),
  1623. new SugarParameter("@DurationMs", ToDurationMs(finishedAt - startedAt)),
  1624. new SugarParameter("@ErrorMessage", Truncate(message, 2000)),
  1625. new SugarParameter("@Id", runLogId));
  1626. }
  1627. catch (Exception ex)
  1628. {
  1629. // 写库自身失败兜底(典型场景:远端 MySQL 瞬断导致 MarkFailed 自身也连不上):
  1630. // 避免再抛二次异常掩盖原错;遗留 RUNNING 行可由运维手动清理。
  1631. Console.Error.WriteLine($"[S1MdpSyncTransform] MarkTransformRunFailed write failed (runLogId={runLogId}): {ex.Message}");
  1632. }
  1633. }
  1634. private static S1MdpSqlCommand Cmd(string sql, S1MdpRunScope scope, string batchId, DateTime now)
  1635. {
  1636. return new S1MdpSqlCommand(sql, new[]
  1637. {
  1638. new SugarParameter("@BatchId", batchId),
  1639. new SugarParameter("@Now", now),
  1640. new SugarParameter("@StatDate", now.Date),
  1641. new SugarParameter("@TenantId", scope.TenantId),
  1642. new SugarParameter("@FactoryId", scope.FactoryId)
  1643. });
  1644. }
  1645. private static string BuildJsonObjectExpression(IEnumerable<string> columns)
  1646. {
  1647. var parts = columns.SelectMany(c => new[] { $"'{c.Replace("'", "''")}'", $"s.`{c}`" });
  1648. return $"JSON_OBJECT({string.Join(",", parts)})";
  1649. }
  1650. private static string BuildOptionalColumnExpr(IReadOnlyCollection<string> columns, string expected, string fallback)
  1651. {
  1652. return columns.Any(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase))
  1653. ? $"s.`{FindColumn(columns, expected)}`"
  1654. : fallback;
  1655. }
  1656. private static string FindColumn(IEnumerable<string> columns, string expected)
  1657. {
  1658. return columns.First(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase));
  1659. }
  1660. private static async Task KeepLeaseAliveAsync(IS1MdpFullRunLease lease, CancellationToken ct)
  1661. {
  1662. while (!ct.IsCancellationRequested)
  1663. {
  1664. try
  1665. {
  1666. await Task.Delay(TimeSpan.FromSeconds(30), ct);
  1667. await lease.HeartbeatAsync(CancellationToken.None);
  1668. }
  1669. catch (OperationCanceledException)
  1670. {
  1671. return;
  1672. }
  1673. }
  1674. }
  1675. private static int ToDurationMs(TimeSpan elapsed)
  1676. {
  1677. var ms = elapsed.TotalMilliseconds;
  1678. if (double.IsNaN(ms) || ms <= 0) return 0;
  1679. return ms >= int.MaxValue ? int.MaxValue : (int)ms;
  1680. }
  1681. private static string NormalizeTriggerType(string? triggerType)
  1682. {
  1683. return string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  1684. }
  1685. /// <summary>
  1686. /// 运行摘要。改用 JsonSerializer 而非字符串插值:batchId 今天不可能含引号,
  1687. /// 但拼 JSON 是会被后来者继承的脆弱写法,且本次要加的 publish 是嵌套对象。
  1688. /// </summary>
  1689. private static string BuildRunSummaryJson(S1MdpSyncTransformResult result)
  1690. {
  1691. return System.Text.Json.JsonSerializer.Serialize(new
  1692. {
  1693. batchId = result.BatchId,
  1694. stageRows = result.StageRows,
  1695. standardRows = result.StandardRows,
  1696. dwdRows = result.DwdRows,
  1697. kpiRows = result.KpiRows,
  1698. publish = result.Publication is null
  1699. ? null
  1700. : new
  1701. {
  1702. table = result.Publication.Table,
  1703. batchId = result.Publication.BatchId,
  1704. factoryId = result.Publication.FactoryId,
  1705. currentRows = result.Publication.CurrentRows
  1706. }
  1707. });
  1708. }
  1709. private static string ResolveKpiValueTable(int metricLevel)
  1710. {
  1711. return metricLevel switch
  1712. {
  1713. 1 => "ado_s9_kpi_value_l1_day",
  1714. 2 => "ado_s9_kpi_value_l2_day",
  1715. 3 => "ado_s9_kpi_value_l3_day",
  1716. 4 => "ado_s9_kpi_value_l4_day",
  1717. _ => "ado_s9_kpi_value_l2_day"
  1718. };
  1719. }
  1720. private static decimal DefaultS1Target(string metricCode) => LegacyKpiCodeTargets.GetOrZero(metricCode);
  1721. private static string ResolveKpiStatus(decimal actual, decimal target, string? direction, decimal? yellowThreshold, decimal? redThreshold)
  1722. {
  1723. if (target <= 0) return "gray";
  1724. var ratio = actual / target * 100m;
  1725. if (string.Equals(direction, "lower_is_better", StringComparison.OrdinalIgnoreCase))
  1726. {
  1727. if (actual <= target) return "green";
  1728. if (ratio <= (yellowThreshold ?? 110m)) return "yellow";
  1729. return ratio >= (redThreshold ?? 120m) ? "red" : "yellow";
  1730. }
  1731. if (actual >= target) return "green";
  1732. if (ratio >= (yellowThreshold ?? 95m)) return "yellow";
  1733. return ratio <= (redThreshold ?? 80m) ? "red" : "yellow";
  1734. }
  1735. private static string ResolveTrendFlag(decimal actual, decimal? previous)
  1736. {
  1737. if (previous == null) return "flat";
  1738. if (actual > previous.Value) return "up";
  1739. if (actual < previous.Value) return "down";
  1740. return "flat";
  1741. }
  1742. private static string Truncate(string? raw, int maxLength)
  1743. {
  1744. if (string.IsNullOrEmpty(raw)) return string.Empty;
  1745. return raw.Length <= maxLength ? raw : raw[..maxLength];
  1746. }
  1747. private sealed class S1ColumnRow
  1748. {
  1749. public string ColumnName { get; set; } = string.Empty;
  1750. }
  1751. private sealed class S1MdpEntityRow
  1752. {
  1753. public long Id { get; set; }
  1754. public string EntityName { get; set; } = string.Empty;
  1755. }
  1756. private sealed class S1KpiCalcRow
  1757. {
  1758. public long TenantId { get; set; }
  1759. public long FactoryId { get; set; }
  1760. public string MetricCode { get; set; } = string.Empty;
  1761. public decimal? MetricValue { get; set; }
  1762. }
  1763. private sealed class S1KpiMetaRow
  1764. {
  1765. public int MetricLevel { get; set; }
  1766. public string Direction { get; set; } = "higher_is_better";
  1767. public decimal? YellowThreshold { get; set; }
  1768. public decimal? RedThreshold { get; set; }
  1769. }
  1770. private sealed class S1KpiValueRow
  1771. {
  1772. public long Id { get; set; }
  1773. public decimal? MetricValue { get; set; }
  1774. public decimal? TargetValue { get; set; }
  1775. }
  1776. }
  1777. public sealed class S1MdpSyncTransformResult
  1778. {
  1779. public long RunLogId { get; set; }
  1780. public string BatchId { get; set; } = string.Empty;
  1781. public int StageRows { get; set; }
  1782. public int StandardRows { get; set; }
  1783. public int DwdRows { get; set; }
  1784. public int KpiRows { get; set; }
  1785. public int AtomicRows { get; set; }
  1786. /// <summary>
  1787. /// 当前快照发布证据。为 null 表示本轮没走到发布这一步。
  1788. ///
  1789. /// <para><b>为什么需要它</b>:发布语句在本轮零行时命中 0 行,在库里不留任何痕迹,
  1790. /// 于是「源侧本轮确实没有数据」与「发布压根没跑成」在 DWD 里字面等同。
  1791. /// 下游的 Authority 健康判定因此只能保守地判 UNKNOWN、拦下恢复 ——
  1792. /// 哪怕这是一个完全正常的空快照。</para>
  1793. ///
  1794. /// <para>把发布结果记进运行日志,就把这个隐式事实变成了可观测事实:
  1795. /// 「本轮成功发布了批次 X,其中当前行数为 N(N 可以是 0)」。</para>
  1796. /// </summary>
  1797. public S1MdpSnapshotPublication? Publication { get; set; }
  1798. }
  1799. /// <summary>
  1800. /// 一次当前快照发布的结果。与 <c>status='SUCCESS'</c> 写在同一条 UPDATE 里,
  1801. /// 因此不存在「成功了但没有发布证据」的中间态。
  1802. /// </summary>
  1803. public sealed class S1MdpSnapshotPublication
  1804. {
  1805. /// <summary>被发布为当前快照的表。</summary>
  1806. public string Table { get; set; } = string.Empty;
  1807. /// <summary>本轮批次号。读者据此确认「当前快照」出自哪一次运行。</summary>
  1808. public string BatchId { get; set; } = string.Empty;
  1809. /// <summary>发布作用域的工厂号。运行日志表本身没有工厂列,只能由这里承载。</summary>
  1810. public long FactoryId { get; set; }
  1811. /// <summary>发布后该批次的当前行数。<b>0 是合法值</b>,表示成功发布了一个空快照。</summary>
  1812. public int CurrentRows { get; set; }
  1813. }
  1814. internal sealed record S1MdpSqlCommand(string Sql, SugarParameter[] Parameters);
  1815. internal sealed record S1MdpEntityConfig(
  1816. string EntityCode,
  1817. string SourceTable,
  1818. string TargetTable,
  1819. string SourceRowIdExpression,
  1820. string SourceBizKeyExpression)
  1821. {
  1822. public static readonly IReadOnlyList<S1MdpEntityConfig> All = new List<S1MdpEntityConfig>
  1823. {
  1824. new("S1_SEORDER", "crm_seorder", "mdp_stg_so", "Id", "COALESCE(s.`bill_no`, CAST(s.`Id` AS CHAR))"),
  1825. new("S1_SEORDER_ENTRY", "crm_seorderentry", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`entry_seq`, CAST(s.`Id` AS CHAR)))"),
  1826. new("S1_SEORDER_CHANGE", "crm_seorder_change", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', CAST(s.`Id` AS CHAR))"),
  1827. new("S1_CONTRACT_REVIEW", "ado_contract_review", "mdp_stg_so", "RecID", "COALESCE(s.`BillNo`, CAST(s.`RecID` AS CHAR))"),
  1828. new("S1_CONTRACT_REVIEW_FLOW", "ado_contract_review_flow", "mdp_stg_so", "RecID", "CONCAT(IFNULL(s.`ReviewBillNo`,''), ':', IFNULL(s.`StageNo`,''), ':', CAST(s.`RecID` AS CHAR))"),
  1829. new("S1_PRODUCT_DESIGN", "ado_product_design", "mdp_stg_so", "Id", "COALESCE(s.`BillNo`, CAST(s.`Id` AS CHAR))"),
  1830. new("S1_PRODUCT_DESIGN_BOM", "ado_product_design_bom", "mdp_stg_so", "Id", "CONCAT(CAST(s.`ProductDesignId` AS CHAR), ':', CAST(s.`Id` AS CHAR))"),
  1831. new("S1_PRODUCT_DESIGN_ROUTING", "ado_product_design_routing", "mdp_stg_so", "Id", "CONCAT(CAST(s.`ProductDesignId` AS CHAR), ':', CAST(s.`Id` AS CHAR))"),
  1832. new("S1_REQUIREMENT_EXAMINE_RESULT", "b_examine_result", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`morder_no`,''), ':', CAST(s.`Id` AS CHAR))"),
  1833. new("S1_REQUIREMENT_EXAMINE_DETAIL", "b_bom_child_examine", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`examine_id`,''), ':', IFNULL(s.`item_number`,''), ':', CAST(s.`Id` AS CHAR))"),
  1834. new("S1_SHIPPING_PLAN", "ShippingPlan", "mdp_stg_ship_trans", "RecID", "COALESCE(s.`LotSerial`, CAST(s.`RecID` AS CHAR))"),
  1835. new("S1_SHIPPING_PLAN_DETAIL", "ShippingPlanDetail", "mdp_stg_ship_trans", "RecID", "CONCAT(IFNULL(s.`plan_id`,''), ':', IFNULL(s.`OrdNbr`,''), ':', CAST(s.`RecID` AS CHAR))"),
  1836. new("S1_ASN_SHIPPER_MASTER", "ASNBOLShipperMaster", "mdp_stg_ship_trans", "RecID", "COALESCE(s.`Id`, CONCAT(IFNULL(s.`OrdNbr`,''), ':', CAST(s.`RecID` AS CHAR)))"),
  1837. new("S1_ASN_SHIPPER_DETAIL", "ASNBOLShipperDetail", "mdp_stg_ship_trans", "RecID", "CONCAT(IFNULL(s.`Id`,''), ':', IFNULL(s.`Line`, CAST(s.`RecID` AS CHAR)))"),
  1838. new("S1_LINKAGE_PLAN", "LinkagePlan", "mdp_stg_ship_trans", "id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`item_number`,''), ':', CAST(s.`id` AS CHAR))")
  1839. };
  1840. }