S2MdpSyncTransformService.cs 56 KB

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