S2MdpSyncTransformService.cs 55 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985
  1. using Admin.NET.Plugin.AiDOP.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  3. using Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  4. using Admin.NET.Plugin.AiDOP.SmartOps;
  5. namespace Admin.NET.Plugin.AiDOP.Production;
  6. /// <summary>
  7. /// S2 生产排程 MDP 同步、标准化、DWD 与 KPI 计算服务。
  8. /// Phase2:源→stg 经统一执行器;std/DWD/KPI 仍读 stg。
  9. /// </summary>
  10. public class S2MdpSyncTransformService : ITransient
  11. {
  12. private const string JobCode = "S2_MDP_SYNC_TRANSFORM";
  13. private readonly ISqlSugarClient _db;
  14. private readonly TransformRunLogFinalizer _runLogFinalizer;
  15. private readonly SmartOpsKpiAtomicBuildService _atomicBuild;
  16. private readonly MdpModuleStagingPuller _stagingPuller;
  17. private readonly IKpiTargetResolver _kpiTargetResolver;
  18. private readonly S9CompositeKpiWriter _s9Composite;
  19. private readonly MdpNeutralSourceGate _neutralGate;
  20. private readonly NativeNeutralProjectionService _nativeProjection;
  21. private MdpRebuildScope _runScope = null!;
  22. public S2MdpSyncTransformService(
  23. ISqlSugarClient db,
  24. SmartOpsKpiAtomicBuildService atomicBuild,
  25. MdpModuleStagingPuller stagingPuller,
  26. IKpiTargetResolver kpiTargetResolver,
  27. TransformRunLogFinalizer runLogFinalizer,
  28. S9CompositeKpiWriter s9Composite,
  29. MdpNeutralSourceGate neutralGate,
  30. NativeNeutralProjectionService nativeProjection)
  31. {
  32. _db = db;
  33. _runLogFinalizer = runLogFinalizer;
  34. _atomicBuild = atomicBuild;
  35. _stagingPuller = stagingPuller;
  36. _kpiTargetResolver = kpiTargetResolver;
  37. _s9Composite = s9Composite;
  38. _neutralGate = neutralGate;
  39. _nativeProjection = nativeProjection;
  40. }
  41. public Task<S2MdpSyncTransformResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO") =>
  42. throw new InvalidOperationException("S2 全量重算必须指定租户/工厂作用域,请走 ModuleRebuildService.Enqueue");
  43. public async Task<S2MdpSyncTransformResult> RunFullAsync(
  44. MdpRebuildScope scope,
  45. CancellationToken cancellationToken = default,
  46. string triggerType = "AUTO",
  47. long? rebuildJobId = null,
  48. Func<ModuleProgressUpdate, Task>? reportProgress = null)
  49. {
  50. _ = rebuildJobId;
  51. _runScope = MdpRebuildScope.Create("S2", scope.TenantId, scope.FactoryId);
  52. cancellationToken.ThrowIfCancellationRequested();
  53. var now = DateTime.Now;
  54. var batchId = $"S2_MDP_FULL_{_runScope.TenantId}_{_runScope.FactoryId}_{now:yyyyMMddHHmmss}";
  55. var runLogId = await InsertTransformRunLogAsync(batchId, now, triggerType);
  56. var result = new S2MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  57. try
  58. {
  59. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Preparing, 0, 5, "准备运行环境"));
  60. await EnsureS2RuntimeObjectsAsync();
  61. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Staging, 1, 10, "正在拉取源数据"));
  62. result.StageRows = await SyncStagingAsync(batchId, now, cancellationToken, triggerType);
  63. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Staging, 1, 35, "拉取源数据完成", result.StageRows, ModuleRebuildStages.Staging));
  64. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Standard, 2, 35, "正在标准化数据"));
  65. result.StandardRows = await TransformStandardAsync(batchId, now, cancellationToken);
  66. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Standard, 2, 50, "标准化数据完成", result.StandardRows, ModuleRebuildStages.Standard));
  67. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Dwd, 3, 50, "正在生成 DWD 明细"));
  68. result.DwdRows = await BuildDwdAsync(batchId, now, cancellationToken);
  69. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Dwd, 3, 65, "生成 DWD 明细完成", result.DwdRows, ModuleRebuildStages.Dwd));
  70. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Kpi, 4, 65, "正在重算 KPI"));
  71. result.KpiRows = await BuildS2KpiValuesAsync(batchId, now, cancellationToken);
  72. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Kpi, 4, 80, "重算 KPI 完成", result.KpiRows, ModuleRebuildStages.Kpi));
  73. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Atomic, 5, 82, "正在重建原子数据"));
  74. result.AtomicRows = await _atomicBuild.BuildWorkScheduleDomainForAllDatesAsync(
  75. _runScope.TenantId, _runScope.FactoryId, batchId, cancellationToken);
  76. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Atomic, 5, 98, "重建原子数据完成", result.AtomicRows, ModuleRebuildStages.Atomic));
  77. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Finalizing, 5, 99, "正在收尾"));
  78. await _nativeProjection.ProjectAsync(_runScope.TenantId, batchId, cancellationToken);
  79. result.KpiRows += await _s9Composite.WriteRecentAsync(_runScope.TenantId, _runScope.FactoryId, now, cancellationToken);
  80. await MarkTransformRunSuccessAsync(runLogId, now, result);
  81. return result;
  82. }
  83. catch (Exception ex)
  84. {
  85. if (ex is ModuleRebuildCancelledException)
  86. throw;
  87. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  88. if (!_runLogFinalizer.IsHostStopping)
  89. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  90. throw;
  91. }
  92. finally
  93. {
  94. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  95. }
  96. }
  97. private static async Task ReportAsync(Func<ModuleProgressUpdate, Task>? report, ModuleProgressUpdate update)
  98. {
  99. if (report == null) return;
  100. await report(update);
  101. }
  102. private async Task EnsureS2RuntimeObjectsAsync()
  103. {
  104. foreach (var table in new[]
  105. {
  106. "mdp_stg_schedule",
  107. "mdp_std_work_order_schedule",
  108. "mdp_std_operation_schedule",
  109. "dwd_order_schedule_trans"
  110. })
  111. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl(table));
  112. foreach (var sql in S2MdpDdl.SqlBlocks)
  113. {
  114. await MdpSchemaAligner.ExecuteAsync(_db, sql);
  115. }
  116. }
  117. private async Task<int> SyncStagingAsync(string batchId, DateTime now, CancellationToken cancellationToken, string triggerType = "AUTO")
  118. {
  119. _ = now;
  120. EnsureRunScope();
  121. return await _stagingPuller.PullEntitiesAsync(
  122. S2MdpEntityConfig.All.Select(x => x.EntityCode),
  123. batchId, _runScope.TenantId, fullRefresh: ModuleRebuildTriggerType.ShouldPullFull(triggerType), taskCode: "S2_MDP_INBOUND", cancellationToken,
  124. factoryId: _runScope.FactoryId, requireMatchingSourceTenant: true);
  125. }
  126. public async Task<S2MdpSyncTransformResult> RunInboundAsync(
  127. IEnumerable<string>? entityCodes = null,
  128. long tenantId = 0,
  129. bool fullRefresh = false,
  130. CancellationToken cancellationToken = default)
  131. {
  132. cancellationToken.ThrowIfCancellationRequested();
  133. if (tenantId <= 0)
  134. throw new InvalidOperationException("S2 inbound 必须指定有效 tenantId");
  135. _runScope = MdpRebuildScope.Create("S2", tenantId, 1);
  136. var now = DateTime.Now;
  137. var batchId = $"S2_MDP_IN_{_runScope.TenantId}_{now:yyyyMMddHHmmss}";
  138. var runLogId = await InsertTransformRunLogAsync(batchId, now, "INBOUND");
  139. var result = new S2MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  140. try
  141. {
  142. await EnsureS2RuntimeObjectsAsync();
  143. var codes = entityCodes?.ToList() ?? S2MdpEntityConfig.All.Select(x => x.EntityCode).ToList();
  144. result.StageRows = await _stagingPuller.PullEntitiesAsync(
  145. codes, batchId, tenantId, fullRefresh, "S2_MDP_INBOUND", cancellationToken);
  146. result.StandardRows = await TransformStandardAsync(batchId, now, cancellationToken);
  147. result.DwdRows = await BuildDwdAsync(batchId, now, cancellationToken);
  148. result.KpiRows = await BuildS2KpiValuesAsync(batchId, now, cancellationToken);
  149. result.AtomicRows = await _atomicBuild.BuildWorkScheduleDomainForAllDatesAsync(
  150. _runScope.TenantId, _runScope.FactoryId, batchId, cancellationToken);
  151. await _nativeProjection.ProjectAsync(_runScope.TenantId, batchId, cancellationToken);
  152. result.KpiRows += await _s9Composite.WriteRecentAsync(_runScope.TenantId, _runScope.FactoryId, now, cancellationToken);
  153. await MarkTransformRunSuccessAsync(runLogId, now, result);
  154. return result;
  155. }
  156. catch (Exception ex)
  157. {
  158. if (ex is ModuleRebuildCancelledException)
  159. throw;
  160. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  161. if (!_runLogFinalizer.IsHostStopping)
  162. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  163. throw;
  164. }
  165. finally
  166. {
  167. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  168. }
  169. }
  170. [Obsolete("Phase3: use MdpModuleStagingPuller / executors")]
  171. private async Task<int> SyncOneEntityAsync(S2MdpEntityConfig entity, string batchId, DateTime now)
  172. {
  173. var entityRow = await _db.Ado.SqlQuerySingleAsync<S2MdpEntityRow>(
  174. "SELECT id AS Id, entity_name AS EntityName FROM mdp_entity WHERE tenant_id=0 AND entity_code=@EntityCode LIMIT 1",
  175. new SugarParameter("@EntityCode", entity.EntityCode));
  176. if (entityRow == null) throw Oops.Oh($"未找到 MDP 实体配置:{entity.EntityCode}");
  177. var columns = await _db.Ado.SqlQueryAsync<S2ColumnRow>(
  178. """
  179. SELECT COLUMN_NAME AS ColumnName
  180. FROM information_schema.COLUMNS
  181. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@TableName
  182. ORDER BY ORDINAL_POSITION
  183. """,
  184. new SugarParameter("@TableName", entity.SourceTable));
  185. if (columns.Count == 0) throw Oops.Oh($"未找到源表:{entity.SourceTable}");
  186. var names = columns.Select(u => u.ColumnName).ToList();
  187. var tenantExpr = names.Any(u => string.Equals(u, "tenant_id", StringComparison.OrdinalIgnoreCase))
  188. ? $"IFNULL(s.`{FindColumn(names, "tenant_id")}`,0)"
  189. : "0";
  190. var factoryExpr = names.Any(u => string.Equals(u, "factory_id", StringComparison.OrdinalIgnoreCase))
  191. ? $"s.`{FindColumn(names, "factory_id")}`"
  192. : "NULL";
  193. var domainExpr = names.Any(u => string.Equals(u, "Domain", StringComparison.OrdinalIgnoreCase))
  194. ? $"s.`{FindColumn(names, "Domain")}`"
  195. : "NULL";
  196. var sourceRowExpr = names.Any(u => string.Equals(u, entity.SourceRowIdExpression, StringComparison.OrdinalIgnoreCase))
  197. ? $"s.`{FindColumn(names, entity.SourceRowIdExpression)}`"
  198. : entity.SourceRowIdExpression;
  199. var rawDataExpr = BuildJsonObjectExpression(names);
  200. var rowsRead = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM `{entity.SourceTable}`");
  201. var logId = await InsertSyncLogAsync(entityRow.Id, entityRow.EntityName, batchId, rowsRead);
  202. var started = DateTime.Now;
  203. try
  204. {
  205. var affected = await MdpSchemaAligner.ExecuteAsync(_db,
  206. $"""
  207. INSERT INTO `{entity.TargetTable}`
  208. (tenant_id, factory_id, source_system, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time, process_status, raw_data)
  209. SELECT
  210. {tenantExpr},
  211. COALESCE({factoryExpr}, {domainExpr}),
  212. COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = @SourceTable LIMIT 1), NULLIF('',''), ''),
  213. @SourceTable,
  214. CAST({sourceRowExpr} AS CHAR),
  215. CAST(COALESCE({entity.SourceBizKeyExpression}, CAST({sourceRowExpr} AS CHAR)) AS CHAR),
  216. @BatchId,
  217. @Now,
  218. 'PENDING',
  219. {rawDataExpr}
  220. FROM `{entity.SourceTable}` s
  221. ON DUPLICATE KEY UPDATE
  222. tenant_id=VALUES(tenant_id),
  223. factory_id=VALUES(factory_id),
  224. sync_batch_id=VALUES(sync_batch_id),
  225. sync_time=VALUES(sync_time),
  226. process_status=VALUES(process_status),
  227. raw_data=VALUES(raw_data),
  228. update_time=CURRENT_TIMESTAMP
  229. """,
  230. new SugarParameter("@SourceTable", entity.SourceTable),
  231. new SugarParameter("@BatchId", batchId),
  232. new SugarParameter("@Now", now));
  233. await MarkSyncLogSuccessAsync(logId, started, affected);
  234. return rowsRead;
  235. }
  236. catch (Exception ex)
  237. {
  238. await MarkSyncLogFailedAsync(logId, started, ex.Message);
  239. throw;
  240. }
  241. }
  242. private async Task<int> TransformStandardAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  243. {
  244. if (!await _neutralGate.AllowsAsync(_runScope.TenantId, "WO_SCHEDULE", "DOPDEMORQ_SQLSERVER"))
  245. return 0;
  246. var total = 0;
  247. foreach (var command in BuildStandardCommands(batchId, now))
  248. {
  249. cancellationToken.ThrowIfCancellationRequested();
  250. total += await MdpSchemaAligner.ExecuteAsync(_db, command.Sql, command.Parameters);
  251. }
  252. return total;
  253. }
  254. private async Task<int> BuildDwdAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  255. {
  256. var total = 0;
  257. foreach (var command in BuildDwdCommands(batchId, now))
  258. {
  259. cancellationToken.ThrowIfCancellationRequested();
  260. total += await MdpSchemaAligner.ExecuteAsync(_db, command.Sql, command.Parameters);
  261. }
  262. return total;
  263. }
  264. private async Task<int> BuildS2KpiValuesAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  265. {
  266. var affected = 0;
  267. for (var dayOffset = 13; dayOffset >= 0; dayOffset--)
  268. {
  269. var statDate = now.Date.AddDays(-dayOffset);
  270. var rows = await CalculateS2KpiValuesAsync(batchId, statDate);
  271. foreach (var row in rows)
  272. {
  273. cancellationToken.ThrowIfCancellationRequested();
  274. affected += await UpsertS2KpiValueAsync(row, statDate, now);
  275. }
  276. }
  277. return affected;
  278. }
  279. private IEnumerable<S2MdpSqlCommand> BuildStandardCommands(string batchId, DateTime now)
  280. {
  281. yield return Cmd(
  282. """
  283. INSERT INTO mdp_std_work_order_schedule
  284. (tenant_id, factory_id, source_system, doc_type, work_order, sales_order_no, sales_order_entry_id, item_code, item_name, site_code, status,
  285. priority, urgent_flag, qty_ordered, qty_completed, order_date, due_date, release_date, prod_line,
  286. specification, lot_serial, drawing_no, project, work_order_type, labor_variance,
  287. source_biz_key, sync_batch_id, sync_time)
  288. SELECT s.tenant_id,
  289. CASE WHEN s.factory_id REGEXP '^[0-9]+$' THEN CAST(s.factory_id AS UNSIGNED) ELSE 1 END,
  290. COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = s.source_table LIMIT 1), NULLIF(s.source_system,''), ''),
  291. 'PROD_TASK',
  292. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrd')),
  293. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SalesJob')),
  294. -- 订单行 ID 桥:WorkOrdMaster.BusinessID,>0 时恒等于 crm_seorderentry.Id
  295. -- (全库实测 28/28 命中同租户订单行、0 孤儿)。源值 0 表示「无关联订单行」,
  296. -- 必须归一为 NULL —— 0 不是合法订单行 Id,落 0 会在桥上造出假关联。
  297. -- sales_order_no(=SalesJob) 只是文本,禁止用它做归因。
  298. NULLIF(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.BusinessID')) REGEXP '^[0-9]+$'
  299. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.BusinessID')) AS SIGNED) END, 0),
  300. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemNum')),
  301. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemName')),
  302. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Site')),
  303. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Status')),
  304. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Priority')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Priority')) AS DECIMAL(18,6)) END,
  305. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Urgent')) IN ('1','true','True') THEN 1 ELSE 0 END,
  306. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyOrded')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyOrded')) AS DECIMAL(18,6)) ELSE 0 END,
  307. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyCompleted')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyCompleted')) AS DECIMAL(18,6)) ELSE 0 END,
  308. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdDate')), 'null'), ''),
  309. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DueDate')), 'null'), ''),
  310. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReleaseDate')), 'null'), ''),
  311. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProdLine')),
  312. -- Specification(型号)= ItemMaster.Descr1 中台 enrich((Domain,ItemNum) 去重取一,实测同键 Descr1 无冲突;
  313. -- 严格按 Domain+ItemNum 关联,不做 ItemNum-only 去重,避免跨工厂串型号):
  314. im.Descr1,
  315. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LotSerial')), 'null'), ''),
  316. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Drawing')), 'null'), ''),
  317. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Project')), 'null'), ''),
  318. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WoTyped')), 'null'), ''),
  319. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LbrVar')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LbrVar')) AS DECIMAL(18,6)) END,
  320. s.source_biz_key, @BatchId, @Now
  321. FROM mdp_stg_schedule s
  322. LEFT JOIN (
  323. SELECT `Domain`, ItemNum, MIN(Descr1) AS Descr1
  324. FROM ItemMaster
  325. GROUP BY `Domain`, ItemNum
  326. ) im ON im.`Domain` = JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain'))
  327. AND im.ItemNum = JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemNum'))
  328. WHERE s.source_table='WorkOrdMaster' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrd')), '') <> ''
  329. ON DUPLICATE KEY UPDATE
  330. sales_order_no=VALUES(sales_order_no), sales_order_entry_id=VALUES(sales_order_entry_id),
  331. item_code=VALUES(item_code), item_name=VALUES(item_name),
  332. site_code=VALUES(site_code), status=VALUES(status), priority=VALUES(priority), urgent_flag=VALUES(urgent_flag),
  333. qty_ordered=VALUES(qty_ordered), qty_completed=VALUES(qty_completed), order_date=VALUES(order_date),
  334. due_date=VALUES(due_date), release_date=VALUES(release_date), prod_line=VALUES(prod_line),
  335. specification=VALUES(specification), lot_serial=VALUES(lot_serial), drawing_no=VALUES(drawing_no),
  336. project=VALUES(project), work_order_type=VALUES(work_order_type), labor_variance=VALUES(labor_variance),
  337. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  338. """, batchId, now);
  339. yield return Cmd(
  340. """
  341. INSERT INTO mdp_std_operation_schedule
  342. (tenant_id, factory_id, source_system, work_order, op_no, work_center, line_code, item_code,
  343. plan_date, prod_date, start_time, end_time, ord_qty, comp_qty, run_crew, employee, status,
  344. source_biz_key, sync_batch_id, sync_time)
  345. SELECT tenant_id,
  346. CASE WHEN factory_id REGEXP '^[0-9]+$' THEN CAST(factory_id AS UNSIGNED) ELSE 1 END,
  347. COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = mdp_stg_schedule.source_table LIMIT 1), NULLIF(mdp_stg_schedule.source_system,''), ''),
  348. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrds')),
  349. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Op')),
  350. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkCtr')),
  351. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Line')),
  352. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ItemNum')),
  353. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.PlanDate')), 'null'), ''),
  354. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ProdDate')), 'null'), ''),
  355. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''),
  356. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.EndTime')), 'null'), ''),
  357. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.OrdQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.OrdQty')) AS DECIMAL(18,6)) ELSE 0 END,
  358. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompQty')) AS DECIMAL(18,6)) ELSE 0 END,
  359. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RunCrew')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RunCrew')) AS DECIMAL(18,6)) END,
  360. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Employee')),
  361. -- 日计划下达状态:'Y'=已下达(可执行),空=未下达。原样投影 PeriodSequenceDet.Status。
  362. -- 它只是【状态佐证】—— Stage-2 End 的时间 Authority 一律取
  363. -- S2_DAILY_PLAN_RELEASE SUCCESS 的 end_time,绝不能用 update_time 当下达时间
  364. -- (update_time 会被后续任意一次 UPDATE 覆盖)。
  365. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Status')),
  366. source_biz_key, @BatchId, @Now
  367. FROM mdp_stg_schedule
  368. WHERE source_table='PeriodSequenceDet' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrds')), '') <> ''
  369. ON DUPLICATE KEY UPDATE
  370. work_center=VALUES(work_center), line_code=VALUES(line_code), item_code=VALUES(item_code),
  371. plan_date=VALUES(plan_date), prod_date=VALUES(prod_date), start_time=VALUES(start_time), end_time=VALUES(end_time),
  372. ord_qty=VALUES(ord_qty), comp_qty=VALUES(comp_qty), run_crew=VALUES(run_crew), employee=VALUES(employee),
  373. status=VALUES(status),
  374. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  375. """, batchId, now);
  376. // ⚠️ 同表的另一个来源:ScheduleResultOpMaster —— 它没有「日计划下达」这个概念。
  377. // 因此既不投影 status,ODKU 里也【刻意不列】status,
  378. // 否则会把 PeriodSequenceDet 分支刚投影进来的下达状态覆盖成 NULL。
  379. // (两来源 source_biz_key 格式不同,uk=(tenant_id,source_biz_key) 实测不碰撞;此为纵深防御。)
  380. yield return Cmd(
  381. """
  382. INSERT INTO mdp_std_operation_schedule
  383. (tenant_id, factory_id, source_system, work_order, op_no, work_center, line_code, item_code,
  384. plan_date, prod_date, start_time, end_time, ord_qty, comp_qty, run_crew, employee,
  385. source_biz_key, sync_batch_id, sync_time)
  386. SELECT tenant_id,
  387. CASE WHEN factory_id REGEXP '^[0-9]+$' THEN CAST(factory_id AS UNSIGNED) ELSE 1 END,
  388. COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = mdp_stg_schedule.source_table LIMIT 1), NULLIF(mdp_stg_schedule.source_system,''), ''),
  389. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrd')),
  390. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Op')),
  391. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkCtr')),
  392. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Line')),
  393. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ItemNum')),
  394. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkDate')), 'null'), ''),
  395. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkDate')), 'null'), ''),
  396. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkStartTime')), 'null'), ''),
  397. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkEndTime')), 'null'), ''),
  398. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkQty')) AS DECIMAL(18,6)) ELSE 0 END,
  399. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkQty')) AS DECIMAL(18,6)) ELSE 0 END,
  400. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.AssignedPersonnelCount')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.AssignedPersonnelCount')) AS DECIMAL(18,6)) END,
  401. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.AssignedEmployeeID')),
  402. source_biz_key, @BatchId, @Now
  403. FROM mdp_stg_schedule
  404. WHERE source_table='ScheduleResultOpMaster' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrd')), '') <> ''
  405. ON DUPLICATE KEY UPDATE
  406. work_center=VALUES(work_center), line_code=VALUES(line_code), item_code=VALUES(item_code),
  407. plan_date=VALUES(plan_date), prod_date=VALUES(prod_date), start_time=VALUES(start_time), end_time=VALUES(end_time),
  408. ord_qty=VALUES(ord_qty), comp_qty=VALUES(comp_qty), run_crew=VALUES(run_crew), employee=VALUES(employee),
  409. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  410. """, batchId, now);
  411. }
  412. private IEnumerable<S2MdpSqlCommand> BuildDwdCommands(string batchId, DateTime now)
  413. {
  414. yield return Cmd(
  415. """
  416. INSERT INTO dwd_order_schedule_trans
  417. (tenant_id, factory_id, stat_date, doc_type, work_order, sales_order_no, item_code, item_name, site_code, prod_line,
  418. status, urgent_flag, qty_ordered, qty_completed, order_date, due_date, release_date, first_plan_date,
  419. last_plan_date, first_start_time, last_end_time, operation_count, scheduled_qty, completed_op_qty,
  420. schedule_cycle_days, schedule_satisfaction_flag, wip_qty, resource_person_count, calc_batch_id, calc_time)
  421. SELECT w.tenant_id, COALESCE(w.factory_id, 1), @StatDate, w.doc_type,
  422. w.work_order, w.sales_order_no, w.item_code, w.item_name, w.site_code, w.prod_line,
  423. w.status, w.urgent_flag, w.qty_ordered, w.qty_completed, w.order_date, w.due_date, w.release_date,
  424. MIN(o.plan_date), MAX(o.plan_date), MIN(o.start_time), MAX(o.end_time),
  425. COUNT(o.id), SUM(IFNULL(o.ord_qty, 0)), SUM(IFNULL(o.comp_qty, 0)),
  426. CASE
  427. WHEN COALESCE(MIN(o.plan_date), MAX(o.plan_date)) IS NOT NULL
  428. AND COALESCE(w.release_date, w.order_date) IS NOT NULL
  429. THEN TIMESTAMPDIFF(HOUR, COALESCE(w.release_date, w.order_date), COALESCE(MAX(o.plan_date), MIN(o.plan_date))) / 24
  430. END,
  431. CASE
  432. WHEN w.due_date IS NOT NULL AND COALESCE(MAX(o.plan_date), MIN(o.plan_date)) IS NOT NULL
  433. AND DATE(COALESCE(MAX(o.plan_date), MIN(o.plan_date))) <= DATE(w.due_date) THEN 1
  434. ELSE 0
  435. END,
  436. GREATEST(IFNULL(w.qty_ordered, 0) - IFNULL(w.qty_completed, 0), 0),
  437. SUM(IFNULL(o.run_crew, 0)),
  438. @BatchId, @Now
  439. FROM mdp_std_work_order_schedule w
  440. LEFT JOIN mdp_std_operation_schedule o
  441. ON o.tenant_id=w.tenant_id AND IFNULL(o.work_order,'')=IFNULL(w.work_order,'')
  442. WHERE IFNULL(w.work_order, '') <> ''
  443. GROUP BY w.tenant_id, COALESCE(w.factory_id, 1), w.doc_type, w.work_order, w.sales_order_no, w.item_code, w.item_name,
  444. w.site_code, w.prod_line, w.status, w.urgent_flag, w.qty_ordered, w.qty_completed,
  445. w.order_date, w.due_date, w.release_date
  446. ON DUPLICATE KEY UPDATE
  447. sales_order_no=VALUES(sales_order_no), item_code=VALUES(item_code), item_name=VALUES(item_name),
  448. site_code=VALUES(site_code), prod_line=VALUES(prod_line), status=VALUES(status), urgent_flag=VALUES(urgent_flag),
  449. qty_ordered=VALUES(qty_ordered), qty_completed=VALUES(qty_completed), order_date=VALUES(order_date),
  450. due_date=VALUES(due_date), release_date=VALUES(release_date), first_plan_date=VALUES(first_plan_date),
  451. last_plan_date=VALUES(last_plan_date), first_start_time=VALUES(first_start_time), last_end_time=VALUES(last_end_time),
  452. operation_count=VALUES(operation_count), scheduled_qty=VALUES(scheduled_qty), completed_op_qty=VALUES(completed_op_qty),
  453. schedule_cycle_days=VALUES(schedule_cycle_days), schedule_satisfaction_flag=VALUES(schedule_satisfaction_flag),
  454. wip_qty=VALUES(wip_qty), resource_person_count=VALUES(resource_person_count), calc_time=VALUES(calc_time),
  455. update_time=CURRENT_TIMESTAMP
  456. """, batchId, now);
  457. }
  458. private async Task<List<S2KpiCalcRow>> CalculateS2KpiValuesAsync(string batchId, DateTime statDate)
  459. {
  460. var scope = RequireScope();
  461. return await _db.Ado.SqlQueryAsync<S2KpiCalcRow>(MdpSqlScope.InjectTenantFactory(
  462. """
  463. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L1_001' AS MetricCode,
  464. ROUND(AVG(schedule_cycle_days), 4) AS MetricValue
  465. FROM dwd_order_schedule_trans
  466. WHERE calc_batch_id=@BatchId AND doc_type='PROD_TASK' AND schedule_cycle_days IS NOT NULL AND schedule_cycle_days >= 0
  467. GROUP BY tenant_id, factory_id
  468. UNION ALL
  469. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L1_002' AS MetricCode,
  470. ROUND(100 * SUM(schedule_satisfaction_flag) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  471. FROM dwd_order_schedule_trans
  472. WHERE calc_batch_id=@BatchId AND doc_type='PROD_TASK'
  473. GROUP BY tenant_id, factory_id
  474. UNION ALL
  475. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L1_003' AS MetricCode,
  476. ROUND(COUNT(1) / NULLIF(SUM(CASE WHEN IFNULL(resource_person_count, 0) > 0 THEN resource_person_count ELSE 1 END), 0), 4) AS MetricValue
  477. FROM dwd_order_schedule_trans
  478. WHERE calc_batch_id=@BatchId AND doc_type='PROD_TASK'
  479. GROUP BY tenant_id, factory_id
  480. UNION ALL
  481. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L1_004' AS MetricCode,
  482. ROUND(SUM(wip_qty) / NULLIF(SUM(CASE WHEN qty_completed > 0 THEN qty_completed ELSE completed_op_qty END), 0) * 30, 4) AS MetricValue
  483. FROM dwd_order_schedule_trans
  484. WHERE calc_batch_id=@BatchId AND doc_type='PROD_TASK'
  485. GROUP BY tenant_id, factory_id
  486. UNION ALL
  487. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L2_001' AS MetricCode,
  488. ROUND(AVG(schedule_cycle_days), 4) AS MetricValue
  489. FROM dwd_order_schedule_trans
  490. WHERE calc_batch_id=@BatchId AND doc_type='PROD_TASK' AND schedule_cycle_days IS NOT NULL AND schedule_cycle_days >= 0
  491. GROUP BY tenant_id, factory_id
  492. UNION ALL
  493. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L2_002' AS MetricCode,
  494. ROUND(100 * SUM(schedule_satisfaction_flag) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  495. FROM dwd_order_schedule_trans
  496. WHERE calc_batch_id=@BatchId AND doc_type='PROD_TASK'
  497. GROUP BY tenant_id, factory_id
  498. UNION ALL
  499. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L2_003' AS MetricCode,
  500. ROUND(COUNT(1) / NULLIF(SUM(CASE WHEN IFNULL(resource_person_count, 0) > 0 THEN resource_person_count ELSE 1 END), 0), 4) AS MetricValue
  501. FROM dwd_order_schedule_trans
  502. WHERE calc_batch_id=@BatchId AND doc_type='PROD_TASK'
  503. GROUP BY tenant_id, factory_id
  504. UNION ALL
  505. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_001' AS MetricCode,
  506. ROUND(AVG(
  507. CASE
  508. WHEN start_time IS NOT NULL AND end_time IS NOT NULL AND end_time >= start_time
  509. THEN TIMESTAMPDIFF(HOUR, start_time, end_time) / 24
  510. WHEN plan_date IS NOT NULL AND prod_date IS NOT NULL
  511. THEN ABS(TIMESTAMPDIFF(DAY, plan_date, prod_date))
  512. END
  513. ), 4) AS MetricValue
  514. FROM mdp_std_operation_schedule
  515. WHERE IFNULL(work_order, '') <> ''
  516. GROUP BY tenant_id, factory_id
  517. UNION ALL
  518. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_002' AS MetricCode,
  519. ROUND(100 * SUM(CASE WHEN plan_date IS NOT NULL AND (prod_date IS NULL OR prod_date >= plan_date) THEN 1 ELSE 0 END) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  520. FROM mdp_std_operation_schedule
  521. WHERE IFNULL(work_order, '') <> ''
  522. GROUP BY tenant_id, factory_id
  523. UNION ALL
  524. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_003' AS MetricCode,
  525. ROUND(COUNT(1) / NULLIF(SUM(CASE WHEN IFNULL(run_crew, 0) > 0 THEN run_crew ELSE 1 END), 0), 4) AS MetricValue
  526. FROM mdp_std_operation_schedule
  527. WHERE IFNULL(work_order, '') <> ''
  528. GROUP BY tenant_id, factory_id
  529. UNION ALL
  530. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_004' AS MetricCode,
  531. ROUND(AVG(
  532. CASE
  533. WHEN start_time IS NOT NULL AND end_time IS NOT NULL AND end_time >= start_time
  534. THEN TIMESTAMPDIFF(HOUR, start_time, end_time) / 24
  535. WHEN plan_date IS NOT NULL AND prod_date IS NOT NULL
  536. THEN ABS(TIMESTAMPDIFF(DAY, plan_date, prod_date))
  537. END
  538. ), 4) AS MetricValue
  539. FROM mdp_std_operation_schedule
  540. WHERE IFNULL(COALESCE(NULLIF(work_center,''), NULLIF(line_code,'')), '') <> ''
  541. GROUP BY tenant_id, factory_id
  542. UNION ALL
  543. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_005' AS MetricCode,
  544. ROUND(100 * SUM(CASE WHEN plan_date IS NOT NULL AND (prod_date IS NULL OR prod_date >= plan_date) THEN 1 ELSE 0 END) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  545. FROM mdp_std_operation_schedule
  546. WHERE IFNULL(COALESCE(NULLIF(work_center,''), NULLIF(line_code,'')), '') <> ''
  547. GROUP BY tenant_id, factory_id
  548. UNION ALL
  549. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_006' AS MetricCode,
  550. ROUND(COUNT(DISTINCT COALESCE(NULLIF(work_center,''), NULLIF(line_code,''), CONCAT(work_order, ':', op_no)))
  551. / NULLIF(SUM(CASE WHEN IFNULL(run_crew, 0) > 0 THEN run_crew ELSE 1 END), 0), 4) AS MetricValue
  552. FROM mdp_std_operation_schedule
  553. WHERE IFNULL(COALESCE(NULLIF(work_center,''), NULLIF(line_code,'')), '') <> ''
  554. GROUP BY tenant_id, factory_id
  555. """),
  556. new SugarParameter("@BatchId", batchId),
  557. new SugarParameter("@StatDate", statDate),
  558. new SugarParameter("@TenantId", scope.TenantId),
  559. new SugarParameter("@FactoryId", scope.FactoryId));
  560. }
  561. private async Task<int> UpsertS2KpiValueAsync(S2KpiCalcRow row, DateTime statDate, DateTime now)
  562. {
  563. if (row.MetricValue == null) return 0;
  564. var meta = await _db.Ado.SqlQuerySingleAsync<S2KpiMetaRow>(
  565. """
  566. SELECT MetricLevel, Direction, YellowThreshold, RedThreshold
  567. FROM ado_smart_ops_kpi_master
  568. WHERE TenantId=@TenantId AND ModuleCode='S2' AND MetricCode=@MetricCode AND IsEnabled=1
  569. LIMIT 1
  570. """,
  571. new SugarParameter("@TenantId", row.TenantId),
  572. new SugarParameter("@MetricCode", row.MetricCode));
  573. if (meta == null) return 0;
  574. var table = ResolveKpiValueTable(meta.MetricLevel);
  575. var current = await _db.Ado.SqlQuerySingleAsync<S2KpiValueRow>(
  576. $"""
  577. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  578. FROM {table}
  579. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S2'
  580. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  581. ORDER BY id
  582. LIMIT 1
  583. """,
  584. new SugarParameter("@TenantId", row.TenantId),
  585. new SugarParameter("@FactoryId", row.FactoryId),
  586. new SugarParameter("@MetricCode", row.MetricCode),
  587. new SugarParameter("@BizDate", statDate));
  588. var prior = await _db.Ado.SqlQuerySingleAsync<S2KpiValueRow>(
  589. $"""
  590. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  591. FROM {table}
  592. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S2'
  593. AND metric_code=@MetricCode AND biz_date<@BizDate AND is_deleted=0
  594. ORDER BY biz_date DESC, id DESC
  595. LIMIT 1
  596. """,
  597. new SugarParameter("@TenantId", row.TenantId),
  598. new SugarParameter("@FactoryId", row.FactoryId),
  599. new SugarParameter("@MetricCode", row.MetricCode),
  600. new SugarParameter("@BizDate", statDate));
  601. var actual = Math.Round(row.MetricValue.Value, 4);
  602. var snap = await _kpiTargetResolver.ResolveAsync(row.TenantId, row.FactoryId, row.MetricCode, "S2", statDate);
  603. var target = snap.TargetValue;
  604. var status = AidopS4KpiMerge.AchievementLevel(actual, target, meta.Direction, meta.YellowThreshold, meta.RedThreshold);
  605. var trend = ResolveTrendFlag(actual, prior?.MetricValue);
  606. if (current != null)
  607. {
  608. return await MdpSchemaAligner.ExecuteAsync(_db,
  609. $"""
  610. UPDATE {table}
  611. SET {Admin.NET.Plugin.AiDOP.SmartOps.KpiResultWrite.FromMetricSet}metric_value=@MetricValue, target_value=@TargetValue, status_color=@StatusColor, trend_flag=@TrendFlag,
  612. target_config_id=@TargetConfigId, target_source=@TargetSource, target_resolved_at=@TargetResolvedAt,
  613. is_active=1, status='ACTIVE', calc_time=@CalcTime, update_time=@CalcTime
  614. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S2'
  615. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  616. """,
  617. new SugarParameter("@MetricValue", actual),
  618. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  619. new SugarParameter("@StatusColor", status),
  620. new SugarParameter("@TrendFlag", trend),
  621. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  622. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  623. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt),
  624. new SugarParameter("@CalcTime", now),
  625. new SugarParameter("@TenantId", row.TenantId),
  626. new SugarParameter("@FactoryId", row.FactoryId),
  627. new SugarParameter("@MetricCode", row.MetricCode),
  628. new SugarParameter("@BizDate", statDate));
  629. }
  630. var nextId = Yitter.IdGenerator.YitIdHelper.NextId();
  631. return await MdpSchemaAligner.ExecuteAsync(_db,
  632. $"""
  633. INSERT INTO {table}
  634. (id, tenant_id, factory_id, status, biz_date, create_time, update_time, is_deleted, is_active,
  635. module_code, metric_code, result_status, result_reason, metric_value, target_value, status_color, trend_flag, calc_time,
  636. target_config_id, target_source, target_resolved_at)
  637. VALUES
  638. (@Id, @TenantId, @FactoryId, 'ACTIVE', @BizDate, @CalcTime, @CalcTime, 0, 1,
  639. 'S2', @MetricCode, {Admin.NET.Plugin.AiDOP.SmartOps.KpiResultWrite.FromMetricVals}@MetricValue, @TargetValue, @StatusColor, @TrendFlag, @CalcTime,
  640. @TargetConfigId, @TargetSource, @TargetResolvedAt)
  641. """,
  642. new SugarParameter("@Id", nextId),
  643. new SugarParameter("@TenantId", row.TenantId),
  644. new SugarParameter("@FactoryId", row.FactoryId),
  645. new SugarParameter("@BizDate", statDate),
  646. new SugarParameter("@CalcTime", now),
  647. new SugarParameter("@MetricCode", row.MetricCode),
  648. new SugarParameter("@MetricValue", actual),
  649. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  650. new SugarParameter("@StatusColor", status),
  651. new SugarParameter("@TrendFlag", trend),
  652. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  653. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  654. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt));
  655. }
  656. private async Task<long> InsertSyncLogAsync(long entityId, string entityName, string batchId, int rowsRead)
  657. {
  658. EnsureRunScope();
  659. await MdpSchemaAligner.ExecuteAsync(_db,
  660. """
  661. INSERT INTO mdp_sync_log
  662. (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type, sync_start, rows_read, status)
  663. VALUES (@TenantId, @EntityId, 'AIDOPDEV_MYSQL', @EntityName, @BatchId, 'FULL', 'AUTO', NOW(), @RowsRead, 'RUNNING')
  664. """,
  665. new SugarParameter("@TenantId", _runScope.TenantId),
  666. new SugarParameter("@EntityId", entityId),
  667. new SugarParameter("@EntityName", entityName),
  668. new SugarParameter("@BatchId", batchId),
  669. new SugarParameter("@RowsRead", rowsRead));
  670. return await _db.Ado.GetLongAsync(
  671. "SELECT id FROM mdp_sync_log WHERE sync_batch_id=@BatchId AND entity_id=@EntityId ORDER BY id DESC LIMIT 1",
  672. new List<SugarParameter> { new("@BatchId", batchId), new("@EntityId", entityId) });
  673. }
  674. private async Task MarkSyncLogSuccessAsync(long logId, DateTime started, int affected)
  675. {
  676. await MdpSchemaAligner.ExecuteAsync(_db,
  677. """
  678. UPDATE mdp_sync_log
  679. SET sync_end=NOW(), duration_ms=@DurationMs, rows_insert=@RowsInsert, rows_update=0, rows_skip=0, rows_error=0, status='SUCCESS'
  680. WHERE id=@Id
  681. """,
  682. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  683. new SugarParameter("@RowsInsert", affected),
  684. new SugarParameter("@Id", logId));
  685. }
  686. private async Task MarkSyncLogFailedAsync(long logId, DateTime started, string message)
  687. {
  688. try
  689. {
  690. await MdpSchemaAligner.ExecuteAsync(_db,
  691. """
  692. UPDATE mdp_sync_log
  693. SET sync_end=NOW(), duration_ms=@DurationMs, rows_error=1, status='FAILED', error_msg=@ErrorMsg
  694. WHERE id=@Id
  695. """,
  696. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  697. new SugarParameter("@ErrorMsg", Truncate(message, 1000)),
  698. new SugarParameter("@Id", logId));
  699. }
  700. catch (Exception ex)
  701. {
  702. Console.Error.WriteLine($"[S2MdpSyncTransform] MarkSyncLogFailed write failed (syncLogId={logId}): {ex.Message}");
  703. }
  704. }
  705. private async Task<long> InsertTransformRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  706. {
  707. await MdpSchemaAligner.ExecuteAsync(_db,
  708. """
  709. INSERT INTO mdp_transform_run_log
  710. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  711. VALUES (@TenantId, 'S2_MDP_SYNC_TRANSFORM', 'S2 MDP同步与KPI计算', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  712. """,
  713. new SugarParameter("@TenantId", RequireScope().TenantId),
  714. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  715. new SugarParameter("@BatchId", batchId),
  716. new SugarParameter("@StartTime", startedAt));
  717. return await _db.Ado.GetLongAsync(
  718. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  719. new List<SugarParameter> { new("@BatchId", batchId) });
  720. }
  721. private async Task MarkTransformRunSuccessAsync(long runLogId, DateTime startedAt, S2MdpSyncTransformResult result)
  722. {
  723. var finishedAt = DateTime.Now;
  724. await MdpSchemaAligner.ExecuteAsync(_db,
  725. """
  726. UPDATE mdp_transform_run_log
  727. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  728. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=@DwdRows,
  729. summary_json=@SummaryJson, update_time=CURRENT_TIMESTAMP
  730. WHERE id=@Id
  731. """,
  732. new SugarParameter("@EndTime", finishedAt),
  733. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  734. new SugarParameter("@StageRows", result.StageRows),
  735. new SugarParameter("@StandardRows", result.StandardRows),
  736. new SugarParameter("@DwdRows", result.DwdRows),
  737. new SugarParameter("@SummaryJson", BuildRunSummaryJson(result)),
  738. new SugarParameter("@Id", runLogId));
  739. }
  740. private async Task MarkTransformRunFailedAsync(long runLogId, DateTime startedAt, string message)
  741. {
  742. try
  743. {
  744. var finishedAt = DateTime.Now;
  745. await MdpSchemaAligner.ExecuteAsync(_db,
  746. """
  747. UPDATE mdp_transform_run_log
  748. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  749. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  750. WHERE id=@Id
  751. """,
  752. new SugarParameter("@EndTime", finishedAt),
  753. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  754. new SugarParameter("@ErrorMessage", Truncate(message, 2000)),
  755. new SugarParameter("@Id", runLogId));
  756. }
  757. catch (Exception ex)
  758. {
  759. Console.Error.WriteLine($"[S2MdpSyncTransform] MarkTransformRunFailed write failed (runLogId={runLogId}): {ex.Message}");
  760. }
  761. }
  762. private S2MdpSqlCommand Cmd(string sql, string batchId, DateTime now)
  763. {
  764. var scope = RequireScope();
  765. return new S2MdpSqlCommand(MdpSqlScope.InjectTenantFactory(sql), new[]
  766. {
  767. new SugarParameter("@BatchId", batchId),
  768. new SugarParameter("@Now", now),
  769. new SugarParameter("@StatDate", now.Date),
  770. new SugarParameter("@TenantId", scope.TenantId),
  771. new SugarParameter("@FactoryId", scope.FactoryId)
  772. });
  773. }
  774. private MdpRebuildScope RequireScope()
  775. {
  776. EnsureRunScope();
  777. return _runScope;
  778. }
  779. private void EnsureRunScope()
  780. {
  781. if (_runScope == null)
  782. throw new InvalidOperationException("S2 全量重算必须指定有效 tenantId/factoryId");
  783. }
  784. private static string BuildJsonObjectExpression(IEnumerable<string> columns)
  785. {
  786. var parts = columns.SelectMany(c => new[] { $"'{c.Replace("'", "''")}'", $"s.`{c}`" });
  787. return $"JSON_OBJECT({string.Join(",", parts)})";
  788. }
  789. private static string FindColumn(IEnumerable<string> columns, string expected)
  790. {
  791. return columns.First(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase));
  792. }
  793. private static string NormalizeTriggerType(string? triggerType)
  794. {
  795. return string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  796. }
  797. private static string BuildRunSummaryJson(S2MdpSyncTransformResult result)
  798. {
  799. return $$"""{"batchId":"{{result.BatchId}}","stageRows":{{result.StageRows}},"standardRows":{{result.StandardRows}},"dwdRows":{{result.DwdRows}},"kpiRows":{{result.KpiRows}}}""";
  800. }
  801. private static string ResolveKpiValueTable(int metricLevel)
  802. {
  803. return metricLevel switch
  804. {
  805. 1 => "ado_s9_kpi_value_l1_day",
  806. 2 => "ado_s9_kpi_value_l2_day",
  807. 3 => "ado_s9_kpi_value_l3_day",
  808. 4 => "ado_s9_kpi_value_l4_day",
  809. _ => "ado_s9_kpi_value_l2_day"
  810. };
  811. }
  812. private static decimal DefaultS2Target(string metricCode) => LegacyKpiCodeTargets.GetOrZero(metricCode);
  813. private static string ResolveKpiStatus(decimal actual, decimal target, string? direction, decimal? yellowThreshold, decimal? redThreshold)
  814. {
  815. if (target <= 0) return "gray";
  816. var ratio = actual / target * 100m;
  817. if (string.Equals(direction, "lower_is_better", StringComparison.OrdinalIgnoreCase))
  818. {
  819. if (actual <= target) return "green";
  820. if (ratio <= (yellowThreshold ?? 110m)) return "yellow";
  821. return ratio >= (redThreshold ?? 120m) ? "red" : "yellow";
  822. }
  823. if (actual >= target) return "green";
  824. if (ratio >= (yellowThreshold ?? 95m)) return "yellow";
  825. return ratio <= (redThreshold ?? 80m) ? "red" : "yellow";
  826. }
  827. private static string ResolveTrendFlag(decimal actual, decimal? previous)
  828. {
  829. if (previous == null) return "flat";
  830. if (actual > previous.Value) return "up";
  831. if (actual < previous.Value) return "down";
  832. return "flat";
  833. }
  834. private static string Truncate(string? raw, int maxLength)
  835. {
  836. if (string.IsNullOrEmpty(raw)) return string.Empty;
  837. return raw.Length <= maxLength ? raw : raw[..maxLength];
  838. }
  839. private sealed class S2ColumnRow
  840. {
  841. public string ColumnName { get; set; } = string.Empty;
  842. }
  843. private sealed class S2MdpEntityRow
  844. {
  845. public long Id { get; set; }
  846. public string EntityName { get; set; } = string.Empty;
  847. }
  848. private sealed class S2KpiCalcRow
  849. {
  850. public long TenantId { get; set; }
  851. public long FactoryId { get; set; }
  852. public string MetricCode { get; set; } = string.Empty;
  853. public decimal? MetricValue { get; set; }
  854. }
  855. private sealed class S2KpiMetaRow
  856. {
  857. public int MetricLevel { get; set; }
  858. public string Direction { get; set; } = "higher_is_better";
  859. public decimal? YellowThreshold { get; set; }
  860. public decimal? RedThreshold { get; set; }
  861. }
  862. private sealed class S2KpiValueRow
  863. {
  864. public long Id { get; set; }
  865. public decimal? MetricValue { get; set; }
  866. public decimal? TargetValue { get; set; }
  867. }
  868. }
  869. public sealed class S2MdpSyncTransformResult
  870. {
  871. public long RunLogId { get; set; }
  872. public string BatchId { get; set; } = string.Empty;
  873. public int StageRows { get; set; }
  874. public int StandardRows { get; set; }
  875. public int DwdRows { get; set; }
  876. public int KpiRows { get; set; }
  877. public int AtomicRows { get; set; }
  878. }
  879. internal sealed record S2MdpSqlCommand(string Sql, SugarParameter[] Parameters);
  880. internal sealed record S2MdpEntityConfig(
  881. string EntityCode,
  882. string SourceTable,
  883. string TargetTable,
  884. string SourceRowIdExpression,
  885. string SourceBizKeyExpression)
  886. {
  887. public static readonly IReadOnlyList<S2MdpEntityConfig> All = new List<S2MdpEntityConfig>
  888. {
  889. new("S2_WORK_ORDER_MASTER", "WorkOrdMaster", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrd`,''))"),
  890. new("S2_WORK_ORDER_ROUTING", "WorkOrdRouting", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrd`,''), ':', IFNULL(s.`OP`,''))"),
  891. new("S2_WORK_ORDER_DETAIL", "WorkOrdDetail", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrd`,''), ':', IFNULL(s.`Op`,''), ':', IFNULL(s.`ItemNum`,''))"),
  892. new("S2_PERIOD_SEQUENCE_DET", "PeriodSequenceDet", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrds`,''), ':', IFNULL(s.`Op`,''), ':', IFNULL(s.`Sequence`,''))"),
  893. new("S2_SCHEDULE_RESULT_OP", "ScheduleResultOpMaster", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrd`,''), ':', IFNULL(s.`Op`,''), ':', IFNULL(s.`WorkDate`,''))")
  894. };
  895. }
  896. internal static class S2MdpDdl
  897. {
  898. public static readonly IReadOnlyList<string> SqlBlocks = new[]
  899. {
  900. """
  901. INSERT INTO mdp_entity
  902. (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)
  903. SELECT 0, s.id, v.entity_code, v.entity_name, 'TABLE', v.source_table_name, 'mdp_stg_schedule', 'FULL', 5000, 1, v.remark, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  904. FROM mdp_source s
  905. JOIN (
  906. SELECT 'S2_WORK_ORDER_MASTER' AS entity_code, 'S2工单主数据' AS entity_name, 'WorkOrdMaster' AS source_table_name, '工单主数据进入 S2 贴源层' AS remark
  907. UNION ALL SELECT 'S2_WORK_ORDER_ROUTING', 'S2工单工艺路线', 'WorkOrdRouting', '工单工艺路线进入 S2 贴源层'
  908. UNION ALL SELECT 'S2_WORK_ORDER_DETAIL', 'S2工单物料明细', 'WorkOrdDetail', '工单物料需求进入 S2 贴源层'
  909. UNION ALL SELECT 'S2_PERIOD_SEQUENCE_DET', 'S2工序排程计划', 'PeriodSequenceDet', '工序间衔接与排程计划进入 S2 贴源层'
  910. UNION ALL SELECT 'S2_SCHEDULE_RESULT_OP', 'S2工序排产结果', 'ScheduleResultOpMaster', '工序排产结果进入 S2 贴源层'
  911. ) v
  912. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  913. ON DUPLICATE KEY UPDATE
  914. source_id=VALUES(source_id), entity_name=VALUES(entity_name), entity_type=VALUES(entity_type),
  915. source_table_name=VALUES(source_table_name), target_table_name=VALUES(target_table_name),
  916. sync_mode=VALUES(sync_mode), batch_size=VALUES(batch_size), status=VALUES(status),
  917. remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  918. """
  919. };
  920. }