S1MdpSyncTransformService.cs 122 KB

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