S1MdpSyncTransformService.cs 118 KB

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