S2MdpSyncTransformService.cs 50 KB

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