S1MdpSyncTransformService.cs 97 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591
  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. need_count DECIMAL(18,4) NULL,
  150. sqty DECIMAL(18,4) NULL,
  151. use_qty DECIMAL(18,4) NULL,
  152. self_lack_qty DECIMAL(18,4) NULL,
  153. lack_qty DECIMAL(18,4) NULL,
  154. mo_qty DECIMAL(18,4) NULL,
  155. make_qty DECIMAL(18,4) NULL,
  156. purchase_qty DECIMAL(18,4) NULL,
  157. purchase_occupy_qty DECIMAL(18,4) NULL,
  158. satisfy_time DATE NULL,
  159. have_ic_subs VARCHAR(10) NULL,
  160. substitute_code VARCHAR(100) NULL,
  161. create_time DATETIME NULL,
  162. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  163. sync_batch_id VARCHAR(100) NOT NULL,
  164. calc_batch_id VARCHAR(100) NOT NULL,
  165. calc_time DATETIME NOT NULL,
  166. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  167. UNIQUE KEY uk_dwd_req_exam_detail (tenant_id, row_id, calc_batch_id),
  168. KEY idx_req_exam_tenant_batch (tenant_id, calc_batch_id),
  169. KEY idx_req_exam_bill (tenant_id, bill_no),
  170. KEY idx_req_exam_item (tenant_id, item_number)
  171. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S1需求明细核验DWD';
  172. """);
  173. await _db.Ado.ExecuteCommandAsync(
  174. """
  175. INSERT INTO mdp_entity
  176. (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)
  177. SELECT 0, s.id, 'S1_REQUIREMENT_EXAMINE_RESULT', 'S1需求核验结果主表', 'TABLE',
  178. 'b_examine_result', 'mdp_stg_so', 'FULL', 5000, 1, '需求核验结果主表,进入 S1 贴源层', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  179. FROM mdp_source s
  180. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  181. LIMIT 1
  182. ON DUPLICATE KEY UPDATE
  183. source_id=VALUES(source_id), entity_name=VALUES(entity_name), source_table_name=VALUES(source_table_name),
  184. target_table_name=VALUES(target_table_name), sync_mode=VALUES(sync_mode), status=VALUES(status),
  185. remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  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_DETAIL', 'S1需求核验BOM明细', 'TABLE',
  192. 'b_bom_child_examine', 'mdp_stg_so', 'FULL', 5000, 1, '需求核验BOM明细,进入 S1 贴源层和 DWD', 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, incr_column, batch_size, status, remark, create_time, update_time)
  205. 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
  206. FROM mdp_source s
  207. JOIN (
  208. 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
  209. UNION ALL SELECT 'S1_PRODUCT_DESIGN_BOM', 'S1产品设计BOM', 'ado_product_design_bom', 'FULL', NULL, '产品设计BOM,表无时间戳列,只能全量抽取'
  210. UNION ALL SELECT 'S1_PRODUCT_DESIGN_ROUTING', 'S1产品设计工艺路线', 'ado_product_design_routing', 'FULL', NULL, '产品设计工艺路线,表无时间戳列,只能全量抽取'
  211. ) v
  212. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  213. ON DUPLICATE KEY UPDATE
  214. source_id=VALUES(source_id), entity_name=VALUES(entity_name), source_table_name=VALUES(source_table_name),
  215. target_table_name=VALUES(target_table_name), sync_mode=VALUES(sync_mode), incr_column=VALUES(incr_column),
  216. status=VALUES(status), remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  217. """);
  218. }
  219. /// <summary>Phase2/3:源→stg 走执行器(DB/API);遗留 SyncOneEntityAsync 不再调用。</summary>
  220. private async Task<int> SyncStagingAsync(S1MdpRunScope scope, string batchId, DateTime now, CancellationToken cancellationToken)
  221. {
  222. _ = now;
  223. var codes = S1MdpEntityConfig.All.Select(x => x.EntityCode);
  224. return await _stagingPuller.PullEntitiesAsync(
  225. codes, batchId, scope.TenantId, fullRefresh: true, taskCode: "S1_MDP_INBOUND",
  226. cancellationToken, scope.FactoryId, requireMatchingSourceTenant: true);
  227. }
  228. /// <summary>
  229. /// 双模式入站:仅抽 stg(可指定 DB 实体码或 *_API),再跑 std/DWD/KPI。
  230. /// </summary>
  231. public async Task<S1MdpSyncTransformResult> RunInboundAsync(
  232. S1MdpRunScope scope,
  233. IEnumerable<string>? entityCodes = null,
  234. bool fullRefresh = false,
  235. CancellationToken cancellationToken = default)
  236. {
  237. scope = S1MdpRunScope.Create(scope.TenantId, scope.FactoryId);
  238. cancellationToken.ThrowIfCancellationRequested();
  239. var now = DateTime.Now;
  240. var batchId = $"S1_MDP_IN_{scope.TenantId}_{scope.FactoryId}_{now:yyyyMMddHHmmss}";
  241. var runLogId = await InsertTransformRunLogAsync(scope, batchId, now, "INBOUND");
  242. var result = new S1MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  243. try
  244. {
  245. await EnsureS1RuntimeObjectsAsync();
  246. var codes = entityCodes?.ToList() ?? S1MdpEntityConfig.All.Select(x => x.EntityCode).ToList();
  247. result.StageRows = await _stagingPuller.PullEntitiesAsync(
  248. codes, batchId, scope.TenantId, fullRefresh, "S1_MDP_INBOUND", cancellationToken,
  249. scope.FactoryId, requireMatchingSourceTenant: true);
  250. result.StandardRows = await TransformStandardAsync(scope, batchId, now, cancellationToken);
  251. result.DwdRows = await BuildDwdAsync(scope, batchId, now, cancellationToken);
  252. result.KpiRows = await BuildS1KpiValuesAsync(scope, batchId, now, cancellationToken);
  253. result.AtomicRows = await _atomicBuild.BuildOrderDeliveryDomainForAllDatesAsync(
  254. scope.TenantId, scope.FactoryId, batchId, cancellationToken);
  255. await MarkTransformRunSuccessAsync(runLogId, now, result);
  256. return result;
  257. }
  258. catch (Exception ex)
  259. {
  260. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  261. throw;
  262. }
  263. }
  264. /// <summary>Phase3:已停用。源→stg 改由执行器路径;保留方法供紧急回滚对照,勿再调用。</summary>
  265. [Obsolete("Phase3: use MdpModuleStagingPuller / executors")]
  266. private async Task<int> SyncOneEntityAsync(S1MdpEntityConfig entity, string batchId, DateTime now, long tenantId)
  267. {
  268. if (tenantId <= 0) throw new InvalidOperationException("S1 同步日志必须指定有效 tenantId");
  269. var entityRow = await _db.Ado.SqlQuerySingleAsync<S1MdpEntityRow>(
  270. "SELECT id AS Id, entity_name AS EntityName FROM mdp_entity WHERE tenant_id=0 AND entity_code=@EntityCode LIMIT 1",
  271. new SugarParameter("@EntityCode", entity.EntityCode));
  272. if (entityRow == null) throw Oops.Oh($"未找到 MDP 实体配置:{entity.EntityCode}");
  273. var columns = await _db.Ado.SqlQueryAsync<S1ColumnRow>(
  274. """
  275. SELECT COLUMN_NAME AS ColumnName
  276. FROM information_schema.COLUMNS
  277. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@TableName
  278. ORDER BY ORDINAL_POSITION
  279. """,
  280. new SugarParameter("@TableName", entity.SourceTable));
  281. if (columns.Count == 0) throw Oops.Oh($"未找到源表:{entity.SourceTable}");
  282. var names = columns.Select(u => u.ColumnName).ToList();
  283. var tenantExpr = BuildOptionalColumnExpr(names, "tenant_id", "0");
  284. var factoryExpr = BuildOptionalColumnExpr(names, "factory_id", "NULL");
  285. var companyExpr = BuildOptionalColumnExpr(names, "company_id", "NULL");
  286. var sourceRowExpr = names.Any(u => string.Equals(u, entity.SourceRowIdExpression, StringComparison.OrdinalIgnoreCase))
  287. ? $"s.`{FindColumn(names, entity.SourceRowIdExpression)}`"
  288. : entity.SourceRowIdExpression;
  289. var rawDataExpr = BuildJsonObjectExpression(names);
  290. var rowsRead = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM `{entity.SourceTable}`");
  291. var logId = await InsertSyncLogAsync(tenantId, entityRow.Id, entityRow.EntityName, batchId, rowsRead);
  292. var started = DateTime.Now;
  293. try
  294. {
  295. var affected = await _db.Ado.ExecuteCommandAsync(
  296. $"""
  297. INSERT INTO `{entity.TargetTable}`
  298. (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)
  299. SELECT
  300. {tenantExpr},
  301. {factoryExpr},
  302. {companyExpr},
  303. 'AIDOP',
  304. @SourceTable,
  305. CAST({sourceRowExpr} AS CHAR),
  306. CAST(COALESCE({entity.SourceBizKeyExpression}, CAST({sourceRowExpr} AS CHAR)) AS CHAR),
  307. @BatchId,
  308. @Now,
  309. 'PENDING',
  310. {rawDataExpr}
  311. FROM `{entity.SourceTable}` s
  312. ON DUPLICATE KEY UPDATE
  313. tenant_id=VALUES(tenant_id),
  314. factory_id=VALUES(factory_id),
  315. company_id=VALUES(company_id),
  316. source_row_id=VALUES(source_row_id),
  317. sync_batch_id=VALUES(sync_batch_id),
  318. sync_time=VALUES(sync_time),
  319. process_status=VALUES(process_status),
  320. raw_data=VALUES(raw_data),
  321. update_time=CURRENT_TIMESTAMP
  322. """,
  323. new SugarParameter("@SourceTable", entity.SourceTable),
  324. new SugarParameter("@BatchId", batchId),
  325. new SugarParameter("@Now", now));
  326. await MarkSyncLogSuccessAsync(logId, started, affected);
  327. return rowsRead;
  328. }
  329. catch (Exception ex)
  330. {
  331. await MarkSyncLogFailedAsync(logId, started, ex.Message);
  332. throw;
  333. }
  334. }
  335. private async Task<int> TransformStandardAsync(S1MdpRunScope scope, string batchId, DateTime now, CancellationToken cancellationToken)
  336. {
  337. using var db = _db.CopyNew();
  338. db.Ado.CommandTimeOut = Math.Max(db.Ado.CommandTimeOut, 180);
  339. var total = 0;
  340. foreach (var command in BuildStandardCommands(scope, batchId, now))
  341. {
  342. cancellationToken.ThrowIfCancellationRequested();
  343. total += await db.Ado.ExecuteCommandAsync(command.Sql, command.Parameters);
  344. }
  345. return total;
  346. }
  347. private async Task<int> BuildDwdAsync(S1MdpRunScope scope, string batchId, DateTime now, CancellationToken cancellationToken)
  348. {
  349. using var db = _db.CopyNew();
  350. db.Ado.CommandTimeOut = Math.Max(db.Ado.CommandTimeOut, 180);
  351. var total = 0;
  352. foreach (var command in BuildDwdCommands(scope, batchId, now))
  353. {
  354. cancellationToken.ThrowIfCancellationRequested();
  355. total += await db.Ado.ExecuteCommandAsync(command.Sql, command.Parameters);
  356. }
  357. return total;
  358. }
  359. private async Task<int> BuildS1KpiValuesAsync(S1MdpRunScope scope, string batchId, DateTime now, CancellationToken cancellationToken)
  360. {
  361. var statDate = now.Date;
  362. var rows = await CalculateS1KpiValuesAsync(scope, batchId, statDate);
  363. var affected = 0;
  364. foreach (var row in rows)
  365. {
  366. cancellationToken.ThrowIfCancellationRequested();
  367. if (row.TenantId != scope.TenantId || row.FactoryId != scope.FactoryId)
  368. continue;
  369. affected += await UpsertS1KpiValueAsync(row, statDate, now);
  370. }
  371. return affected;
  372. }
  373. private async Task<List<S1KpiCalcRow>> CalculateS1KpiValuesAsync(S1MdpRunScope scope, string batchId, DateTime statDate)
  374. {
  375. return await _db.Ado.SqlQueryAsync<S1KpiCalcRow>(
  376. """
  377. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  378. 'S1_L1_001' AS MetricCode,
  379. ROUND(AVG(TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) / 24), 4) AS MetricValue
  380. FROM mdp_std_so
  381. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  382. AND order_date IS NOT NULL
  383. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  384. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) >= order_date
  385. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  386. UNION ALL
  387. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  388. 'S1_L1_002' AS MetricCode,
  389. 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
  390. FROM mdp_std_so
  391. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  392. AND order_date IS NOT NULL
  393. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  394. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  395. UNION ALL
  396. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  397. 'S1_L1_003' AS MetricCode,
  398. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(planner_no, '')), 1), 4) AS MetricValue
  399. FROM mdp_std_so
  400. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  401. AND order_date IS NOT NULL AND order_date <= @StatDate
  402. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  403. UNION ALL
  404. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  405. 'S1_L1_004' AS MetricCode,
  406. ROUND(AVG(GREATEST(IFNULL(remaining_qty, 0), 0)) / GREATEST(SUM(IFNULL(shipped_qty, 0)), 1) * 30, 4) AS MetricValue
  407. FROM dwd_ship_trans
  408. WHERE calc_batch_id=@BatchId
  409. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  410. GROUP BY tenant_id
  411. HAVING SUM(IFNULL(shipped_qty, 0)) > 0
  412. UNION ALL
  413. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  414. 'S1_L2_010' AS MetricCode,
  415. ROUND(AVG(TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) / 24), 4) AS MetricValue
  416. FROM mdp_std_so
  417. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  418. AND order_date IS NOT NULL
  419. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  420. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) >= order_date
  421. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  422. UNION ALL
  423. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  424. 'S1_L2_011' AS MetricCode,
  425. 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
  426. FROM mdp_std_so
  427. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  428. AND order_date IS NOT NULL
  429. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  430. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  431. UNION ALL
  432. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  433. 'S1_L2_012' AS MetricCode,
  434. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(planner_no, '')), 1), 4) AS MetricValue
  435. FROM mdp_std_so
  436. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  437. AND order_date IS NOT NULL AND order_date <= @StatDate
  438. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  439. UNION ALL
  440. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  441. 'S1_L2_013' AS MetricCode,
  442. ROUND(100 * SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  443. FROM dwd_ship_trans
  444. WHERE calc_batch_id=@BatchId
  445. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  446. GROUP BY tenant_id
  447. UNION ALL
  448. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  449. 'S1_L2_014' AS MetricCode,
  450. ROUND(100 * SUM(CASE WHEN shipped_qty >= planned_ship_qty AND planned_ship_qty > 0 THEN 1 ELSE 0 END)
  451. / GREATEST(SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END), 1), 4) AS MetricValue
  452. FROM dwd_ship_trans
  453. WHERE calc_batch_id=@BatchId
  454. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  455. GROUP BY tenant_id
  456. HAVING SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) > 0
  457. UNION ALL
  458. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  459. 'S1_L2_015' AS MetricCode,
  460. ROUND(100 * SUM(CASE WHEN linkage_status = 'LINKED' THEN 1 ELSE 0 END)
  461. / GREATEST(SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END), 1), 4) AS MetricValue
  462. FROM dwd_ship_trans
  463. WHERE calc_batch_id=@BatchId
  464. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  465. GROUP BY tenant_id
  466. HAVING SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) > 0
  467. UNION ALL
  468. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  469. 'S1_L2_001' AS MetricCode,
  470. ROUND(AVG(TIMESTAMPDIFF(HOUR, create_time, update_time) / 24), 4) AS MetricValue
  471. FROM (
  472. SELECT tenant_id,
  473. 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,
  474. COALESCE(
  475. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  476. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  477. ) AS update_time
  478. FROM mdp_stg_so
  479. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  480. AND source_table = 'ado_contract_review'
  481. ) t
  482. WHERE create_time IS NOT NULL AND update_time IS NOT NULL AND update_time >= create_time
  483. GROUP BY tenant_id
  484. UNION ALL
  485. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  486. 'S1_L2_002' AS MetricCode,
  487. ROUND(100 * SUM(CASE WHEN TIMESTAMPDIFF(HOUR, create_time, update_time) <= 72 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  488. FROM (
  489. SELECT tenant_id,
  490. 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,
  491. COALESCE(
  492. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  493. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  494. ) AS update_time
  495. FROM mdp_stg_so
  496. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  497. AND source_table = 'ado_contract_review'
  498. ) t
  499. WHERE create_time IS NOT NULL AND update_time IS NOT NULL AND update_time >= create_time
  500. GROUP BY tenant_id
  501. UNION ALL
  502. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  503. 'S1_L2_003' AS MetricCode,
  504. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(owner_account, '')), 1), 4) AS MetricValue
  505. FROM (
  506. SELECT tenant_id,
  507. COALESCE(
  508. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ResponsibleAccount')),
  509. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))
  510. ) AS owner_account
  511. FROM mdp_stg_so
  512. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  513. AND source_table = 'ado_contract_review'
  514. ) t
  515. GROUP BY tenant_id
  516. UNION ALL
  517. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  518. 'S1_L2_004' AS MetricCode,
  519. ROUND(AVG(cycle_hours / 24), 4) AS MetricValue
  520. FROM (
  521. SELECT tenant_id,
  522. COALESCE(
  523. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingDesignCycle')) REGEXP '^-?[0-9]+$'
  524. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingDesignCycle')) AS DECIMAL(18,6)) END,
  525. TIMESTAMPDIFF(HOUR,
  526. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualStart')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  527. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualEnd')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  528. )
  529. ) AS cycle_hours
  530. FROM mdp_stg_so
  531. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  532. AND source_table = 'ado_product_design'
  533. ) t
  534. WHERE cycle_hours IS NOT NULL AND cycle_hours >= 0
  535. GROUP BY tenant_id
  536. UNION ALL
  537. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  538. 'S1_L2_005' AS MetricCode,
  539. ROUND(100 * SUM(CASE WHEN actual_end <= plan_end THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  540. FROM (
  541. SELECT tenant_id,
  542. 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,
  543. 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
  544. FROM mdp_stg_so
  545. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  546. AND source_table = 'ado_product_design'
  547. ) t
  548. WHERE plan_end IS NOT NULL AND actual_end IS NOT NULL
  549. GROUP BY tenant_id
  550. UNION ALL
  551. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  552. 'S1_L2_006' AS MetricCode,
  553. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(design_lead, '')), 1), 4) AS MetricValue
  554. FROM (
  555. SELECT tenant_id,
  556. COALESCE(
  557. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DesignLeadAccount')),
  558. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DesignLeadName')),
  559. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))
  560. ) AS design_lead
  561. FROM mdp_stg_so
  562. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  563. AND source_table = 'ado_product_design'
  564. ) t
  565. GROUP BY tenant_id
  566. UNION ALL
  567. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  568. CASE stage_no
  569. WHEN 1 THEN 'S1_L3_001'
  570. WHEN 2 THEN 'S1_L3_002'
  571. WHEN 3 THEN 'S1_L3_003'
  572. WHEN 4 THEN 'S1_L3_004'
  573. WHEN 5 THEN 'S1_L3_005'
  574. END AS MetricCode,
  575. ROUND(AVG(cycle_hours), 4) AS MetricValue
  576. FROM (
  577. SELECT tenant_id,
  578. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) REGEXP '^-?[0-9]+$'
  579. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) END AS stage_no,
  580. COALESCE(
  581. TIMESTAMPDIFF(HOUR,
  582. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  583. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  584. ),
  585. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  586. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  587. ) AS cycle_hours
  588. FROM mdp_stg_so
  589. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  590. AND source_table = 'ado_contract_review_flow'
  591. ) t
  592. WHERE stage_no BETWEEN 1 AND 5
  593. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  594. GROUP BY tenant_id, stage_no
  595. UNION ALL
  596. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  597. CASE stage_no
  598. WHEN 1 THEN 'S1_L3_101'
  599. WHEN 2 THEN 'S1_L3_102'
  600. WHEN 3 THEN 'S1_L3_103'
  601. WHEN 4 THEN 'S1_L3_104'
  602. WHEN 5 THEN 'S1_L3_105'
  603. END AS MetricCode,
  604. ROUND(AVG(cycle_hours), 4) AS MetricValue
  605. FROM (
  606. SELECT tenant_id,
  607. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) REGEXP '^-?[0-9]+$'
  608. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) END AS stage_no,
  609. COALESCE(
  610. TIMESTAMPDIFF(HOUR,
  611. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  612. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  613. ),
  614. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  615. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  616. ) AS cycle_hours
  617. FROM mdp_stg_so
  618. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  619. AND source_table = 'ado_contract_review_flow'
  620. ) t
  621. WHERE stage_no BETWEEN 1 AND 5
  622. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  623. GROUP BY tenant_id, stage_no
  624. UNION ALL
  625. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  626. metric_code AS MetricCode,
  627. ROUND(AVG(cycle_hours), 4) AS MetricValue
  628. FROM (
  629. SELECT tenant_id,
  630. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')), ''), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  631. WHEN 'LAW' THEN 'S1_L4_001'
  632. WHEN 'PRE_SALES' THEN 'S1_L4_002'
  633. WHEN 'MPS' THEN 'S1_L4_003'
  634. WHEN 'TEST' THEN 'S1_L4_004'
  635. WHEN '法律事务部' THEN 'S1_L4_001'
  636. WHEN '技术售前组' THEN 'S1_L4_002'
  637. WHEN '综合主计划' THEN 'S1_L4_003'
  638. WHEN '试验站' THEN 'S1_L4_004'
  639. END AS metric_code,
  640. COALESCE(
  641. TIMESTAMPDIFF(HOUR,
  642. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  643. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  644. ),
  645. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  646. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  647. ) AS cycle_hours
  648. FROM mdp_stg_so
  649. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  650. AND source_table = 'ado_contract_review_flow'
  651. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) = '1'
  652. ) t
  653. WHERE metric_code IS NOT NULL
  654. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  655. GROUP BY tenant_id, metric_code
  656. UNION ALL
  657. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  658. metric_code AS MetricCode,
  659. ROUND(AVG(cycle_hours), 4) AS MetricValue
  660. FROM (
  661. SELECT tenant_id,
  662. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')), ''), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  663. WHEN 'LAW' THEN 'S1_L4_101'
  664. WHEN 'PRE_SALES' THEN 'S1_L4_102'
  665. WHEN 'MPS' THEN 'S1_L4_103'
  666. WHEN 'TEST' THEN 'S1_L4_104'
  667. WHEN '法律事务部' THEN 'S1_L4_101'
  668. WHEN '技术售前组' THEN 'S1_L4_102'
  669. WHEN '综合主计划' THEN 'S1_L4_103'
  670. WHEN '试验站' THEN 'S1_L4_104'
  671. END AS metric_code,
  672. COALESCE(
  673. TIMESTAMPDIFF(HOUR,
  674. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  675. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  676. ),
  677. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  678. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  679. ) AS cycle_hours
  680. FROM mdp_stg_so
  681. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  682. AND source_table = 'ado_contract_review_flow'
  683. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) = '1'
  684. ) t
  685. WHERE metric_code IS NOT NULL
  686. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  687. GROUP BY tenant_id, metric_code
  688. """,
  689. new SugarParameter("@BatchId", batchId),
  690. new SugarParameter("@StatDate", statDate),
  691. new SugarParameter("@TenantId", scope.TenantId),
  692. new SugarParameter("@FactoryId", scope.FactoryId));
  693. }
  694. private async Task<int> UpsertS1KpiValueAsync(S1KpiCalcRow row, DateTime statDate, DateTime now)
  695. {
  696. var meta = await _db.Ado.SqlQuerySingleAsync<S1KpiMetaRow>(
  697. """
  698. SELECT MetricLevel, Direction, YellowThreshold, RedThreshold
  699. FROM ado_smart_ops_kpi_master
  700. WHERE TenantId=@TenantId AND ModuleCode='S1' AND MetricCode=@MetricCode AND IsEnabled=1
  701. LIMIT 1
  702. """,
  703. new SugarParameter("@TenantId", row.TenantId),
  704. new SugarParameter("@MetricCode", row.MetricCode));
  705. if (meta == null || row.MetricValue == null)
  706. return 0;
  707. var table = ResolveKpiValueTable(meta.MetricLevel);
  708. var current = await _db.Ado.SqlQuerySingleAsync<S1KpiValueRow>(
  709. $"""
  710. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  711. FROM {table}
  712. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  713. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  714. ORDER BY id
  715. LIMIT 1
  716. """,
  717. new SugarParameter("@TenantId", row.TenantId),
  718. new SugarParameter("@FactoryId", row.FactoryId),
  719. new SugarParameter("@MetricCode", row.MetricCode),
  720. new SugarParameter("@BizDate", statDate));
  721. var prior = await _db.Ado.SqlQuerySingleAsync<S1KpiValueRow>(
  722. $"""
  723. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  724. FROM {table}
  725. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  726. AND metric_code=@MetricCode AND biz_date<@BizDate AND is_deleted=0
  727. ORDER BY biz_date DESC, id DESC
  728. LIMIT 1
  729. """,
  730. new SugarParameter("@TenantId", row.TenantId),
  731. new SugarParameter("@FactoryId", row.FactoryId),
  732. new SugarParameter("@MetricCode", row.MetricCode),
  733. new SugarParameter("@BizDate", statDate));
  734. var actual = Math.Round(row.MetricValue.Value, 4);
  735. var snap = await _kpiTargetResolver.ResolveAsync(row.TenantId, row.FactoryId, row.MetricCode, "S1", statDate);
  736. var target = snap.TargetValue;
  737. var status = AidopS4KpiMerge.AchievementLevel(actual, target, meta.Direction, meta.YellowThreshold, meta.RedThreshold);
  738. var trend = ResolveTrendFlag(actual, prior?.MetricValue);
  739. if (current != null)
  740. {
  741. return await _db.Ado.ExecuteCommandAsync(
  742. $"""
  743. UPDATE {table}
  744. SET metric_value=@MetricValue, target_value=@TargetValue, status_color=@StatusColor, trend_flag=@TrendFlag,
  745. target_config_id=@TargetConfigId, target_source=@TargetSource, target_resolved_at=@TargetResolvedAt,
  746. is_active=1, status='ACTIVE', calc_time=@CalcTime, update_time=@CalcTime
  747. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  748. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  749. """,
  750. new SugarParameter("@MetricValue", actual),
  751. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  752. new SugarParameter("@StatusColor", status),
  753. new SugarParameter("@TrendFlag", trend),
  754. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  755. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  756. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt),
  757. new SugarParameter("@CalcTime", now),
  758. new SugarParameter("@TenantId", row.TenantId),
  759. new SugarParameter("@FactoryId", row.FactoryId),
  760. new SugarParameter("@MetricCode", row.MetricCode),
  761. new SugarParameter("@BizDate", statDate));
  762. }
  763. var nextId = await _db.Ado.GetLongAsync($"SELECT COALESCE(MAX(id), 0) + 1 FROM {table}");
  764. return await _db.Ado.ExecuteCommandAsync(
  765. $"""
  766. INSERT INTO {table}
  767. (id, tenant_id, factory_id, status, biz_date, create_time, update_time, is_deleted, is_active,
  768. module_code, metric_code, metric_value, target_value, status_color, trend_flag, calc_time,
  769. target_config_id, target_source, target_resolved_at)
  770. VALUES
  771. (@Id, @TenantId, @FactoryId, 'ACTIVE', @BizDate, @CalcTime, @CalcTime, 0, 1,
  772. 'S1', @MetricCode, @MetricValue, @TargetValue, @StatusColor, @TrendFlag, @CalcTime,
  773. @TargetConfigId, @TargetSource, @TargetResolvedAt)
  774. """,
  775. new SugarParameter("@Id", nextId),
  776. new SugarParameter("@TenantId", row.TenantId),
  777. new SugarParameter("@FactoryId", row.FactoryId),
  778. new SugarParameter("@BizDate", statDate),
  779. new SugarParameter("@CalcTime", now),
  780. new SugarParameter("@MetricCode", row.MetricCode),
  781. new SugarParameter("@MetricValue", actual),
  782. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  783. new SugarParameter("@StatusColor", status),
  784. new SugarParameter("@TrendFlag", trend),
  785. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  786. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  787. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt));
  788. }
  789. private IEnumerable<S1MdpSqlCommand> BuildStandardCommands(S1MdpRunScope scope, string batchId, DateTime now)
  790. {
  791. yield return Cmd(
  792. """
  793. INSERT INTO mdp_std_so
  794. (tenant_id, factory_id, company_id, source_system, order_id, order_entry_id, order_no, order_line, order_type,
  795. customer_id, customer_no, customer_name, customer_order_no, country, item_code, item_name, item_spec,
  796. map_number, map_name, bom_number, unit, order_qty, delivered_notice_qty, delivered_qty, price, tax_price,
  797. amount, total_amount, order_date, customer_request_date, plan_delivery_date, promised_delivery_date,
  798. capacity_date, material_ready_date, planner_no, planner_name, order_status, review_status, review_stage,
  799. flow_state, progress, urgent, closed, deleted_flag, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  800. SELECT
  801. COALESCE(NULLIF(e.tenant_id, 0), NULLIF(h.tenant_id, 0),
  802. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  803. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) AS SIGNED) END),
  804. COALESCE(NULLIF(NULLIF(e.factory_id, 0), e.tenant_id),
  805. NULLIF(NULLIF(h.factory_id, 0), h.tenant_id)),
  806. NULLIF(e.company_id, 0),
  807. 'AIDOP',
  808. 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,
  809. 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,
  810. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no')), e.source_biz_key),
  811. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.entry_seq')), JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.Id'))) AS CHAR),
  812. CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.order_type')) AS CHAR),
  813. 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,
  814. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_no')),
  815. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_name')),
  816. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.custom_order_bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_from'))),
  817. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.country')),
  818. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.item_number')),
  819. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.item_name')),
  820. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.specification')),
  821. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.map_number')),
  822. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.map_name')),
  823. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bom_number')),
  824. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.unit')),
  825. 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,
  826. 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,
  827. 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,
  828. 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,
  829. 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,
  830. 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,
  831. 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,
  832. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.date')), 'null'), ''),
  833. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.rdate')), 'null'), ''),
  834. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.plan_date')), 'null'), ''),
  835. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.date')), 'null'), ''),
  836. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.sys_capacity_date')), 'null'), ''),
  837. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.sys_material_date')), 'null'), ''),
  838. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.planner_no')),
  839. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.planner_name')),
  840. CASE WHEN JSON_EXTRACT(h.raw_data,'$.closed') IN (1, true) THEN 'CLOSED' ELSE 'OPEN' END,
  841. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.FlowStatus')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.flowstate'))),
  842. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.CurrentDept')), JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.CurrentStage'))),
  843. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.flowstate')),
  844. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.progress')),
  845. 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,
  846. 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,
  847. 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),
  848. e.source_table,
  849. e.source_row_id,
  850. e.source_biz_key,
  851. @BatchId,
  852. @Now
  853. FROM mdp_stg_so e
  854. LEFT JOIN mdp_stg_so h ON h.source_table='crm_seorder'
  855. AND h.tenant_id = e.tenant_id
  856. AND JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.Id')) = JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.seorder_id'))
  857. LEFT JOIN mdp_stg_so r ON r.source_table='ado_contract_review'
  858. AND r.tenant_id = e.tenant_id
  859. 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')))
  860. WHERE e.source_table='crm_seorderentry'
  861. AND e.tenant_id=@TenantId
  862. AND COALESCE(NULLIF(e.factory_id, 0), 1)=@FactoryId
  863. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no'))), '') <> ''
  864. AND COALESCE(NULLIF(e.tenant_id, 0), NULLIF(h.tenant_id, 0),
  865. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  866. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  867. ON DUPLICATE KEY UPDATE
  868. customer_no=VALUES(customer_no), customer_name=VALUES(customer_name), item_code=VALUES(item_code),
  869. item_name=VALUES(item_name), item_spec=VALUES(item_spec), order_qty=VALUES(order_qty),
  870. delivered_notice_qty=VALUES(delivered_notice_qty), delivered_qty=VALUES(delivered_qty),
  871. plan_delivery_date=VALUES(plan_delivery_date), promised_delivery_date=VALUES(promised_delivery_date),
  872. capacity_date=VALUES(capacity_date), material_ready_date=VALUES(material_ready_date),
  873. order_status=VALUES(order_status), review_status=VALUES(review_status), review_stage=VALUES(review_stage),
  874. flow_state=VALUES(flow_state), progress=VALUES(progress), deleted_flag=VALUES(deleted_flag),
  875. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  876. """, scope, batchId, now);
  877. yield return Cmd(
  878. """
  879. INSERT INTO mdp_std_ship_trans
  880. (tenant_id, factory_id, company_id, source_system, trans_type, plan_id, plan_no, plan_line, order_id, order_entry_id,
  881. order_no, order_line, customer_no, customer_name, country, item_code, item_name, item_spec, qty, plan_qty,
  882. weight, volume, order_date, plan_ship_date, shipping_site, shipping_address, consignee, telephone,
  883. status, confirm_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  884. SELECT
  885. COALESCE(NULLIF(d.tenant_id, 0),
  886. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  887. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  888. NULLIF(d.factory_id, 0),
  889. NULLIF(d.company_id, 0),
  890. 'AIDOP',
  891. 'SHIP_PLAN',
  892. 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,
  893. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.LotSerial')), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id')) AS CHAR)),
  894. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS CHAR),
  895. 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,
  896. 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,
  897. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr'))),
  898. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sentry_id')) AS CHAR),
  899. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CustomNo')),
  900. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CustomName')),
  901. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Country')),
  902. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')),
  903. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemName')),
  904. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Specification')),
  905. 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,
  906. 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,
  907. 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,
  908. 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,
  909. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdDate')), 'null'), ''),
  910. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingDate')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CreateTime'))), 'null'), ''),
  911. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingSite')),
  912. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingAddress')),
  913. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Consignee')),
  914. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Telephone')),
  915. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Status')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Status'))),
  916. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.IsConfirm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.IsConfirm'))) AS CHAR),
  917. d.source_table,
  918. d.source_row_id,
  919. d.source_biz_key,
  920. @BatchId,
  921. @Now
  922. FROM mdp_stg_ship_trans d
  923. LEFT JOIN mdp_stg_ship_trans m ON m.source_table='ShippingPlan'
  924. AND m.tenant_id = d.tenant_id
  925. AND JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.RecID')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id'))
  926. WHERE d.source_table='ShippingPlanDetail'
  927. AND d.tenant_id=@TenantId
  928. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  929. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr'))), '') <> ''
  930. AND COALESCE(NULLIF(d.tenant_id, 0),
  931. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  932. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  933. ON DUPLICATE KEY UPDATE
  934. plan_no=VALUES(plan_no), order_no=VALUES(order_no), customer_no=VALUES(customer_no), customer_name=VALUES(customer_name),
  935. item_code=VALUES(item_code), item_name=VALUES(item_name), qty=VALUES(qty), plan_qty=VALUES(plan_qty),
  936. plan_ship_date=VALUES(plan_ship_date), shipping_site=VALUES(shipping_site), shipping_address=VALUES(shipping_address),
  937. status=VALUES(status), confirm_status=VALUES(confirm_status), sync_batch_id=VALUES(sync_batch_id),
  938. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  939. """, scope, batchId, now);
  940. yield return Cmd(
  941. """
  942. INSERT INTO mdp_std_ship_trans
  943. (tenant_id, factory_id, source_system, trans_type, shipper_rec_id, shipper_no, shipper_line, order_no, order_line,
  944. customer_no, item_code, item_name, qty_to_ship, picking_qty, real_qty, gross_weight, net_weight, volume,
  945. plan_ship_date, actual_ship_date, site, status, confirm_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  946. SELECT
  947. COALESCE(NULLIF(d.tenant_id, 0),
  948. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  949. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  950. NULLIF(d.factory_id, 0),
  951. 'AIDOP',
  952. 'ASN_SHIPPER',
  953. 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,
  954. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Id')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Id'))),
  955. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS CHAR),
  956. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.OrdNbr'))),
  957. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdLine')) AS CHAR),
  958. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.SoldTo')),
  959. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ContainerItem')),
  960. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Descr')),
  961. 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,
  962. 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,
  963. 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,
  964. 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,
  965. 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,
  966. 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,
  967. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ShipDate')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShipDate'))), 'null'), ''),
  968. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ShipDate')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShipDate'))), 'null'), ''),
  969. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Site')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Site'))),
  970. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Status')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Status'))),
  971. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.IsConfirm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.IsConfirm'))) AS CHAR),
  972. d.source_table,
  973. d.source_row_id,
  974. d.source_biz_key,
  975. @BatchId,
  976. @Now
  977. FROM mdp_stg_ship_trans d
  978. LEFT JOIN mdp_stg_ship_trans m ON m.source_table='ASNBOLShipperMaster'
  979. AND m.tenant_id = d.tenant_id
  980. AND JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.RecID')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ASNBOLShipperRecID'))
  981. WHERE d.source_table='ASNBOLShipperDetail'
  982. AND d.tenant_id=@TenantId
  983. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  984. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Id')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Id'))), '') <> ''
  985. AND COALESCE(NULLIF(d.tenant_id, 0),
  986. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  987. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  988. ON DUPLICATE KEY UPDATE
  989. shipper_no=VALUES(shipper_no), order_no=VALUES(order_no), order_line=VALUES(order_line),
  990. customer_no=VALUES(customer_no), item_code=VALUES(item_code), item_name=VALUES(item_name),
  991. qty_to_ship=VALUES(qty_to_ship), picking_qty=VALUES(picking_qty), real_qty=VALUES(real_qty),
  992. actual_ship_date=VALUES(actual_ship_date), site=VALUES(site), status=VALUES(status),
  993. confirm_status=VALUES(confirm_status), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  994. update_time=CURRENT_TIMESTAMP
  995. """, scope, batchId, now);
  996. yield return Cmd(
  997. """
  998. INSERT INTO mdp_std_ship_trans
  999. (tenant_id, factory_id, source_system, trans_type, order_no, customer_no, item_code, item_name, qty,
  1000. plan_ship_date, status, linkage_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  1001. SELECT
  1002. COALESCE(NULLIF(d.tenant_id, 0),
  1003. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1004. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1005. NULLIF(d.factory_id, 0),
  1006. 'AIDOP',
  1007. 'LINKAGE_PLAN',
  1008. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')),
  1009. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.custom_no')),
  1010. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_number')),
  1011. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Descr')),
  1012. 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,
  1013. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sys_capacity_date')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.fystarttime'))), 'null'), ''),
  1014. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.type')),
  1015. CASE
  1016. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.isuse')) = '1' THEN 'LINKED'
  1017. ELSE 'INACTIVE'
  1018. END,
  1019. d.source_table,
  1020. d.source_row_id,
  1021. d.source_biz_key,
  1022. @BatchId,
  1023. @Now
  1024. FROM mdp_stg_ship_trans d
  1025. WHERE d.source_table='LinkagePlan'
  1026. AND d.tenant_id=@TenantId
  1027. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1028. AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), '') <> ''
  1029. AND COALESCE(NULLIF(d.tenant_id, 0),
  1030. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1031. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1032. ON DUPLICATE KEY UPDATE
  1033. order_no=VALUES(order_no), customer_no=VALUES(customer_no), item_code=VALUES(item_code),
  1034. item_name=VALUES(item_name), qty=VALUES(qty), plan_ship_date=VALUES(plan_ship_date),
  1035. status=VALUES(status), linkage_status=VALUES(linkage_status), sync_batch_id=VALUES(sync_batch_id),
  1036. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  1037. """, scope, batchId, now);
  1038. }
  1039. private IEnumerable<S1MdpSqlCommand> BuildDwdCommands(S1MdpRunScope scope, string batchId, DateTime now)
  1040. {
  1041. yield return Cmd(
  1042. """
  1043. INSERT INTO dwd_requirement_examine_detail
  1044. (tenant_id, factory_id, stat_date, row_id, parent_row_id, examine_id, order_entry_id, bill_no, morder_no,
  1045. num, item_number, item_name, bom_number, model, kitting_time, item_type, erp_cls_name, qty, wastage,
  1046. need_count, sqty, use_qty, self_lack_qty, lack_qty, mo_qty, make_qty, purchase_qty, purchase_occupy_qty,
  1047. satisfy_time, have_ic_subs, substitute_code, create_time, source_system, sync_batch_id, calc_batch_id, calc_time)
  1048. SELECT
  1049. COALESCE(NULLIF(d.tenant_id, 0), NULLIF(r.tenant_id, 0), NULLIF(so.tenant_id, 0),
  1050. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1051. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1052. COALESCE(NULLIF(d.factory_id, 0), NULLIF(r.factory_id, 0), NULLIF(so.factory_id, 0)),
  1053. @StatDate,
  1054. CAST(d.source_row_id AS SIGNED),
  1055. CASE
  1056. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.num')) = '1' THEN NULL
  1057. ELSE p.parent_row_id
  1058. END,
  1059. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id')) REGEXP '^-?[0-9]+$'
  1060. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id')) AS SIGNED) END,
  1061. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$'
  1062. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) AS SIGNED) END,
  1063. COALESCE(so.order_no, JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.bill_no'))),
  1064. JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.morder_no')),
  1065. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.num')) AS CHAR),
  1066. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_number')),
  1067. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_name')),
  1068. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bom_number')),
  1069. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.model')),
  1070. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.kitting_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  1071. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.type')) = '1' THEN '替代件' ELSE '标准件' END,
  1072. COALESCE(
  1073. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls_name')),
  1074. CASE JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls'))
  1075. WHEN '0' THEN '配置类'
  1076. WHEN '1' THEN '自制'
  1077. WHEN '2' THEN '委外加工'
  1078. WHEN '3' THEN '外购'
  1079. WHEN '4' THEN '虚拟件'
  1080. END
  1081. ),
  1082. 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,
  1083. 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,
  1084. 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,
  1085. 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,
  1086. 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,
  1087. 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,
  1088. 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,
  1089. 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,
  1090. CASE
  1091. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) = '0' AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls')) = '1'
  1092. 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
  1093. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  1094. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) AS DECIMAL(18,4))
  1095. END,
  1096. 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,
  1097. 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,
  1098. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.satisfy_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  1099. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.haveicsubs')) = '1' THEN '是' ELSE '否' END,
  1100. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.substitute_code')),
  1101. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.create_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  1102. 'AIDOP',
  1103. d.sync_batch_id,
  1104. @BatchId,
  1105. @Now
  1106. FROM mdp_stg_so d
  1107. INNER JOIN mdp_stg_so r ON r.source_table='b_examine_result'
  1108. AND r.tenant_id = d.tenant_id
  1109. AND r.source_row_id = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id'))
  1110. LEFT JOIN mdp_std_so so ON so.tenant_id = COALESCE(r.tenant_id, d.tenant_id)
  1111. AND COALESCE(NULLIF(so.factory_id, 0), 1)=@FactoryId
  1112. AND so.order_entry_id = CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$'
  1113. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) AS SIGNED) END
  1114. LEFT JOIN (
  1115. SELECT examine_id, MIN(source_row_id) AS parent_row_id
  1116. FROM (
  1117. SELECT JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.examine_id')) AS examine_id, CAST(source_row_id AS SIGNED) AS source_row_id
  1118. FROM mdp_stg_so
  1119. WHERE source_table='b_bom_child_examine'
  1120. AND tenant_id=@TenantId
  1121. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1122. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.num')) = '1'
  1123. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.is_use')) IN ('1', 'true', 'True', 'base64:type16:AQ==')
  1124. ) x
  1125. GROUP BY examine_id
  1126. ) p ON p.examine_id = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id'))
  1127. WHERE d.source_table='b_bom_child_examine'
  1128. AND d.tenant_id=@TenantId
  1129. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1130. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.is_use')) IN ('1', 'true', 'True', 'base64:type16:AQ==')
  1131. AND COALESCE(NULLIF(d.tenant_id, 0), NULLIF(r.tenant_id, 0), NULLIF(so.tenant_id, 0),
  1132. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1133. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1134. ON DUPLICATE KEY UPDATE
  1135. parent_row_id=VALUES(parent_row_id), order_entry_id=VALUES(order_entry_id), bill_no=VALUES(bill_no),
  1136. morder_no=VALUES(morder_no), num=VALUES(num), item_number=VALUES(item_number), item_name=VALUES(item_name),
  1137. bom_number=VALUES(bom_number), model=VALUES(model), kitting_time=VALUES(kitting_time),
  1138. item_type=VALUES(item_type), erp_cls_name=VALUES(erp_cls_name), qty=VALUES(qty), wastage=VALUES(wastage),
  1139. need_count=VALUES(need_count), sqty=VALUES(sqty), use_qty=VALUES(use_qty), self_lack_qty=VALUES(self_lack_qty),
  1140. lack_qty=VALUES(lack_qty), mo_qty=VALUES(mo_qty), make_qty=VALUES(make_qty), purchase_qty=VALUES(purchase_qty),
  1141. purchase_occupy_qty=VALUES(purchase_occupy_qty), satisfy_time=VALUES(satisfy_time), have_ic_subs=VALUES(have_ic_subs),
  1142. substitute_code=VALUES(substitute_code), create_time=VALUES(create_time), sync_batch_id=VALUES(sync_batch_id),
  1143. calc_time=VALUES(calc_time), update_time=CURRENT_TIMESTAMP
  1144. """, scope, batchId, now);
  1145. yield return Cmd(
  1146. """
  1147. INSERT INTO dwd_ship_trans
  1148. (tenant_id, factory_id, company_id, stat_date, order_id, order_entry_id, order_no, order_line, customer_no,
  1149. customer_name, country, item_code, item_name, item_spec, order_qty, planned_ship_qty, shipped_qty,
  1150. remaining_qty, order_date, customer_request_date, plan_delivery_date, promised_delivery_date,
  1151. plan_ship_date, actual_ship_date, review_status, order_status, delivery_status, linkage_status, risk_level,
  1152. source_system, source_table, source_row_id, source_biz_key, sync_batch_id, calc_batch_id, calc_time)
  1153. SELECT
  1154. so.tenant_id,
  1155. so.factory_id,
  1156. so.company_id,
  1157. @StatDate,
  1158. so.order_id,
  1159. so.order_entry_id,
  1160. so.order_no,
  1161. IFNULL(so.order_line, ''),
  1162. so.customer_no,
  1163. so.customer_name,
  1164. so.country,
  1165. IFNULL(so.item_code, ''),
  1166. so.item_name,
  1167. so.item_spec,
  1168. IFNULL(so.order_qty, 0),
  1169. IFNULL(p.plan_qty, 0),
  1170. IFNULL(a.real_qty, 0),
  1171. GREATEST(IFNULL(so.order_qty, 0) - IFNULL(a.real_qty, 0), 0),
  1172. so.order_date,
  1173. so.customer_request_date,
  1174. so.plan_delivery_date,
  1175. so.promised_delivery_date,
  1176. p.plan_ship_date,
  1177. a.actual_ship_date,
  1178. so.review_status,
  1179. so.order_status,
  1180. CASE
  1181. WHEN IFNULL(so.order_qty, 0) > 0 AND IFNULL(a.real_qty, 0) >= IFNULL(so.order_qty, 0) THEN 'COMPLETED'
  1182. WHEN COALESCE(p.plan_ship_date, so.promised_delivery_date, so.plan_delivery_date) < @Now THEN 'DELAYED'
  1183. WHEN IFNULL(p.plan_qty, 0) > 0 THEN 'PLANNED'
  1184. ELSE 'OPEN'
  1185. END,
  1186. l.linkage_status,
  1187. CASE
  1188. WHEN COALESCE(p.plan_ship_date, so.promised_delivery_date, so.plan_delivery_date) < @Now
  1189. AND IFNULL(a.real_qty, 0) < IFNULL(so.order_qty, 0) THEN 'HIGH'
  1190. WHEN IFNULL(a.real_qty, 0) < IFNULL(so.order_qty, 0) THEN 'MEDIUM'
  1191. ELSE 'LOW'
  1192. END,
  1193. 'AIDOP',
  1194. so.source_table,
  1195. so.source_row_id,
  1196. so.source_biz_key,
  1197. so.sync_batch_id,
  1198. @BatchId,
  1199. @Now
  1200. FROM mdp_std_so so
  1201. LEFT JOIN (
  1202. SELECT tenant_id, order_no, order_entry_id, IFNULL(order_line, '') AS order_line, IFNULL(item_code, '') AS item_code,
  1203. SUM(IFNULL(plan_qty, IFNULL(qty, 0))) AS plan_qty,
  1204. MIN(plan_ship_date) AS plan_ship_date
  1205. FROM mdp_std_ship_trans
  1206. WHERE trans_type='SHIP_PLAN'
  1207. AND tenant_id=@TenantId
  1208. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1209. GROUP BY tenant_id, order_no, order_entry_id, IFNULL(order_line, ''), IFNULL(item_code, '')
  1210. ) p ON so.tenant_id=p.tenant_id
  1211. AND so.order_no=p.order_no
  1212. AND IFNULL(so.item_code, '')=p.item_code
  1213. AND (
  1214. (p.order_entry_id IS NOT NULL AND so.order_entry_id=p.order_entry_id)
  1215. OR (p.order_entry_id IS NULL AND IFNULL(so.order_line, '')=p.order_line)
  1216. )
  1217. LEFT JOIN (
  1218. SELECT tenant_id, order_no, IFNULL(order_line, '') AS order_line, IFNULL(item_code, '') AS item_code,
  1219. SUM(IFNULL(real_qty, IFNULL(qty_to_ship, 0))) AS real_qty,
  1220. MAX(actual_ship_date) AS actual_ship_date
  1221. FROM mdp_std_ship_trans
  1222. WHERE trans_type='ASN_SHIPPER'
  1223. AND tenant_id=@TenantId
  1224. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1225. GROUP BY tenant_id, order_no, IFNULL(order_line, ''), IFNULL(item_code, '')
  1226. ) 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
  1227. LEFT JOIN (
  1228. SELECT tenant_id, order_no, item_code, MAX(linkage_status) AS linkage_status
  1229. FROM mdp_std_ship_trans
  1230. WHERE IFNULL(linkage_status, '') <> ''
  1231. AND tenant_id=@TenantId
  1232. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1233. GROUP BY tenant_id, order_no, item_code
  1234. ) l ON so.tenant_id=l.tenant_id AND so.order_no=l.order_no AND IFNULL(so.item_code, '')=IFNULL(l.item_code, '')
  1235. WHERE so.tenant_id=@TenantId
  1236. AND COALESCE(NULLIF(so.factory_id, 0), 1)=@FactoryId
  1237. AND IFNULL(so.order_no, '') <> ''
  1238. ON DUPLICATE KEY UPDATE
  1239. customer_no=VALUES(customer_no), customer_name=VALUES(customer_name), item_name=VALUES(item_name),
  1240. order_qty=VALUES(order_qty), planned_ship_qty=VALUES(planned_ship_qty), shipped_qty=VALUES(shipped_qty),
  1241. remaining_qty=VALUES(remaining_qty), plan_ship_date=VALUES(plan_ship_date), actual_ship_date=VALUES(actual_ship_date),
  1242. review_status=VALUES(review_status), order_status=VALUES(order_status), delivery_status=VALUES(delivery_status),
  1243. linkage_status=VALUES(linkage_status), risk_level=VALUES(risk_level), sync_batch_id=VALUES(sync_batch_id),
  1244. calc_batch_id=VALUES(calc_batch_id), calc_time=VALUES(calc_time), update_time=CURRENT_TIMESTAMP
  1245. """, scope, batchId, now);
  1246. }
  1247. private async Task<long> InsertSyncLogAsync(long tenantId, long entityId, string entityName, string batchId, int rowsRead)
  1248. {
  1249. await _db.Ado.ExecuteCommandAsync(
  1250. """
  1251. INSERT INTO mdp_sync_log
  1252. (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type, sync_start, rows_read, status)
  1253. VALUES (@TenantId, @EntityId, 'AIDOPDEV_MYSQL', @EntityName, @BatchId, 'FULL', 'AUTO', NOW(), @RowsRead, 'RUNNING')
  1254. """,
  1255. new SugarParameter("@TenantId", tenantId),
  1256. new SugarParameter("@EntityId", entityId),
  1257. new SugarParameter("@EntityName", entityName),
  1258. new SugarParameter("@BatchId", batchId),
  1259. new SugarParameter("@RowsRead", rowsRead));
  1260. return await _db.Ado.GetLongAsync(
  1261. "SELECT id FROM mdp_sync_log WHERE sync_batch_id=@BatchId AND entity_id=@EntityId ORDER BY id DESC LIMIT 1",
  1262. new List<SugarParameter>
  1263. {
  1264. new("@BatchId", batchId),
  1265. new("@EntityId", entityId)
  1266. });
  1267. }
  1268. private async Task MarkSyncLogSuccessAsync(long logId, DateTime started, int affected)
  1269. {
  1270. await _db.Ado.ExecuteCommandAsync(
  1271. """
  1272. UPDATE mdp_sync_log
  1273. SET sync_end=NOW(), duration_ms=@DurationMs, rows_insert=@RowsInsert, rows_update=0, rows_skip=0, rows_error=0, status='SUCCESS'
  1274. WHERE id=@Id
  1275. """,
  1276. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  1277. new SugarParameter("@RowsInsert", affected),
  1278. new SugarParameter("@Id", logId));
  1279. }
  1280. private async Task MarkSyncLogFailedAsync(long logId, DateTime started, string message)
  1281. {
  1282. try
  1283. {
  1284. await _db.Ado.ExecuteCommandAsync(
  1285. """
  1286. UPDATE mdp_sync_log
  1287. SET sync_end=NOW(), duration_ms=@DurationMs, rows_error=1, status='FAILED', error_msg=@ErrorMsg
  1288. WHERE id=@Id
  1289. """,
  1290. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  1291. new SugarParameter("@ErrorMsg", Truncate(message, 1000)),
  1292. new SugarParameter("@Id", logId));
  1293. }
  1294. catch (Exception ex)
  1295. {
  1296. // 写库自身失败兜底:避免再抛掩盖原异常;遗留 RUNNING 行可由运维手动清理
  1297. Console.Error.WriteLine($"[S1MdpSyncTransform] MarkSyncLogFailed write failed (syncLogId={logId}): {ex.Message}");
  1298. }
  1299. }
  1300. private async Task<long> InsertTransformRunLogAsync(S1MdpRunScope scope, string batchId, DateTime startedAt, string triggerType)
  1301. {
  1302. await _db.Ado.ExecuteCommandAsync(
  1303. """
  1304. INSERT INTO mdp_transform_run_log
  1305. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  1306. VALUES (@TenantId, @JobCode, 'S1 MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  1307. """,
  1308. new SugarParameter("@TenantId", scope.TenantId),
  1309. new SugarParameter("@JobCode", JobCode),
  1310. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  1311. new SugarParameter("@BatchId", batchId),
  1312. new SugarParameter("@StartTime", startedAt));
  1313. return await _db.Ado.GetLongAsync(
  1314. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  1315. new List<SugarParameter> { new("@BatchId", batchId) });
  1316. }
  1317. private async Task MarkTransformRunSuccessAsync(long runLogId, DateTime startedAt, S1MdpSyncTransformResult result)
  1318. {
  1319. var finishedAt = DateTime.Now;
  1320. await _db.Ado.ExecuteCommandAsync(
  1321. """
  1322. UPDATE mdp_transform_run_log
  1323. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  1324. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=@DwdRows,
  1325. summary_json=@SummaryJson, update_time=CURRENT_TIMESTAMP
  1326. WHERE id=@Id
  1327. """,
  1328. new SugarParameter("@EndTime", finishedAt),
  1329. new SugarParameter("@DurationMs", ToDurationMs(finishedAt - startedAt)),
  1330. new SugarParameter("@StageRows", result.StageRows),
  1331. new SugarParameter("@StandardRows", result.StandardRows),
  1332. new SugarParameter("@DwdRows", result.DwdRows),
  1333. new SugarParameter("@SummaryJson", BuildRunSummaryJson(result)),
  1334. new SugarParameter("@Id", runLogId));
  1335. }
  1336. private async Task MarkTransformRunFailedAsync(long runLogId, DateTime startedAt, string message)
  1337. {
  1338. try
  1339. {
  1340. var finishedAt = DateTime.Now;
  1341. var db = _db.CopyNew();
  1342. await db.Ado.ExecuteCommandAsync(
  1343. """
  1344. UPDATE mdp_transform_run_log
  1345. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  1346. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  1347. WHERE id=@Id
  1348. """,
  1349. new SugarParameter("@EndTime", finishedAt),
  1350. new SugarParameter("@DurationMs", ToDurationMs(finishedAt - startedAt)),
  1351. new SugarParameter("@ErrorMessage", Truncate(message, 2000)),
  1352. new SugarParameter("@Id", runLogId));
  1353. }
  1354. catch (Exception ex)
  1355. {
  1356. // 写库自身失败兜底(典型场景:远端 MySQL 瞬断导致 MarkFailed 自身也连不上):
  1357. // 避免再抛二次异常掩盖原错;遗留 RUNNING 行可由运维手动清理。
  1358. Console.Error.WriteLine($"[S1MdpSyncTransform] MarkTransformRunFailed write failed (runLogId={runLogId}): {ex.Message}");
  1359. }
  1360. }
  1361. private static S1MdpSqlCommand Cmd(string sql, S1MdpRunScope scope, string batchId, DateTime now)
  1362. {
  1363. return new S1MdpSqlCommand(sql, new[]
  1364. {
  1365. new SugarParameter("@BatchId", batchId),
  1366. new SugarParameter("@Now", now),
  1367. new SugarParameter("@StatDate", now.Date),
  1368. new SugarParameter("@TenantId", scope.TenantId),
  1369. new SugarParameter("@FactoryId", scope.FactoryId)
  1370. });
  1371. }
  1372. private static string BuildJsonObjectExpression(IEnumerable<string> columns)
  1373. {
  1374. var parts = columns.SelectMany(c => new[] { $"'{c.Replace("'", "''")}'", $"s.`{c}`" });
  1375. return $"JSON_OBJECT({string.Join(",", parts)})";
  1376. }
  1377. private static string BuildOptionalColumnExpr(IReadOnlyCollection<string> columns, string expected, string fallback)
  1378. {
  1379. return columns.Any(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase))
  1380. ? $"s.`{FindColumn(columns, expected)}`"
  1381. : fallback;
  1382. }
  1383. private static string FindColumn(IEnumerable<string> columns, string expected)
  1384. {
  1385. return columns.First(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase));
  1386. }
  1387. private static async Task KeepLeaseAliveAsync(IS1MdpFullRunLease lease, CancellationToken ct)
  1388. {
  1389. while (!ct.IsCancellationRequested)
  1390. {
  1391. try
  1392. {
  1393. await Task.Delay(TimeSpan.FromSeconds(30), ct);
  1394. await lease.HeartbeatAsync(CancellationToken.None);
  1395. }
  1396. catch (OperationCanceledException)
  1397. {
  1398. return;
  1399. }
  1400. }
  1401. }
  1402. private static int ToDurationMs(TimeSpan elapsed)
  1403. {
  1404. var ms = elapsed.TotalMilliseconds;
  1405. if (double.IsNaN(ms) || ms <= 0) return 0;
  1406. return ms >= int.MaxValue ? int.MaxValue : (int)ms;
  1407. }
  1408. private static string NormalizeTriggerType(string? triggerType)
  1409. {
  1410. return string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  1411. }
  1412. private static string BuildRunSummaryJson(S1MdpSyncTransformResult result)
  1413. {
  1414. return $$"""{"batchId":"{{result.BatchId}}","stageRows":{{result.StageRows}},"standardRows":{{result.StandardRows}},"dwdRows":{{result.DwdRows}},"kpiRows":{{result.KpiRows}}}""";
  1415. }
  1416. private static string ResolveKpiValueTable(int metricLevel)
  1417. {
  1418. return metricLevel switch
  1419. {
  1420. 1 => "ado_s9_kpi_value_l1_day",
  1421. 2 => "ado_s9_kpi_value_l2_day",
  1422. 3 => "ado_s9_kpi_value_l3_day",
  1423. 4 => "ado_s9_kpi_value_l4_day",
  1424. _ => "ado_s9_kpi_value_l2_day"
  1425. };
  1426. }
  1427. private static decimal DefaultS1Target(string metricCode) => LegacyKpiCodeTargets.GetOrZero(metricCode);
  1428. private static string ResolveKpiStatus(decimal actual, decimal target, string? direction, decimal? yellowThreshold, decimal? redThreshold)
  1429. {
  1430. if (target <= 0) return "gray";
  1431. var ratio = actual / target * 100m;
  1432. if (string.Equals(direction, "lower_is_better", StringComparison.OrdinalIgnoreCase))
  1433. {
  1434. if (actual <= target) return "green";
  1435. if (ratio <= (yellowThreshold ?? 110m)) return "yellow";
  1436. return ratio >= (redThreshold ?? 120m) ? "red" : "yellow";
  1437. }
  1438. if (actual >= target) return "green";
  1439. if (ratio >= (yellowThreshold ?? 95m)) return "yellow";
  1440. return ratio <= (redThreshold ?? 80m) ? "red" : "yellow";
  1441. }
  1442. private static string ResolveTrendFlag(decimal actual, decimal? previous)
  1443. {
  1444. if (previous == null) return "flat";
  1445. if (actual > previous.Value) return "up";
  1446. if (actual < previous.Value) return "down";
  1447. return "flat";
  1448. }
  1449. private static string Truncate(string? raw, int maxLength)
  1450. {
  1451. if (string.IsNullOrEmpty(raw)) return string.Empty;
  1452. return raw.Length <= maxLength ? raw : raw[..maxLength];
  1453. }
  1454. private sealed class S1ColumnRow
  1455. {
  1456. public string ColumnName { get; set; } = string.Empty;
  1457. }
  1458. private sealed class S1MdpEntityRow
  1459. {
  1460. public long Id { get; set; }
  1461. public string EntityName { get; set; } = string.Empty;
  1462. }
  1463. private sealed class S1KpiCalcRow
  1464. {
  1465. public long TenantId { get; set; }
  1466. public long FactoryId { get; set; }
  1467. public string MetricCode { get; set; } = string.Empty;
  1468. public decimal? MetricValue { get; set; }
  1469. }
  1470. private sealed class S1KpiMetaRow
  1471. {
  1472. public int MetricLevel { get; set; }
  1473. public string Direction { get; set; } = "higher_is_better";
  1474. public decimal? YellowThreshold { get; set; }
  1475. public decimal? RedThreshold { get; set; }
  1476. }
  1477. private sealed class S1KpiValueRow
  1478. {
  1479. public long Id { get; set; }
  1480. public decimal? MetricValue { get; set; }
  1481. public decimal? TargetValue { get; set; }
  1482. }
  1483. }
  1484. public sealed class S1MdpSyncTransformResult
  1485. {
  1486. public long RunLogId { get; set; }
  1487. public string BatchId { get; set; } = string.Empty;
  1488. public int StageRows { get; set; }
  1489. public int StandardRows { get; set; }
  1490. public int DwdRows { get; set; }
  1491. public int KpiRows { get; set; }
  1492. public int AtomicRows { get; set; }
  1493. }
  1494. internal sealed record S1MdpSqlCommand(string Sql, SugarParameter[] Parameters);
  1495. internal sealed record S1MdpEntityConfig(
  1496. string EntityCode,
  1497. string SourceTable,
  1498. string TargetTable,
  1499. string SourceRowIdExpression,
  1500. string SourceBizKeyExpression)
  1501. {
  1502. public static readonly IReadOnlyList<S1MdpEntityConfig> All = new List<S1MdpEntityConfig>
  1503. {
  1504. new("S1_SEORDER", "crm_seorder", "mdp_stg_so", "Id", "COALESCE(s.`bill_no`, CAST(s.`Id` AS CHAR))"),
  1505. new("S1_SEORDER_ENTRY", "crm_seorderentry", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`entry_seq`, CAST(s.`Id` AS CHAR)))"),
  1506. new("S1_SEORDER_CHANGE", "crm_seorder_change", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', CAST(s.`Id` AS CHAR))"),
  1507. new("S1_CONTRACT_REVIEW", "ado_contract_review", "mdp_stg_so", "RecID", "COALESCE(s.`BillNo`, CAST(s.`RecID` AS CHAR))"),
  1508. 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))"),
  1509. new("S1_PRODUCT_DESIGN", "ado_product_design", "mdp_stg_so", "Id", "COALESCE(s.`BillNo`, CAST(s.`Id` AS CHAR))"),
  1510. new("S1_PRODUCT_DESIGN_BOM", "ado_product_design_bom", "mdp_stg_so", "Id", "CONCAT(CAST(s.`ProductDesignId` AS CHAR), ':', CAST(s.`Id` AS CHAR))"),
  1511. new("S1_PRODUCT_DESIGN_ROUTING", "ado_product_design_routing", "mdp_stg_so", "Id", "CONCAT(CAST(s.`ProductDesignId` AS CHAR), ':', CAST(s.`Id` AS CHAR))"),
  1512. 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))"),
  1513. 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))"),
  1514. new("S1_SHIPPING_PLAN", "ShippingPlan", "mdp_stg_ship_trans", "RecID", "COALESCE(s.`LotSerial`, CAST(s.`RecID` AS CHAR))"),
  1515. new("S1_SHIPPING_PLAN_DETAIL", "ShippingPlanDetail", "mdp_stg_ship_trans", "RecID", "CONCAT(IFNULL(s.`plan_id`,''), ':', IFNULL(s.`OrdNbr`,''), ':', CAST(s.`RecID` AS CHAR))"),
  1516. new("S1_ASN_SHIPPER_MASTER", "ASNBOLShipperMaster", "mdp_stg_ship_trans", "RecID", "COALESCE(s.`Id`, CONCAT(IFNULL(s.`OrdNbr`,''), ':', CAST(s.`RecID` AS CHAR)))"),
  1517. new("S1_ASN_SHIPPER_DETAIL", "ASNBOLShipperDetail", "mdp_stg_ship_trans", "RecID", "CONCAT(IFNULL(s.`Id`,''), ':', IFNULL(s.`Line`, CAST(s.`RecID` AS CHAR)))"),
  1518. new("S1_LINKAGE_PLAN", "LinkagePlan", "mdp_stg_ship_trans", "id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`item_number`,''), ':', CAST(s.`id` AS CHAR))")
  1519. };
  1520. }