S1MdpSyncTransformService.cs 106 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710
  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 affected = 0;
  362. for (var dayOffset = 13; dayOffset >= 0; dayOffset--)
  363. {
  364. var statDate = now.Date.AddDays(-dayOffset);
  365. var rows = await CalculateS1KpiValuesAsync(scope, batchId, statDate);
  366. foreach (var row in rows)
  367. {
  368. cancellationToken.ThrowIfCancellationRequested();
  369. if (row.TenantId != scope.TenantId || row.FactoryId != scope.FactoryId)
  370. continue;
  371. affected += await UpsertS1KpiValueAsync(row, statDate, now);
  372. }
  373. }
  374. return affected;
  375. }
  376. private async Task<List<S1KpiCalcRow>> CalculateS1KpiValuesAsync(S1MdpRunScope scope, string batchId, DateTime statDate)
  377. {
  378. return await _db.Ado.SqlQueryAsync<S1KpiCalcRow>(
  379. """
  380. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  381. 'S1_L1_001' AS MetricCode,
  382. ROUND(AVG(TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) / 24), 4) AS MetricValue
  383. FROM mdp_std_so
  384. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  385. AND order_date IS NOT NULL
  386. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  387. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) >= order_date
  388. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  389. UNION ALL
  390. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  391. 'S1_L1_002' AS MetricCode,
  392. 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
  393. FROM mdp_std_so
  394. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  395. AND order_date IS NOT NULL
  396. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  397. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  398. UNION ALL
  399. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  400. 'S1_L1_003' AS MetricCode,
  401. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(planner_no, '')), 1), 4) AS MetricValue
  402. FROM mdp_std_so
  403. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  404. AND order_date IS NOT NULL AND order_date <= @StatDate
  405. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  406. UNION ALL
  407. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  408. 'S1_L1_004' AS MetricCode,
  409. ROUND(AVG(GREATEST(IFNULL(remaining_qty, 0), 0)) / GREATEST(SUM(IFNULL(shipped_qty, 0)), 1) * 30, 4) AS MetricValue
  410. FROM dwd_ship_trans
  411. WHERE calc_batch_id=@BatchId
  412. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  413. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  414. HAVING SUM(IFNULL(shipped_qty, 0)) > 0
  415. UNION ALL
  416. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  417. 'S1_L2_010' AS MetricCode,
  418. ROUND(AVG(TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) / 24), 4) AS MetricValue
  419. FROM mdp_std_so
  420. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  421. AND order_date IS NOT NULL
  422. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  423. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) >= order_date
  424. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  425. UNION ALL
  426. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  427. 'S1_L2_011' AS MetricCode,
  428. 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
  429. FROM mdp_std_so
  430. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  431. AND order_date IS NOT NULL
  432. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  433. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  434. UNION ALL
  435. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  436. 'S1_L2_012' AS MetricCode,
  437. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(planner_no, '')), 1), 4) AS MetricValue
  438. FROM mdp_std_so
  439. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  440. AND order_date IS NOT NULL AND order_date <= @StatDate
  441. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  442. UNION ALL
  443. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  444. 'S1_L2_013' AS MetricCode,
  445. ROUND(100 * SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  446. FROM dwd_ship_trans
  447. WHERE calc_batch_id=@BatchId
  448. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  449. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  450. UNION ALL
  451. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  452. 'S1_L2_014' AS MetricCode,
  453. ROUND(100 * SUM(CASE WHEN shipped_qty >= planned_ship_qty AND planned_ship_qty > 0 THEN 1 ELSE 0 END)
  454. / GREATEST(SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END), 1), 4) AS MetricValue
  455. FROM dwd_ship_trans
  456. WHERE calc_batch_id=@BatchId
  457. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  458. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  459. HAVING SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) > 0
  460. UNION ALL
  461. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  462. 'S1_L2_015' AS MetricCode,
  463. ROUND(100 * SUM(CASE WHEN linkage_status = 'LINKED' THEN 1 ELSE 0 END)
  464. / GREATEST(SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END), 1), 4) AS MetricValue
  465. FROM dwd_ship_trans
  466. WHERE calc_batch_id=@BatchId
  467. AND tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  468. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  469. HAVING SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) > 0
  470. UNION ALL
  471. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  472. 'S1_L2_001' AS MetricCode,
  473. ROUND(AVG(TIMESTAMPDIFF(HOUR, create_time, update_time) / 24), 4) AS MetricValue
  474. FROM (
  475. SELECT tenant_id,
  476. 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,
  477. COALESCE(
  478. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  479. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  480. ) AS update_time
  481. FROM mdp_stg_so
  482. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  483. AND source_table = 'ado_contract_review'
  484. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-PROC-%'
  485. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-DEL-%'
  486. ) t
  487. WHERE create_time IS NOT NULL AND update_time IS NOT NULL AND update_time >= create_time
  488. GROUP BY tenant_id
  489. UNION ALL
  490. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  491. 'S1_L2_002' AS MetricCode,
  492. ROUND(100 * SUM(CASE WHEN TIMESTAMPDIFF(HOUR, create_time, update_time) <= 72 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  493. FROM (
  494. SELECT tenant_id,
  495. 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,
  496. COALESCE(
  497. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  498. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  499. ) AS update_time
  500. FROM mdp_stg_so
  501. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  502. AND source_table = 'ado_contract_review'
  503. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-PROC-%'
  504. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-DEL-%'
  505. ) t
  506. WHERE create_time IS NOT NULL AND update_time IS NOT NULL AND update_time >= create_time
  507. GROUP BY tenant_id
  508. UNION ALL
  509. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  510. 'S1_L2_003' AS MetricCode,
  511. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(owner_account, '')), 1), 4) AS MetricValue
  512. FROM (
  513. SELECT tenant_id,
  514. COALESCE(
  515. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ResponsibleAccount')),
  516. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))
  517. ) AS owner_account
  518. FROM mdp_stg_so
  519. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  520. AND source_table = 'ado_contract_review'
  521. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-PROC-%'
  522. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') NOT LIKE 'UAT-DEL-%'
  523. ) t
  524. GROUP BY tenant_id
  525. UNION ALL
  526. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  527. 'S1_L2_007' AS MetricCode,
  528. ROUND(AVG(TIMESTAMPDIFF(HOUR,create_time,update_time)/24),4) AS MetricValue
  529. FROM (
  530. SELECT tenant_id,
  531. 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,
  532. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f') AS update_time
  533. FROM mdp_stg_so
  534. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  535. AND source_table='ado_contract_review'
  536. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') LIKE 'UAT-PROC-%'
  537. ) t
  538. WHERE create_time IS NOT NULL AND update_time>=create_time
  539. GROUP BY tenant_id
  540. UNION ALL
  541. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  542. 'S1_L2_008' AS MetricCode,
  543. ROUND(100*SUM(CASE WHEN TIMESTAMPDIFF(HOUR,create_time,update_time)<=72 THEN 1 ELSE 0 END)
  544. / NULLIF(COUNT(*),0),4) AS MetricValue
  545. FROM (
  546. SELECT tenant_id,
  547. 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,
  548. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f') AS update_time
  549. FROM mdp_stg_so
  550. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  551. AND source_table='ado_contract_review'
  552. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') LIKE 'UAT-PROC-%'
  553. ) t
  554. WHERE create_time IS NOT NULL AND update_time IS NOT NULL
  555. GROUP BY tenant_id
  556. UNION ALL
  557. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  558. 'S1_L2_009' AS MetricCode,
  559. ROUND(COUNT(*)/NULLIF(COUNT(DISTINCT NULLIF(owner_account,'')),0),4) AS MetricValue
  560. FROM (
  561. SELECT tenant_id,
  562. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ResponsibleAccount')),
  563. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))) AS owner_account
  564. FROM mdp_stg_so
  565. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  566. AND source_table='ado_contract_review'
  567. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.BillNo')),'') LIKE 'UAT-PROC-%'
  568. ) t
  569. GROUP BY tenant_id
  570. UNION ALL
  571. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  572. 'S1_L2_004' AS MetricCode,
  573. ROUND(AVG(cycle_hours / 24), 4) AS MetricValue
  574. FROM (
  575. SELECT tenant_id,
  576. COALESCE(
  577. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingDesignCycle')) REGEXP '^-?[0-9]+$'
  578. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingDesignCycle')) AS DECIMAL(18,6)) END,
  579. TIMESTAMPDIFF(HOUR,
  580. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualStart')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  581. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualEnd')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  582. )
  583. ) AS cycle_hours
  584. FROM mdp_stg_so
  585. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  586. AND source_table = 'ado_product_design'
  587. ) t
  588. WHERE cycle_hours IS NOT NULL AND cycle_hours >= 0
  589. GROUP BY tenant_id
  590. UNION ALL
  591. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  592. 'S1_L2_005' AS MetricCode,
  593. ROUND(100 * SUM(CASE WHEN actual_end <= plan_end THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  594. FROM (
  595. SELECT tenant_id,
  596. 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,
  597. 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
  598. FROM mdp_stg_so
  599. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  600. AND source_table = 'ado_product_design'
  601. ) t
  602. WHERE plan_end IS NOT NULL AND actual_end IS NOT NULL
  603. GROUP BY tenant_id
  604. UNION ALL
  605. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  606. 'S1_L2_006' AS MetricCode,
  607. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(design_lead, '')), 1), 4) AS MetricValue
  608. FROM (
  609. SELECT tenant_id,
  610. COALESCE(
  611. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DesignLeadAccount')),
  612. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DesignLeadName')),
  613. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))
  614. ) AS design_lead
  615. FROM mdp_stg_so
  616. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  617. AND source_table = 'ado_product_design'
  618. ) t
  619. GROUP BY tenant_id
  620. UNION ALL
  621. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  622. CASE stage_no
  623. WHEN 1 THEN 'S1_L3_001'
  624. WHEN 2 THEN 'S1_L3_002'
  625. WHEN 3 THEN 'S1_L3_003'
  626. WHEN 4 THEN 'S1_L3_004'
  627. WHEN 5 THEN 'S1_L3_005'
  628. END AS MetricCode,
  629. ROUND(AVG(cycle_hours), 4) AS MetricValue
  630. FROM (
  631. SELECT tenant_id,
  632. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) REGEXP '^-?[0-9]+$'
  633. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) END AS stage_no,
  634. COALESCE(
  635. TIMESTAMPDIFF(HOUR,
  636. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  637. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  638. ),
  639. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  640. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  641. ) AS cycle_hours
  642. FROM mdp_stg_so
  643. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  644. AND source_table = 'ado_contract_review_flow'
  645. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-PROC-%'
  646. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-DEL-%'
  647. ) t
  648. WHERE stage_no BETWEEN 1 AND 5
  649. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  650. GROUP BY tenant_id, stage_no
  651. UNION ALL
  652. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  653. CASE stage_no
  654. WHEN 1 THEN 'S1_L3_101'
  655. WHEN 2 THEN 'S1_L3_102'
  656. WHEN 3 THEN 'S1_L3_103'
  657. WHEN 4 THEN 'S1_L3_104'
  658. WHEN 5 THEN 'S1_L3_105'
  659. END AS MetricCode,
  660. ROUND(AVG(cycle_hours), 4) AS MetricValue
  661. FROM (
  662. SELECT tenant_id,
  663. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) REGEXP '^-?[0-9]+$'
  664. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) END AS stage_no,
  665. COALESCE(
  666. TIMESTAMPDIFF(HOUR,
  667. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  668. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  669. ),
  670. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  671. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  672. ) AS cycle_hours
  673. FROM mdp_stg_so
  674. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  675. AND source_table = 'ado_contract_review_flow'
  676. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-PROC-%'
  677. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-DEL-%'
  678. ) t
  679. WHERE stage_no BETWEEN 1 AND 5
  680. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  681. GROUP BY tenant_id, stage_no
  682. UNION ALL
  683. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  684. metric_code AS MetricCode,
  685. ROUND(AVG(cycle_hours), 4) AS MetricValue
  686. FROM (
  687. SELECT tenant_id,
  688. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')), ''), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  689. WHEN 'LAW' THEN 'S1_L4_001'
  690. WHEN 'PRE_SALES' THEN 'S1_L4_002'
  691. WHEN 'MPS' THEN 'S1_L4_003'
  692. WHEN 'TEST' THEN 'S1_L4_004'
  693. WHEN '法律事务部' THEN 'S1_L4_001'
  694. WHEN '技术售前组' THEN 'S1_L4_002'
  695. WHEN '综合主计划' THEN 'S1_L4_003'
  696. WHEN '试验站' THEN 'S1_L4_004'
  697. END AS metric_code,
  698. COALESCE(
  699. TIMESTAMPDIFF(HOUR,
  700. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  701. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  702. ),
  703. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  704. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  705. ) AS cycle_hours
  706. FROM mdp_stg_so
  707. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  708. AND source_table = 'ado_contract_review_flow'
  709. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) = '1'
  710. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-PROC-%'
  711. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-DEL-%'
  712. ) t
  713. WHERE metric_code IS NOT NULL
  714. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  715. GROUP BY tenant_id, metric_code
  716. UNION ALL
  717. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  718. metric_code AS MetricCode,
  719. ROUND(AVG(cycle_hours), 4) AS MetricValue
  720. FROM (
  721. SELECT tenant_id,
  722. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')), ''), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  723. WHEN 'LAW' THEN 'S1_L4_101'
  724. WHEN 'PRE_SALES' THEN 'S1_L4_102'
  725. WHEN 'MPS' THEN 'S1_L4_103'
  726. WHEN 'TEST' THEN 'S1_L4_104'
  727. WHEN '法律事务部' THEN 'S1_L4_101'
  728. WHEN '技术售前组' THEN 'S1_L4_102'
  729. WHEN '综合主计划' THEN 'S1_L4_103'
  730. WHEN '试验站' THEN 'S1_L4_104'
  731. END AS metric_code,
  732. COALESCE(
  733. TIMESTAMPDIFF(HOUR,
  734. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  735. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  736. ),
  737. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  738. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  739. ) AS cycle_hours
  740. FROM mdp_stg_so
  741. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  742. AND source_table = 'ado_contract_review_flow'
  743. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) = '1'
  744. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-PROC-%'
  745. AND COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') NOT LIKE 'UAT-DEL-%'
  746. ) t
  747. WHERE metric_code IS NOT NULL
  748. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  749. GROUP BY tenant_id, metric_code
  750. UNION ALL
  751. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  752. CONCAT('S1_L3_',branch_code,LPAD(stage_no,2,'0')) AS MetricCode,
  753. ROUND(AVG(cycle_hours),4) AS MetricValue
  754. FROM (
  755. SELECT tenant_id,
  756. CASE
  757. WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-PROC-%' THEN '3'
  758. WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-DEL-%' THEN '4'
  759. END AS branch_code,
  760. CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) AS stage_no,
  761. TIMESTAMPDIFF(
  762. HOUR,
  763. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f'),
  764. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f')
  765. ) AS cycle_hours
  766. FROM mdp_stg_so
  767. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  768. AND source_table='ado_contract_review_flow'
  769. AND (COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-PROC-%'
  770. OR COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-DEL-%')
  771. ) t
  772. WHERE branch_code IS NOT NULL AND stage_no BETWEEN 1 AND 5
  773. AND cycle_hours IS NOT NULL AND cycle_hours>=0
  774. GROUP BY tenant_id,branch_code,stage_no
  775. UNION ALL
  776. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  777. CONCAT('S1_L4_',branch_code,dept_code) AS MetricCode,
  778. ROUND(AVG(cycle_hours),4) AS MetricValue
  779. FROM (
  780. SELECT tenant_id,
  781. CASE
  782. WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-PROC-%' THEN '3'
  783. WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-DEL-%' THEN '4'
  784. END AS branch_code,
  785. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')),''),
  786. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  787. WHEN 'LAW' THEN '01' WHEN '法律事务部' THEN '01'
  788. WHEN 'PRE_SALES' THEN '02' WHEN '技术售前组' THEN '02'
  789. WHEN 'MPS' THEN '03' WHEN '综合主计划' THEN '03'
  790. WHEN 'TEST' THEN '04' WHEN '试验站' THEN '04'
  791. END AS dept_code,
  792. TIMESTAMPDIFF(
  793. HOUR,
  794. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f'),
  795. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')),'null'),''),'T',' '),'%Y-%m-%d %H:%i:%s.%f')
  796. ) AS cycle_hours
  797. FROM mdp_stg_so
  798. WHERE tenant_id=@TenantId AND COALESCE(NULLIF(factory_id,0),1)=@FactoryId
  799. AND source_table='ado_contract_review_flow'
  800. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo'))='1'
  801. AND (COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-PROC-%'
  802. OR COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ReviewBillNo')),'') LIKE 'UAT-DEL-%')
  803. ) t
  804. WHERE branch_code IS NOT NULL AND dept_code IS NOT NULL
  805. AND cycle_hours IS NOT NULL AND cycle_hours>=0
  806. GROUP BY tenant_id,branch_code,dept_code
  807. """,
  808. new SugarParameter("@BatchId", batchId),
  809. new SugarParameter("@StatDate", statDate),
  810. new SugarParameter("@TenantId", scope.TenantId),
  811. new SugarParameter("@FactoryId", scope.FactoryId));
  812. }
  813. private async Task<int> UpsertS1KpiValueAsync(S1KpiCalcRow row, DateTime statDate, DateTime now)
  814. {
  815. var meta = await _db.Ado.SqlQuerySingleAsync<S1KpiMetaRow>(
  816. """
  817. SELECT MetricLevel, Direction, YellowThreshold, RedThreshold
  818. FROM ado_smart_ops_kpi_master
  819. WHERE TenantId=@TenantId AND ModuleCode='S1' AND MetricCode=@MetricCode AND IsEnabled=1
  820. LIMIT 1
  821. """,
  822. new SugarParameter("@TenantId", row.TenantId),
  823. new SugarParameter("@MetricCode", row.MetricCode));
  824. if (meta == null || row.MetricValue == null)
  825. return 0;
  826. var table = ResolveKpiValueTable(meta.MetricLevel);
  827. var current = await _db.Ado.SqlQuerySingleAsync<S1KpiValueRow>(
  828. $"""
  829. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  830. FROM {table}
  831. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  832. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  833. ORDER BY id
  834. LIMIT 1
  835. """,
  836. new SugarParameter("@TenantId", row.TenantId),
  837. new SugarParameter("@FactoryId", row.FactoryId),
  838. new SugarParameter("@MetricCode", row.MetricCode),
  839. new SugarParameter("@BizDate", statDate));
  840. var prior = await _db.Ado.SqlQuerySingleAsync<S1KpiValueRow>(
  841. $"""
  842. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  843. FROM {table}
  844. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  845. AND metric_code=@MetricCode AND biz_date<@BizDate AND is_deleted=0
  846. ORDER BY biz_date DESC, id DESC
  847. LIMIT 1
  848. """,
  849. new SugarParameter("@TenantId", row.TenantId),
  850. new SugarParameter("@FactoryId", row.FactoryId),
  851. new SugarParameter("@MetricCode", row.MetricCode),
  852. new SugarParameter("@BizDate", statDate));
  853. var actual = Math.Round(row.MetricValue.Value, 4);
  854. var snap = await _kpiTargetResolver.ResolveAsync(row.TenantId, row.FactoryId, row.MetricCode, "S1", statDate);
  855. var target = snap.TargetValue;
  856. var status = AidopS4KpiMerge.AchievementLevel(actual, target, meta.Direction, meta.YellowThreshold, meta.RedThreshold);
  857. var trend = ResolveTrendFlag(actual, prior?.MetricValue);
  858. if (current != null)
  859. {
  860. return await _db.Ado.ExecuteCommandAsync(
  861. $"""
  862. UPDATE {table}
  863. SET metric_value=@MetricValue, target_value=@TargetValue, status_color=@StatusColor, trend_flag=@TrendFlag,
  864. target_config_id=@TargetConfigId, target_source=@TargetSource, target_resolved_at=@TargetResolvedAt,
  865. is_active=1, status='ACTIVE', calc_time=@CalcTime, update_time=@CalcTime
  866. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  867. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  868. """,
  869. new SugarParameter("@MetricValue", actual),
  870. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  871. new SugarParameter("@StatusColor", status),
  872. new SugarParameter("@TrendFlag", trend),
  873. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  874. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  875. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt),
  876. new SugarParameter("@CalcTime", now),
  877. new SugarParameter("@TenantId", row.TenantId),
  878. new SugarParameter("@FactoryId", row.FactoryId),
  879. new SugarParameter("@MetricCode", row.MetricCode),
  880. new SugarParameter("@BizDate", statDate));
  881. }
  882. var nextId = Yitter.IdGenerator.YitIdHelper.NextId();
  883. return await _db.Ado.ExecuteCommandAsync(
  884. $"""
  885. INSERT INTO {table}
  886. (id, tenant_id, factory_id, status, biz_date, create_time, update_time, is_deleted, is_active,
  887. module_code, metric_code, metric_value, target_value, status_color, trend_flag, calc_time,
  888. target_config_id, target_source, target_resolved_at)
  889. VALUES
  890. (@Id, @TenantId, @FactoryId, 'ACTIVE', @BizDate, @CalcTime, @CalcTime, 0, 1,
  891. 'S1', @MetricCode, @MetricValue, @TargetValue, @StatusColor, @TrendFlag, @CalcTime,
  892. @TargetConfigId, @TargetSource, @TargetResolvedAt)
  893. """,
  894. new SugarParameter("@Id", nextId),
  895. new SugarParameter("@TenantId", row.TenantId),
  896. new SugarParameter("@FactoryId", row.FactoryId),
  897. new SugarParameter("@BizDate", statDate),
  898. new SugarParameter("@CalcTime", now),
  899. new SugarParameter("@MetricCode", row.MetricCode),
  900. new SugarParameter("@MetricValue", actual),
  901. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  902. new SugarParameter("@StatusColor", status),
  903. new SugarParameter("@TrendFlag", trend),
  904. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  905. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  906. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt));
  907. }
  908. private IEnumerable<S1MdpSqlCommand> BuildStandardCommands(S1MdpRunScope scope, string batchId, DateTime now)
  909. {
  910. yield return Cmd(
  911. """
  912. INSERT INTO mdp_std_so
  913. (tenant_id, factory_id, company_id, source_system, order_id, order_entry_id, order_no, order_line, order_type,
  914. customer_id, customer_no, customer_name, customer_order_no, country, item_code, item_name, item_spec,
  915. map_number, map_name, bom_number, unit, order_qty, delivered_notice_qty, delivered_qty, price, tax_price,
  916. amount, total_amount, order_date, customer_request_date, plan_delivery_date, promised_delivery_date,
  917. capacity_date, material_ready_date, planner_no, planner_name, order_status, review_status, review_stage,
  918. flow_state, progress, urgent, closed, deleted_flag, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  919. SELECT
  920. COALESCE(NULLIF(e.tenant_id, 0), NULLIF(h.tenant_id, 0),
  921. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  922. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) AS SIGNED) END),
  923. COALESCE(NULLIF(NULLIF(e.factory_id, 0), e.tenant_id),
  924. NULLIF(NULLIF(h.factory_id, 0), h.tenant_id)),
  925. NULLIF(e.company_id, 0),
  926. 'AIDOP',
  927. 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,
  928. 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,
  929. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no')), e.source_biz_key),
  930. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.entry_seq')), JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.Id'))) AS CHAR),
  931. CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.order_type')) AS CHAR),
  932. 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,
  933. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_no')),
  934. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_name')),
  935. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.custom_order_bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_from'))),
  936. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.country')),
  937. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.item_number')),
  938. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.item_name')),
  939. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.specification')),
  940. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.map_number')),
  941. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.map_name')),
  942. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bom_number')),
  943. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.unit')),
  944. 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,
  945. 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,
  946. 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,
  947. 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,
  948. 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,
  949. 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,
  950. 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,
  951. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.date')), 'null'), ''),
  952. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.rdate')), 'null'), ''),
  953. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.plan_date')), 'null'), ''),
  954. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.date')), 'null'), ''),
  955. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.sys_capacity_date')), 'null'), ''),
  956. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.sys_material_date')), 'null'), ''),
  957. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.planner_no')),
  958. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.planner_name')),
  959. CASE WHEN JSON_EXTRACT(h.raw_data,'$.closed') IN (1, true) THEN 'CLOSED' ELSE 'OPEN' END,
  960. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.FlowStatus')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.flowstate'))),
  961. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.CurrentDept')), JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.CurrentStage'))),
  962. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.flowstate')),
  963. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.progress')),
  964. 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,
  965. 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,
  966. 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),
  967. e.source_table,
  968. e.source_row_id,
  969. e.source_biz_key,
  970. @BatchId,
  971. @Now
  972. FROM mdp_stg_so e
  973. LEFT JOIN mdp_stg_so h ON h.source_table='crm_seorder'
  974. AND h.tenant_id = e.tenant_id
  975. AND JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.Id')) = JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.seorder_id'))
  976. LEFT JOIN mdp_stg_so r ON r.source_table='ado_contract_review'
  977. AND r.tenant_id = e.tenant_id
  978. 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')))
  979. WHERE e.source_table='crm_seorderentry'
  980. AND e.tenant_id=@TenantId
  981. AND COALESCE(NULLIF(e.factory_id, 0), 1)=@FactoryId
  982. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no'))), '') <> ''
  983. AND COALESCE(NULLIF(e.tenant_id, 0), NULLIF(h.tenant_id, 0),
  984. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  985. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  986. ON DUPLICATE KEY UPDATE
  987. factory_id=VALUES(factory_id), customer_no=VALUES(customer_no), customer_name=VALUES(customer_name), item_code=VALUES(item_code),
  988. item_name=VALUES(item_name), item_spec=VALUES(item_spec), order_qty=VALUES(order_qty),
  989. delivered_notice_qty=VALUES(delivered_notice_qty), delivered_qty=VALUES(delivered_qty),
  990. plan_delivery_date=VALUES(plan_delivery_date), promised_delivery_date=VALUES(promised_delivery_date),
  991. capacity_date=VALUES(capacity_date), material_ready_date=VALUES(material_ready_date),
  992. order_status=VALUES(order_status), review_status=VALUES(review_status), review_stage=VALUES(review_stage),
  993. flow_state=VALUES(flow_state), progress=VALUES(progress), deleted_flag=VALUES(deleted_flag),
  994. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), 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, company_id, source_system, trans_type, plan_id, plan_no, plan_line, order_id, order_entry_id,
  1000. order_no, order_line, customer_no, customer_name, country, item_code, item_name, item_spec, qty, plan_qty,
  1001. weight, volume, order_date, plan_ship_date, shipping_site, shipping_address, consignee, telephone,
  1002. status, confirm_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  1003. SELECT
  1004. COALESCE(NULLIF(d.tenant_id, 0),
  1005. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1006. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1007. NULLIF(d.factory_id, 0),
  1008. NULLIF(d.company_id, 0),
  1009. 'AIDOP',
  1010. 'SHIP_PLAN',
  1011. 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,
  1012. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.LotSerial')), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id')) AS CHAR)),
  1013. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS CHAR),
  1014. 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,
  1015. 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,
  1016. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr'))),
  1017. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sentry_id')) AS CHAR),
  1018. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CustomNo')),
  1019. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CustomName')),
  1020. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Country')),
  1021. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')),
  1022. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemName')),
  1023. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Specification')),
  1024. 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,
  1025. 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,
  1026. 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,
  1027. 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,
  1028. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdDate')), 'null'), ''),
  1029. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingDate')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CreateTime'))), 'null'), ''),
  1030. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingSite')),
  1031. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingAddress')),
  1032. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Consignee')),
  1033. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Telephone')),
  1034. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Status')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Status'))),
  1035. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.IsConfirm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.IsConfirm'))) AS CHAR),
  1036. d.source_table,
  1037. d.source_row_id,
  1038. d.source_biz_key,
  1039. @BatchId,
  1040. @Now
  1041. FROM mdp_stg_ship_trans d
  1042. LEFT JOIN mdp_stg_ship_trans m ON m.source_table='ShippingPlan'
  1043. AND m.tenant_id = d.tenant_id
  1044. AND JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.RecID')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id'))
  1045. WHERE d.source_table='ShippingPlanDetail'
  1046. AND d.tenant_id=@TenantId
  1047. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1048. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr'))), '') <> ''
  1049. AND COALESCE(NULLIF(d.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) IS NOT NULL
  1052. ON DUPLICATE KEY UPDATE
  1053. factory_id=VALUES(factory_id), plan_no=VALUES(plan_no), order_no=VALUES(order_no), customer_no=VALUES(customer_no), customer_name=VALUES(customer_name),
  1054. item_code=VALUES(item_code), item_name=VALUES(item_name), qty=VALUES(qty), plan_qty=VALUES(plan_qty),
  1055. plan_ship_date=VALUES(plan_ship_date), shipping_site=VALUES(shipping_site), shipping_address=VALUES(shipping_address),
  1056. status=VALUES(status), confirm_status=VALUES(confirm_status), sync_batch_id=VALUES(sync_batch_id),
  1057. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  1058. """, scope, batchId, now);
  1059. yield return Cmd(
  1060. """
  1061. INSERT INTO mdp_std_ship_trans
  1062. (tenant_id, factory_id, source_system, trans_type, shipper_rec_id, shipper_no, shipper_line, order_no, order_line,
  1063. customer_no, item_code, item_name, qty_to_ship, picking_qty, real_qty, gross_weight, net_weight, volume,
  1064. plan_ship_date, actual_ship_date, site, status, confirm_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  1065. SELECT
  1066. COALESCE(NULLIF(d.tenant_id, 0),
  1067. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1068. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1069. NULLIF(d.factory_id, 0),
  1070. 'AIDOP',
  1071. 'ASN_SHIPPER',
  1072. 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,
  1073. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Id')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Id'))),
  1074. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS CHAR),
  1075. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.OrdNbr'))),
  1076. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdLine')) AS CHAR),
  1077. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.SoldTo')),
  1078. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ContainerItem')),
  1079. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Descr')),
  1080. 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,
  1081. 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,
  1082. 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,
  1083. 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,
  1084. 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,
  1085. 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,
  1086. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ShipDate')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShipDate'))), 'null'), ''),
  1087. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ShipDate')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShipDate'))), 'null'), ''),
  1088. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Site')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Site'))),
  1089. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Status')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Status'))),
  1090. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.IsConfirm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.IsConfirm'))) AS CHAR),
  1091. d.source_table,
  1092. d.source_row_id,
  1093. d.source_biz_key,
  1094. @BatchId,
  1095. @Now
  1096. FROM mdp_stg_ship_trans d
  1097. LEFT JOIN mdp_stg_ship_trans m ON m.source_table='ASNBOLShipperMaster'
  1098. AND m.tenant_id = d.tenant_id
  1099. AND JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.RecID')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ASNBOLShipperRecID'))
  1100. WHERE d.source_table='ASNBOLShipperDetail'
  1101. AND d.tenant_id=@TenantId
  1102. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1103. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Id')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Id'))), '') <> ''
  1104. AND COALESCE(NULLIF(d.tenant_id, 0),
  1105. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1106. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1107. ON DUPLICATE KEY UPDATE
  1108. factory_id=VALUES(factory_id), shipper_no=VALUES(shipper_no), order_no=VALUES(order_no), order_line=VALUES(order_line),
  1109. customer_no=VALUES(customer_no), item_code=VALUES(item_code), item_name=VALUES(item_name),
  1110. qty_to_ship=VALUES(qty_to_ship), picking_qty=VALUES(picking_qty), real_qty=VALUES(real_qty),
  1111. actual_ship_date=VALUES(actual_ship_date), site=VALUES(site), status=VALUES(status),
  1112. confirm_status=VALUES(confirm_status), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  1113. update_time=CURRENT_TIMESTAMP
  1114. """, scope, batchId, now);
  1115. yield return Cmd(
  1116. """
  1117. INSERT INTO mdp_std_ship_trans
  1118. (tenant_id, factory_id, source_system, trans_type, order_no, customer_no, item_code, item_name, qty,
  1119. plan_ship_date, status, linkage_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  1120. SELECT
  1121. COALESCE(NULLIF(d.tenant_id, 0),
  1122. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1123. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1124. NULLIF(d.factory_id, 0),
  1125. 'AIDOP',
  1126. 'LINKAGE_PLAN',
  1127. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')),
  1128. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.custom_no')),
  1129. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_number')),
  1130. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Descr')),
  1131. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) AS DECIMAL(18,6)) END,
  1132. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sys_capacity_date')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.fystarttime'))), 'null'), ''),
  1133. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.type')),
  1134. CASE
  1135. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.isuse')) = '1' THEN 'LINKED'
  1136. ELSE 'INACTIVE'
  1137. END,
  1138. d.source_table,
  1139. d.source_row_id,
  1140. d.source_biz_key,
  1141. @BatchId,
  1142. @Now
  1143. FROM mdp_stg_ship_trans d
  1144. WHERE d.source_table='LinkagePlan'
  1145. AND d.tenant_id=@TenantId
  1146. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1147. AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), '') <> ''
  1148. AND COALESCE(NULLIF(d.tenant_id, 0),
  1149. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1150. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1151. ON DUPLICATE KEY UPDATE
  1152. factory_id=VALUES(factory_id), order_no=VALUES(order_no), customer_no=VALUES(customer_no), item_code=VALUES(item_code),
  1153. item_name=VALUES(item_name), qty=VALUES(qty), plan_ship_date=VALUES(plan_ship_date),
  1154. status=VALUES(status), linkage_status=VALUES(linkage_status), sync_batch_id=VALUES(sync_batch_id),
  1155. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  1156. """, scope, batchId, now);
  1157. }
  1158. private IEnumerable<S1MdpSqlCommand> BuildDwdCommands(S1MdpRunScope scope, string batchId, DateTime now)
  1159. {
  1160. yield return Cmd(
  1161. """
  1162. INSERT INTO dwd_requirement_examine_detail
  1163. (tenant_id, factory_id, stat_date, row_id, parent_row_id, examine_id, order_entry_id, bill_no, morder_no,
  1164. num, item_number, item_name, bom_number, model, kitting_time, item_type, erp_cls_name, qty, wastage,
  1165. need_count, sqty, use_qty, self_lack_qty, lack_qty, mo_qty, make_qty, purchase_qty, purchase_occupy_qty,
  1166. satisfy_time, have_ic_subs, substitute_code, create_time, source_system, sync_batch_id, calc_batch_id, calc_time)
  1167. SELECT
  1168. COALESCE(NULLIF(d.tenant_id, 0), NULLIF(r.tenant_id, 0), NULLIF(so.tenant_id, 0),
  1169. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1170. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END),
  1171. COALESCE(NULLIF(d.factory_id, 0), NULLIF(r.factory_id, 0), NULLIF(so.factory_id, 0)),
  1172. @StatDate,
  1173. CAST(d.source_row_id AS SIGNED),
  1174. CASE
  1175. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.num')) = '1' THEN NULL
  1176. ELSE p.parent_row_id
  1177. END,
  1178. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id')) REGEXP '^-?[0-9]+$'
  1179. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id')) AS SIGNED) END,
  1180. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$'
  1181. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) AS SIGNED) END,
  1182. COALESCE(so.order_no, JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.bill_no'))),
  1183. JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.morder_no')),
  1184. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.num')) AS CHAR),
  1185. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_number')),
  1186. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_name')),
  1187. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bom_number')),
  1188. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.model')),
  1189. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.kitting_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  1190. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.type')) = '1' THEN '替代件' ELSE '标准件' END,
  1191. COALESCE(
  1192. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls_name')),
  1193. CASE JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls'))
  1194. WHEN '0' THEN '配置类'
  1195. WHEN '1' THEN '自制'
  1196. WHEN '2' THEN '委外加工'
  1197. WHEN '3' THEN '外购'
  1198. WHEN '4' THEN '虚拟件'
  1199. END
  1200. ),
  1201. 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,
  1202. 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,
  1203. 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,
  1204. 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,
  1205. 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,
  1206. 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,
  1207. 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,
  1208. 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,
  1209. CASE
  1210. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) = '0' AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls')) = '1'
  1211. 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
  1212. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  1213. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) AS DECIMAL(18,4))
  1214. END,
  1215. 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,
  1216. 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,
  1217. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.satisfy_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  1218. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.haveicsubs')) = '1' THEN '是' ELSE '否' END,
  1219. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.substitute_code')),
  1220. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.create_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  1221. 'AIDOP',
  1222. d.sync_batch_id,
  1223. @BatchId,
  1224. @Now
  1225. FROM mdp_stg_so d
  1226. INNER JOIN mdp_stg_so r ON r.source_table='b_examine_result'
  1227. AND r.tenant_id = d.tenant_id
  1228. AND r.source_row_id = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id'))
  1229. LEFT JOIN mdp_std_so so ON so.tenant_id = COALESCE(r.tenant_id, d.tenant_id)
  1230. AND COALESCE(NULLIF(so.factory_id, 0), 1)=@FactoryId
  1231. AND so.order_entry_id = CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$'
  1232. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) AS SIGNED) END
  1233. LEFT JOIN (
  1234. SELECT examine_id, MIN(source_row_id) AS parent_row_id
  1235. FROM (
  1236. SELECT JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.examine_id')) AS examine_id, CAST(source_row_id AS SIGNED) AS source_row_id
  1237. FROM mdp_stg_so
  1238. WHERE source_table='b_bom_child_examine'
  1239. AND tenant_id=@TenantId
  1240. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1241. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.num')) = '1'
  1242. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.is_use')) IN ('1', 'true', 'True', 'base64:type16:AQ==')
  1243. ) x
  1244. GROUP BY examine_id
  1245. ) p ON p.examine_id = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id'))
  1246. WHERE d.source_table='b_bom_child_examine'
  1247. AND d.tenant_id=@TenantId
  1248. AND COALESCE(NULLIF(d.factory_id, 0), 1)=@FactoryId
  1249. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.is_use')) IN ('1', 'true', 'True', 'base64:type16:AQ==')
  1250. AND COALESCE(NULLIF(d.tenant_id, 0), NULLIF(r.tenant_id, 0), NULLIF(so.tenant_id, 0),
  1251. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  1252. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END) IS NOT NULL
  1253. ON DUPLICATE KEY UPDATE
  1254. parent_row_id=VALUES(parent_row_id), order_entry_id=VALUES(order_entry_id), bill_no=VALUES(bill_no),
  1255. morder_no=VALUES(morder_no), num=VALUES(num), item_number=VALUES(item_number), item_name=VALUES(item_name),
  1256. bom_number=VALUES(bom_number), model=VALUES(model), kitting_time=VALUES(kitting_time),
  1257. item_type=VALUES(item_type), erp_cls_name=VALUES(erp_cls_name), qty=VALUES(qty), wastage=VALUES(wastage),
  1258. need_count=VALUES(need_count), sqty=VALUES(sqty), use_qty=VALUES(use_qty), self_lack_qty=VALUES(self_lack_qty),
  1259. lack_qty=VALUES(lack_qty), mo_qty=VALUES(mo_qty), make_qty=VALUES(make_qty), purchase_qty=VALUES(purchase_qty),
  1260. purchase_occupy_qty=VALUES(purchase_occupy_qty), satisfy_time=VALUES(satisfy_time), have_ic_subs=VALUES(have_ic_subs),
  1261. substitute_code=VALUES(substitute_code), create_time=VALUES(create_time), sync_batch_id=VALUES(sync_batch_id),
  1262. calc_time=VALUES(calc_time), update_time=CURRENT_TIMESTAMP
  1263. """, scope, batchId, now);
  1264. yield return Cmd(
  1265. """
  1266. INSERT INTO dwd_ship_trans
  1267. (tenant_id, factory_id, company_id, stat_date, order_id, order_entry_id, order_no, order_line, customer_no,
  1268. customer_name, country, item_code, item_name, item_spec, order_qty, planned_ship_qty, shipped_qty,
  1269. remaining_qty, order_date, customer_request_date, plan_delivery_date, promised_delivery_date,
  1270. plan_ship_date, actual_ship_date, review_status, order_status, delivery_status, linkage_status, risk_level,
  1271. source_system, source_table, source_row_id, source_biz_key, sync_batch_id, calc_batch_id, calc_time)
  1272. SELECT
  1273. so.tenant_id,
  1274. so.factory_id,
  1275. so.company_id,
  1276. @StatDate,
  1277. so.order_id,
  1278. so.order_entry_id,
  1279. so.order_no,
  1280. IFNULL(so.order_line, ''),
  1281. so.customer_no,
  1282. so.customer_name,
  1283. so.country,
  1284. IFNULL(so.item_code, ''),
  1285. so.item_name,
  1286. so.item_spec,
  1287. IFNULL(so.order_qty, 0),
  1288. IFNULL(p.plan_qty, 0),
  1289. IFNULL(a.real_qty, 0),
  1290. GREATEST(IFNULL(so.order_qty, 0) - IFNULL(a.real_qty, 0), 0),
  1291. so.order_date,
  1292. so.customer_request_date,
  1293. so.plan_delivery_date,
  1294. so.promised_delivery_date,
  1295. p.plan_ship_date,
  1296. a.actual_ship_date,
  1297. so.review_status,
  1298. so.order_status,
  1299. CASE
  1300. WHEN IFNULL(so.order_qty, 0) > 0 AND IFNULL(a.real_qty, 0) >= IFNULL(so.order_qty, 0) THEN 'COMPLETED'
  1301. WHEN COALESCE(p.plan_ship_date, so.promised_delivery_date, so.plan_delivery_date) < @Now THEN 'DELAYED'
  1302. WHEN IFNULL(p.plan_qty, 0) > 0 THEN 'PLANNED'
  1303. ELSE 'OPEN'
  1304. END,
  1305. l.linkage_status,
  1306. CASE
  1307. WHEN COALESCE(p.plan_ship_date, so.promised_delivery_date, so.plan_delivery_date) < @Now
  1308. AND IFNULL(a.real_qty, 0) < IFNULL(so.order_qty, 0) THEN 'HIGH'
  1309. WHEN IFNULL(a.real_qty, 0) < IFNULL(so.order_qty, 0) THEN 'MEDIUM'
  1310. ELSE 'LOW'
  1311. END,
  1312. 'AIDOP',
  1313. so.source_table,
  1314. so.source_row_id,
  1315. so.source_biz_key,
  1316. so.sync_batch_id,
  1317. @BatchId,
  1318. @Now
  1319. FROM mdp_std_so so
  1320. LEFT JOIN (
  1321. SELECT tenant_id, order_no, order_entry_id, IFNULL(order_line, '') AS order_line, IFNULL(item_code, '') AS item_code,
  1322. SUM(IFNULL(plan_qty, IFNULL(qty, 0))) AS plan_qty,
  1323. MIN(plan_ship_date) AS plan_ship_date
  1324. FROM mdp_std_ship_trans
  1325. WHERE trans_type='SHIP_PLAN'
  1326. AND tenant_id=@TenantId
  1327. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1328. GROUP BY tenant_id, order_no, order_entry_id, IFNULL(order_line, ''), IFNULL(item_code, '')
  1329. ) p ON so.tenant_id=p.tenant_id
  1330. AND so.order_no=p.order_no
  1331. AND IFNULL(so.item_code, '')=p.item_code
  1332. AND (
  1333. (p.order_entry_id IS NOT NULL AND so.order_entry_id=p.order_entry_id)
  1334. OR (p.order_entry_id IS NULL AND IFNULL(so.order_line, '')=p.order_line)
  1335. )
  1336. LEFT JOIN (
  1337. SELECT tenant_id, order_no, IFNULL(order_line, '') AS order_line, IFNULL(item_code, '') AS item_code,
  1338. SUM(IFNULL(real_qty, IFNULL(qty_to_ship, 0))) AS real_qty,
  1339. MAX(actual_ship_date) AS actual_ship_date
  1340. FROM mdp_std_ship_trans
  1341. WHERE trans_type='ASN_SHIPPER'
  1342. AND tenant_id=@TenantId
  1343. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1344. GROUP BY tenant_id, order_no, IFNULL(order_line, ''), IFNULL(item_code, '')
  1345. ) 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
  1346. LEFT JOIN (
  1347. SELECT tenant_id, order_no, item_code, MAX(linkage_status) AS linkage_status
  1348. FROM mdp_std_ship_trans
  1349. WHERE IFNULL(linkage_status, '') <> ''
  1350. AND tenant_id=@TenantId
  1351. AND COALESCE(NULLIF(factory_id, 0), 1)=@FactoryId
  1352. GROUP BY tenant_id, order_no, item_code
  1353. ) l ON so.tenant_id=l.tenant_id AND so.order_no=l.order_no AND IFNULL(so.item_code, '')=IFNULL(l.item_code, '')
  1354. WHERE so.tenant_id=@TenantId
  1355. AND COALESCE(NULLIF(so.factory_id, 0), 1)=@FactoryId
  1356. AND IFNULL(so.order_no, '') <> ''
  1357. ON DUPLICATE KEY UPDATE
  1358. factory_id=VALUES(factory_id), customer_no=VALUES(customer_no), customer_name=VALUES(customer_name), item_name=VALUES(item_name),
  1359. order_qty=VALUES(order_qty), planned_ship_qty=VALUES(planned_ship_qty), shipped_qty=VALUES(shipped_qty),
  1360. remaining_qty=VALUES(remaining_qty), plan_ship_date=VALUES(plan_ship_date), actual_ship_date=VALUES(actual_ship_date),
  1361. review_status=VALUES(review_status), order_status=VALUES(order_status), delivery_status=VALUES(delivery_status),
  1362. linkage_status=VALUES(linkage_status), risk_level=VALUES(risk_level), sync_batch_id=VALUES(sync_batch_id),
  1363. calc_batch_id=VALUES(calc_batch_id), calc_time=VALUES(calc_time), update_time=CURRENT_TIMESTAMP
  1364. """, scope, batchId, now);
  1365. }
  1366. private async Task<long> InsertSyncLogAsync(long tenantId, long entityId, string entityName, string batchId, int rowsRead)
  1367. {
  1368. await _db.Ado.ExecuteCommandAsync(
  1369. """
  1370. INSERT INTO mdp_sync_log
  1371. (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type, sync_start, rows_read, status)
  1372. VALUES (@TenantId, @EntityId, 'AIDOPDEV_MYSQL', @EntityName, @BatchId, 'FULL', 'AUTO', NOW(), @RowsRead, 'RUNNING')
  1373. """,
  1374. new SugarParameter("@TenantId", tenantId),
  1375. new SugarParameter("@EntityId", entityId),
  1376. new SugarParameter("@EntityName", entityName),
  1377. new SugarParameter("@BatchId", batchId),
  1378. new SugarParameter("@RowsRead", rowsRead));
  1379. return await _db.Ado.GetLongAsync(
  1380. "SELECT id FROM mdp_sync_log WHERE sync_batch_id=@BatchId AND entity_id=@EntityId ORDER BY id DESC LIMIT 1",
  1381. new List<SugarParameter>
  1382. {
  1383. new("@BatchId", batchId),
  1384. new("@EntityId", entityId)
  1385. });
  1386. }
  1387. private async Task MarkSyncLogSuccessAsync(long logId, DateTime started, int affected)
  1388. {
  1389. await _db.Ado.ExecuteCommandAsync(
  1390. """
  1391. UPDATE mdp_sync_log
  1392. SET sync_end=NOW(), duration_ms=@DurationMs, rows_insert=@RowsInsert, rows_update=0, rows_skip=0, rows_error=0, status='SUCCESS'
  1393. WHERE id=@Id
  1394. """,
  1395. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  1396. new SugarParameter("@RowsInsert", affected),
  1397. new SugarParameter("@Id", logId));
  1398. }
  1399. private async Task MarkSyncLogFailedAsync(long logId, DateTime started, string message)
  1400. {
  1401. try
  1402. {
  1403. await _db.Ado.ExecuteCommandAsync(
  1404. """
  1405. UPDATE mdp_sync_log
  1406. SET sync_end=NOW(), duration_ms=@DurationMs, rows_error=1, status='FAILED', error_msg=@ErrorMsg
  1407. WHERE id=@Id
  1408. """,
  1409. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  1410. new SugarParameter("@ErrorMsg", Truncate(message, 1000)),
  1411. new SugarParameter("@Id", logId));
  1412. }
  1413. catch (Exception ex)
  1414. {
  1415. // 写库自身失败兜底:避免再抛掩盖原异常;遗留 RUNNING 行可由运维手动清理
  1416. Console.Error.WriteLine($"[S1MdpSyncTransform] MarkSyncLogFailed write failed (syncLogId={logId}): {ex.Message}");
  1417. }
  1418. }
  1419. private async Task<long> InsertTransformRunLogAsync(S1MdpRunScope scope, string batchId, DateTime startedAt, string triggerType)
  1420. {
  1421. await _db.Ado.ExecuteCommandAsync(
  1422. """
  1423. INSERT INTO mdp_transform_run_log
  1424. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  1425. VALUES (@TenantId, @JobCode, 'S1 MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  1426. """,
  1427. new SugarParameter("@TenantId", scope.TenantId),
  1428. new SugarParameter("@JobCode", JobCode),
  1429. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  1430. new SugarParameter("@BatchId", batchId),
  1431. new SugarParameter("@StartTime", startedAt));
  1432. return await _db.Ado.GetLongAsync(
  1433. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  1434. new List<SugarParameter> { new("@BatchId", batchId) });
  1435. }
  1436. private async Task MarkTransformRunSuccessAsync(long runLogId, DateTime startedAt, S1MdpSyncTransformResult result)
  1437. {
  1438. var finishedAt = DateTime.Now;
  1439. await _db.Ado.ExecuteCommandAsync(
  1440. """
  1441. UPDATE mdp_transform_run_log
  1442. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  1443. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=@DwdRows,
  1444. summary_json=@SummaryJson, update_time=CURRENT_TIMESTAMP
  1445. WHERE id=@Id
  1446. """,
  1447. new SugarParameter("@EndTime", finishedAt),
  1448. new SugarParameter("@DurationMs", ToDurationMs(finishedAt - startedAt)),
  1449. new SugarParameter("@StageRows", result.StageRows),
  1450. new SugarParameter("@StandardRows", result.StandardRows),
  1451. new SugarParameter("@DwdRows", result.DwdRows),
  1452. new SugarParameter("@SummaryJson", BuildRunSummaryJson(result)),
  1453. new SugarParameter("@Id", runLogId));
  1454. }
  1455. private async Task MarkTransformRunFailedAsync(long runLogId, DateTime startedAt, string message)
  1456. {
  1457. try
  1458. {
  1459. var finishedAt = DateTime.Now;
  1460. var db = _db.CopyNew();
  1461. await db.Ado.ExecuteCommandAsync(
  1462. """
  1463. UPDATE mdp_transform_run_log
  1464. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  1465. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  1466. WHERE id=@Id
  1467. """,
  1468. new SugarParameter("@EndTime", finishedAt),
  1469. new SugarParameter("@DurationMs", ToDurationMs(finishedAt - startedAt)),
  1470. new SugarParameter("@ErrorMessage", Truncate(message, 2000)),
  1471. new SugarParameter("@Id", runLogId));
  1472. }
  1473. catch (Exception ex)
  1474. {
  1475. // 写库自身失败兜底(典型场景:远端 MySQL 瞬断导致 MarkFailed 自身也连不上):
  1476. // 避免再抛二次异常掩盖原错;遗留 RUNNING 行可由运维手动清理。
  1477. Console.Error.WriteLine($"[S1MdpSyncTransform] MarkTransformRunFailed write failed (runLogId={runLogId}): {ex.Message}");
  1478. }
  1479. }
  1480. private static S1MdpSqlCommand Cmd(string sql, S1MdpRunScope scope, string batchId, DateTime now)
  1481. {
  1482. return new S1MdpSqlCommand(sql, new[]
  1483. {
  1484. new SugarParameter("@BatchId", batchId),
  1485. new SugarParameter("@Now", now),
  1486. new SugarParameter("@StatDate", now.Date),
  1487. new SugarParameter("@TenantId", scope.TenantId),
  1488. new SugarParameter("@FactoryId", scope.FactoryId)
  1489. });
  1490. }
  1491. private static string BuildJsonObjectExpression(IEnumerable<string> columns)
  1492. {
  1493. var parts = columns.SelectMany(c => new[] { $"'{c.Replace("'", "''")}'", $"s.`{c}`" });
  1494. return $"JSON_OBJECT({string.Join(",", parts)})";
  1495. }
  1496. private static string BuildOptionalColumnExpr(IReadOnlyCollection<string> columns, string expected, string fallback)
  1497. {
  1498. return columns.Any(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase))
  1499. ? $"s.`{FindColumn(columns, expected)}`"
  1500. : fallback;
  1501. }
  1502. private static string FindColumn(IEnumerable<string> columns, string expected)
  1503. {
  1504. return columns.First(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase));
  1505. }
  1506. private static async Task KeepLeaseAliveAsync(IS1MdpFullRunLease lease, CancellationToken ct)
  1507. {
  1508. while (!ct.IsCancellationRequested)
  1509. {
  1510. try
  1511. {
  1512. await Task.Delay(TimeSpan.FromSeconds(30), ct);
  1513. await lease.HeartbeatAsync(CancellationToken.None);
  1514. }
  1515. catch (OperationCanceledException)
  1516. {
  1517. return;
  1518. }
  1519. }
  1520. }
  1521. private static int ToDurationMs(TimeSpan elapsed)
  1522. {
  1523. var ms = elapsed.TotalMilliseconds;
  1524. if (double.IsNaN(ms) || ms <= 0) return 0;
  1525. return ms >= int.MaxValue ? int.MaxValue : (int)ms;
  1526. }
  1527. private static string NormalizeTriggerType(string? triggerType)
  1528. {
  1529. return string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  1530. }
  1531. private static string BuildRunSummaryJson(S1MdpSyncTransformResult result)
  1532. {
  1533. return $$"""{"batchId":"{{result.BatchId}}","stageRows":{{result.StageRows}},"standardRows":{{result.StandardRows}},"dwdRows":{{result.DwdRows}},"kpiRows":{{result.KpiRows}}}""";
  1534. }
  1535. private static string ResolveKpiValueTable(int metricLevel)
  1536. {
  1537. return metricLevel switch
  1538. {
  1539. 1 => "ado_s9_kpi_value_l1_day",
  1540. 2 => "ado_s9_kpi_value_l2_day",
  1541. 3 => "ado_s9_kpi_value_l3_day",
  1542. 4 => "ado_s9_kpi_value_l4_day",
  1543. _ => "ado_s9_kpi_value_l2_day"
  1544. };
  1545. }
  1546. private static decimal DefaultS1Target(string metricCode) => LegacyKpiCodeTargets.GetOrZero(metricCode);
  1547. private static string ResolveKpiStatus(decimal actual, decimal target, string? direction, decimal? yellowThreshold, decimal? redThreshold)
  1548. {
  1549. if (target <= 0) return "gray";
  1550. var ratio = actual / target * 100m;
  1551. if (string.Equals(direction, "lower_is_better", StringComparison.OrdinalIgnoreCase))
  1552. {
  1553. if (actual <= target) return "green";
  1554. if (ratio <= (yellowThreshold ?? 110m)) return "yellow";
  1555. return ratio >= (redThreshold ?? 120m) ? "red" : "yellow";
  1556. }
  1557. if (actual >= target) return "green";
  1558. if (ratio >= (yellowThreshold ?? 95m)) return "yellow";
  1559. return ratio <= (redThreshold ?? 80m) ? "red" : "yellow";
  1560. }
  1561. private static string ResolveTrendFlag(decimal actual, decimal? previous)
  1562. {
  1563. if (previous == null) return "flat";
  1564. if (actual > previous.Value) return "up";
  1565. if (actual < previous.Value) return "down";
  1566. return "flat";
  1567. }
  1568. private static string Truncate(string? raw, int maxLength)
  1569. {
  1570. if (string.IsNullOrEmpty(raw)) return string.Empty;
  1571. return raw.Length <= maxLength ? raw : raw[..maxLength];
  1572. }
  1573. private sealed class S1ColumnRow
  1574. {
  1575. public string ColumnName { get; set; } = string.Empty;
  1576. }
  1577. private sealed class S1MdpEntityRow
  1578. {
  1579. public long Id { get; set; }
  1580. public string EntityName { get; set; } = string.Empty;
  1581. }
  1582. private sealed class S1KpiCalcRow
  1583. {
  1584. public long TenantId { get; set; }
  1585. public long FactoryId { get; set; }
  1586. public string MetricCode { get; set; } = string.Empty;
  1587. public decimal? MetricValue { get; set; }
  1588. }
  1589. private sealed class S1KpiMetaRow
  1590. {
  1591. public int MetricLevel { get; set; }
  1592. public string Direction { get; set; } = "higher_is_better";
  1593. public decimal? YellowThreshold { get; set; }
  1594. public decimal? RedThreshold { get; set; }
  1595. }
  1596. private sealed class S1KpiValueRow
  1597. {
  1598. public long Id { get; set; }
  1599. public decimal? MetricValue { get; set; }
  1600. public decimal? TargetValue { get; set; }
  1601. }
  1602. }
  1603. public sealed class S1MdpSyncTransformResult
  1604. {
  1605. public long RunLogId { get; set; }
  1606. public string BatchId { get; set; } = string.Empty;
  1607. public int StageRows { get; set; }
  1608. public int StandardRows { get; set; }
  1609. public int DwdRows { get; set; }
  1610. public int KpiRows { get; set; }
  1611. public int AtomicRows { get; set; }
  1612. }
  1613. internal sealed record S1MdpSqlCommand(string Sql, SugarParameter[] Parameters);
  1614. internal sealed record S1MdpEntityConfig(
  1615. string EntityCode,
  1616. string SourceTable,
  1617. string TargetTable,
  1618. string SourceRowIdExpression,
  1619. string SourceBizKeyExpression)
  1620. {
  1621. public static readonly IReadOnlyList<S1MdpEntityConfig> All = new List<S1MdpEntityConfig>
  1622. {
  1623. new("S1_SEORDER", "crm_seorder", "mdp_stg_so", "Id", "COALESCE(s.`bill_no`, CAST(s.`Id` AS CHAR))"),
  1624. new("S1_SEORDER_ENTRY", "crm_seorderentry", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`entry_seq`, CAST(s.`Id` AS CHAR)))"),
  1625. new("S1_SEORDER_CHANGE", "crm_seorder_change", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', CAST(s.`Id` AS CHAR))"),
  1626. new("S1_CONTRACT_REVIEW", "ado_contract_review", "mdp_stg_so", "RecID", "COALESCE(s.`BillNo`, CAST(s.`RecID` AS CHAR))"),
  1627. 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))"),
  1628. new("S1_PRODUCT_DESIGN", "ado_product_design", "mdp_stg_so", "Id", "COALESCE(s.`BillNo`, CAST(s.`Id` AS CHAR))"),
  1629. new("S1_PRODUCT_DESIGN_BOM", "ado_product_design_bom", "mdp_stg_so", "Id", "CONCAT(CAST(s.`ProductDesignId` AS CHAR), ':', CAST(s.`Id` AS CHAR))"),
  1630. new("S1_PRODUCT_DESIGN_ROUTING", "ado_product_design_routing", "mdp_stg_so", "Id", "CONCAT(CAST(s.`ProductDesignId` AS CHAR), ':', CAST(s.`Id` AS CHAR))"),
  1631. 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))"),
  1632. 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))"),
  1633. new("S1_SHIPPING_PLAN", "ShippingPlan", "mdp_stg_ship_trans", "RecID", "COALESCE(s.`LotSerial`, CAST(s.`RecID` AS CHAR))"),
  1634. new("S1_SHIPPING_PLAN_DETAIL", "ShippingPlanDetail", "mdp_stg_ship_trans", "RecID", "CONCAT(IFNULL(s.`plan_id`,''), ':', IFNULL(s.`OrdNbr`,''), ':', CAST(s.`RecID` AS CHAR))"),
  1635. new("S1_ASN_SHIPPER_MASTER", "ASNBOLShipperMaster", "mdp_stg_ship_trans", "RecID", "COALESCE(s.`Id`, CONCAT(IFNULL(s.`OrdNbr`,''), ':', CAST(s.`RecID` AS CHAR)))"),
  1636. new("S1_ASN_SHIPPER_DETAIL", "ASNBOLShipperDetail", "mdp_stg_ship_trans", "RecID", "CONCAT(IFNULL(s.`Id`,''), ':', IFNULL(s.`Line`, CAST(s.`RecID` AS CHAR)))"),
  1637. new("S1_LINKAGE_PLAN", "LinkagePlan", "mdp_stg_ship_trans", "id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`item_number`,''), ':', CAST(s.`id` AS CHAR))")
  1638. };
  1639. }