S1MdpSyncTransformService.cs 120 KB

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