S4MdpSyncTransformService.cs 55 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032
  1. using Admin.NET.Plugin.AiDOP.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  3. using Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  4. using Admin.NET.Plugin.AiDOP.Infrastructure;
  5. using Admin.NET.Plugin.AiDOP.SmartOps;
  6. namespace Admin.NET.Plugin.AiDOP.ProcurementExecution;
  7. /// <summary>
  8. /// S4 采购执行 MDP 转换:S4 专属 STG/STD + 消费 S3 共享 DWD,写入 dwd_s4_purchase_execution / dwd_po_trans 与 S4 KPI。
  9. /// Phase2:源→stg 经统一执行器;std/DWD/KPI 仍读 stg。
  10. /// </summary>
  11. public class S4MdpSyncTransformService : ITransient
  12. {
  13. private const string JobCode = "S4_MDP_SYNC_TRANSFORM";
  14. private readonly ISqlSugarClient _db;
  15. private readonly TransformRunLogFinalizer _runLogFinalizer;
  16. private readonly MdpModuleStagingPuller _stagingPuller;
  17. private readonly IKpiTargetResolver _kpiTargetResolver;
  18. private readonly MdpNeutralSourceGate _neutralGate;
  19. private MdpRebuildScope _runScope = null!;
  20. public S4MdpSyncTransformService(ISqlSugarClient db, MdpModuleStagingPuller stagingPuller, IKpiTargetResolver kpiTargetResolver, TransformRunLogFinalizer runLogFinalizer, MdpNeutralSourceGate neutralGate)
  21. {
  22. _db = db;
  23. _runLogFinalizer = runLogFinalizer;
  24. _stagingPuller = stagingPuller;
  25. _kpiTargetResolver = kpiTargetResolver;
  26. _neutralGate = neutralGate;
  27. }
  28. public Task<S4MdpSyncTransformResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO") =>
  29. throw new InvalidOperationException("S4 全量重算必须指定租户/工厂作用域,请走 ModuleRebuildService.Enqueue");
  30. public async Task<S4MdpSyncTransformResult> 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("S4", scope.TenantId, scope.FactoryId);
  39. cancellationToken.ThrowIfCancellationRequested();
  40. await EnsureS4TablesAsync();
  41. var now = DateTime.Now;
  42. var batchId = $"S4_MDP_FULL_{_runScope.TenantId}_{_runScope.FactoryId}_{now:yyyyMMddHHmmss}";
  43. var runLogId = await InsertTransformRunLogAsync(batchId, now, triggerType);
  44. var result = new S4MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  45. try
  46. {
  47. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Preparing, 0, 5, "准备运行环境"));
  48. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Staging, 1, 12, "正在拉取源数据"));
  49. result.StageRows = await SyncStagingAsync(batchId, now, cancellationToken, triggerType);
  50. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Staging, 1, 35, "拉取源数据完成", result.StageRows, ModuleRebuildStages.Staging));
  51. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Standard, 2, 40, "正在标准化数据"));
  52. result.StandardRows = await TransformStandardAsync(batchId, now, cancellationToken);
  53. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Standard, 2, 55, "标准化数据完成", result.StandardRows, ModuleRebuildStages.Standard));
  54. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Dwd, 3, 60, "正在生成 DWD 明细"));
  55. result.DwdRows = await BuildDwdAsync(batchId, now, cancellationToken);
  56. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Dwd, 3, 75, "生成 DWD 明细完成", result.DwdRows, ModuleRebuildStages.Dwd));
  57. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Kpi, 4, 80, "正在重算 KPI"));
  58. result.KpiRows = await BuildS4KpiValuesAsync(now, cancellationToken);
  59. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Kpi, 4, 96, "重算 KPI 完成", result.KpiRows, ModuleRebuildStages.Kpi));
  60. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Finalizing, 4, 99, "正在收尾"));
  61. await MarkTransformRunSuccessAsync(runLogId, now, result);
  62. return result;
  63. }
  64. catch (Exception ex)
  65. {
  66. if (ex is ModuleRebuildCancelledException)
  67. throw;
  68. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  69. if (!_runLogFinalizer.IsHostStopping)
  70. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  71. throw;
  72. }
  73. finally
  74. {
  75. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  76. }
  77. }
  78. private static async Task ReportAsync(Func<ModuleProgressUpdate, Task>? report, ModuleProgressUpdate update)
  79. {
  80. if (report == null) return;
  81. await report(update);
  82. }
  83. private async Task EnsureS4TablesAsync()
  84. {
  85. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_s4_iqc"));
  86. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_s4_shipment"));
  87. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_s4_return"));
  88. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_s4_shortage"));
  89. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_s4_iqc"));
  90. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_s4_shipment"));
  91. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_s4_return"));
  92. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_s4_shortage"));
  93. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("dwd_s4_purchase_execution"));
  94. }
  95. private async Task<int> SyncStagingAsync(string batchId, DateTime now, CancellationToken cancellationToken, string triggerType = "AUTO")
  96. {
  97. _ = now;
  98. EnsureRunScope();
  99. var total = await _stagingPuller.PullEntitiesAsync(
  100. S4MdpEntityConfig.All.Select(x => x.EntityCode),
  101. batchId, _runScope.TenantId, fullRefresh: ModuleRebuildTriggerType.ShouldPullFull(triggerType), taskCode: "S4_MDP_INBOUND", cancellationToken,
  102. factoryId: _runScope.FactoryId, requireMatchingSourceTenant: true);
  103. total += await CountSharedStagingRowsAsync();
  104. return total;
  105. }
  106. public async Task<S4MdpSyncTransformResult> RunInboundAsync(
  107. IEnumerable<string>? entityCodes = null,
  108. long tenantId = 0,
  109. bool fullRefresh = false,
  110. CancellationToken cancellationToken = default)
  111. {
  112. cancellationToken.ThrowIfCancellationRequested();
  113. if (tenantId <= 0)
  114. throw new InvalidOperationException("S4 inbound 必须指定有效 tenantId");
  115. _runScope = MdpRebuildScope.Create("S4", tenantId, 1);
  116. await EnsureS4TablesAsync();
  117. var now = DateTime.Now;
  118. var batchId = $"S4_MDP_IN_{_runScope.TenantId}_{now:yyyyMMddHHmmss}";
  119. var runLogId = await InsertTransformRunLogAsync(batchId, now, "INBOUND");
  120. var result = new S4MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  121. try
  122. {
  123. var codes = entityCodes?.ToList() ?? S4MdpEntityConfig.All.Select(x => x.EntityCode).ToList();
  124. result.StageRows = await _stagingPuller.PullEntitiesAsync(
  125. codes, batchId, tenantId, fullRefresh, "S4_MDP_INBOUND", cancellationToken,
  126. factoryId: _runScope.FactoryId);
  127. result.StageRows += await CountSharedStagingRowsAsync();
  128. result.StandardRows = await TransformStandardAsync(batchId, now, cancellationToken);
  129. result.DwdRows = await BuildDwdAsync(batchId, now, cancellationToken);
  130. result.KpiRows = await BuildS4KpiValuesAsync(now, cancellationToken);
  131. await MarkTransformRunSuccessAsync(runLogId, now, result);
  132. return result;
  133. }
  134. catch (Exception ex)
  135. {
  136. if (ex is ModuleRebuildCancelledException)
  137. throw;
  138. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  139. if (!_runLogFinalizer.IsHostStopping)
  140. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  141. throw;
  142. }
  143. finally
  144. {
  145. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  146. }
  147. }
  148. private async Task<int> CountSharedStagingRowsAsync()
  149. {
  150. try
  151. {
  152. return await _db.Ado.GetIntAsync(
  153. """
  154. SELECT IFNULL(SUM(cnt), 0) FROM (
  155. SELECT COUNT(1) AS cnt FROM mdp_stg_purchase_order
  156. UNION ALL SELECT COUNT(1) FROM mdp_stg_delivery
  157. UNION ALL SELECT COUNT(1) FROM mdp_stg_receipt
  158. ) t
  159. """);
  160. }
  161. catch
  162. {
  163. return 0;
  164. }
  165. }
  166. [Obsolete("Phase3: use MdpModuleStagingPuller / executors")]
  167. private async Task<int> SyncOneEntityAsync(S4MdpEntityConfig entity, string batchId, DateTime now)
  168. {
  169. var entityRow = await _db.Ado.SqlQuerySingleAsync<S4MdpEntityRow>(
  170. "SELECT id AS Id, entity_name AS EntityName FROM mdp_entity WHERE tenant_id=0 AND entity_code=@EntityCode LIMIT 1",
  171. new SugarParameter("@EntityCode", entity.EntityCode));
  172. if (entityRow == null)
  173. throw Oops.Oh($"未找到 MDP 实体配置:{entity.EntityCode},请先执行 1.0.154.sql 或启动迁移。");
  174. var tableExists = await _db.Ado.GetIntAsync(
  175. """
  176. SELECT COUNT(1) FROM information_schema.TABLES
  177. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@TableName
  178. """,
  179. new SugarParameter("@TableName", entity.SourceTable));
  180. if (tableExists == 0)
  181. {
  182. if (entity.Optional)
  183. return 0;
  184. throw Oops.Oh($"未找到 S4 源表:{entity.SourceTable}(实体 {entity.EntityCode})");
  185. }
  186. var columns = await _db.Ado.SqlQueryAsync<S4ColumnRow>(
  187. """
  188. SELECT COLUMN_NAME AS ColumnName
  189. FROM information_schema.COLUMNS
  190. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@TableName
  191. ORDER BY ORDINAL_POSITION
  192. """,
  193. new SugarParameter("@TableName", entity.SourceTable));
  194. if (columns.Count == 0)
  195. throw Oops.Oh($"未找到源表字段:{entity.SourceTable}");
  196. var names = columns.Select(u => u.ColumnName).ToList();
  197. var tenantExpr = names.Any(u => string.Equals(u, "tenant_id", StringComparison.OrdinalIgnoreCase))
  198. ? $"IFNULL(s.`{FindColumn(names, "tenant_id")}`,0)"
  199. : "0";
  200. var sourceRowExpr = names.Any(u => string.Equals(u, entity.SourceRowIdExpression, StringComparison.OrdinalIgnoreCase))
  201. ? $"s.`{FindColumn(names, entity.SourceRowIdExpression)}`"
  202. : entity.SourceRowIdExpression;
  203. var rawDataExpr = BuildJsonObjectExpression(names);
  204. var rowsRead = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM `{entity.SourceTable}`");
  205. var logId = await InsertSyncLogAsync(entityRow.Id, entityRow.EntityName, batchId, rowsRead);
  206. var started = DateTime.Now;
  207. try
  208. {
  209. var affected = await MdpSchemaAligner.ExecuteAsync(_db,
  210. $"""
  211. INSERT INTO `{entity.TargetTable}`
  212. (tenant_id, source_system, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time, process_status, raw_data)
  213. SELECT
  214. {tenantExpr},
  215. COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = @SourceTable LIMIT 1), NULLIF('',''), ''),
  216. @SourceTable,
  217. CAST({sourceRowExpr} AS CHAR),
  218. CAST(COALESCE({entity.SourceBizKeyExpression}, CAST({sourceRowExpr} AS CHAR)) AS CHAR),
  219. @BatchId,
  220. @Now,
  221. 'PENDING',
  222. {rawDataExpr}
  223. FROM `{entity.SourceTable}` s
  224. ON DUPLICATE KEY UPDATE
  225. source_row_id=VALUES(source_row_id),
  226. sync_batch_id=VALUES(sync_batch_id),
  227. sync_time=VALUES(sync_time),
  228. process_status=VALUES(process_status),
  229. raw_data=VALUES(raw_data),
  230. update_time=CURRENT_TIMESTAMP
  231. """,
  232. new SugarParameter("@SourceTable", entity.SourceTable),
  233. new SugarParameter("@BatchId", batchId),
  234. new SugarParameter("@Now", now));
  235. await MarkSyncLogSuccessAsync(logId, started, affected);
  236. return rowsRead;
  237. }
  238. catch (Exception ex)
  239. {
  240. await MarkSyncLogFailedAsync(logId, started, ex.Message);
  241. if (entity.Optional) return 0;
  242. throw;
  243. }
  244. }
  245. private async Task<int> TransformStandardAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  246. {
  247. if (!await _neutralGate.AllowsAsync(_runScope.TenantId, "SUPPLIER_DELIVERY", "DOPDEMORQ_SQLSERVER"))
  248. return await CountSharedStandardRowsAsync();
  249. var total = 0;
  250. foreach (var command in BuildStandardCommands(batchId, now))
  251. {
  252. cancellationToken.ThrowIfCancellationRequested();
  253. try
  254. {
  255. total += await MdpSchemaAligner.ExecuteAsync(_db, command.Sql, command.Parameters);
  256. }
  257. catch
  258. {
  259. // 标准层表可能尚未有贴源数据,跳过单条失败
  260. }
  261. }
  262. total += await CountSharedStandardRowsAsync();
  263. return total;
  264. }
  265. private async Task<int> CountSharedStandardRowsAsync()
  266. {
  267. try
  268. {
  269. return await _db.Ado.GetIntAsync(
  270. """
  271. SELECT IFNULL(SUM(cnt), 0) FROM (
  272. SELECT COUNT(1) AS cnt FROM mdp_std_purchase_order
  273. UNION ALL SELECT COUNT(1) FROM mdp_std_delivery_schedule
  274. UNION ALL SELECT COUNT(1) FROM mdp_std_delivery_result
  275. ) t
  276. """);
  277. }
  278. catch
  279. {
  280. return 0;
  281. }
  282. }
  283. private IEnumerable<S4MdpSqlCommand> BuildStandardCommands(string batchId, DateTime now)
  284. {
  285. yield return Cmd(
  286. """
  287. INSERT INTO mdp_std_s4_iqc
  288. (tenant_id, factory_id, source_system, po_no, po_line, supplier_code, item_code, receipt_qty, sample_qty, defect_qty, qc_result, receipt_date, source_biz_key, sync_batch_id, sync_time)
  289. SELECT tenant_id, COALESCE(NULLIF(factory_id,0),1), COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = mdp_stg_s4_iqc.source_table LIMIT 1), NULLIF(mdp_stg_s4_iqc.source_system,''), ''),
  290. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.PurOrd')), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.OrdNbr'))),
  291. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.OrdLine')), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Line'))) AS CHAR),
  292. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Supp')), ''),
  293. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ItemNum')), ''),
  294. COALESCE(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.QtyReceived')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.QtyReceived')) AS DECIMAL(18,6)) END, 0),
  295. COALESCE(CASE WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.SampleQty')), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.QcQty'))) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.SampleQty')), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.QcQty'))) AS DECIMAL(18,6)) END, 0),
  296. COALESCE(CASE WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RejectQty')), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.QtyReturn'))) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RejectQty')), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.QtyReturn'))) AS DECIMAL(18,6)) END, 0),
  297. CASE WHEN COALESCE(CASE WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RejectQty')), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.QtyReturn'))) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RejectQty')), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.QtyReturn'))) AS DECIMAL(18,6)) END, 0) > 0 THEN 'FAIL' ELSE 'PASS' END,
  298. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RcptDate')), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RctDate'))), 'null'), ''),
  299. source_biz_key, @BatchId, @Now
  300. FROM mdp_stg_s4_iqc
  301. WHERE source_table='PurOrdRctDetail'
  302. ON DUPLICATE KEY UPDATE receipt_qty=VALUES(receipt_qty), sample_qty=VALUES(sample_qty), defect_qty=VALUES(defect_qty),
  303. qc_result=VALUES(qc_result), receipt_date=VALUES(receipt_date), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  304. """, batchId, now);
  305. yield return Cmd(
  306. """
  307. INSERT INTO mdp_std_s4_shipment
  308. (tenant_id, factory_id, source_system, shipment_no, po_no, po_line, supplier_code, item_code, ship_qty, ship_date, source_biz_key, sync_batch_id, sync_time)
  309. SELECT tenant_id, COALESCE(NULLIF(factory_id,0),1), COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = mdp_stg_s4_shipment.source_table LIMIT 1), NULLIF(mdp_stg_s4_shipment.source_system,''), ''),
  310. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.shddh')), source_row_id),
  311. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.po_bill')),
  312. CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.po_billline')) AS CHAR),
  313. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.suppliercode')), ''),
  314. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.itemnum')), ''),
  315. COALESCE(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.sh_delivery_quantity')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.sh_delivery_quantity')) AS DECIMAL(18,6)) END, 0),
  316. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.updatetime')), 'null'), ''),
  317. source_biz_key, @BatchId, @Now
  318. FROM mdp_stg_s4_shipment
  319. WHERE source_table='scm_shdzb'
  320. ON DUPLICATE KEY UPDATE ship_qty=VALUES(ship_qty), ship_date=VALUES(ship_date), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  321. """, batchId, now);
  322. yield return Cmd(
  323. """
  324. INSERT INTO mdp_std_s4_return
  325. (tenant_id, factory_id, source_system, po_no, po_line, supplier_code, item_code, return_qty, return_reason, return_status, source_biz_key, sync_batch_id, sync_time)
  326. SELECT tenant_id, COALESCE(NULLIF(factory_id,0),1), COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = mdp_stg_s4_return.source_table LIMIT 1), NULLIF(mdp_stg_s4_return.source_system,''), ''),
  327. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ponumber')),
  328. CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.poline')) AS CHAR),
  329. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.suppliercode')), ''),
  330. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.itemnum')), ''),
  331. COALESCE(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.returnqty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.returnqty')) AS DECIMAL(18,6)) END, 0),
  332. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.remark')),
  333. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.status')),
  334. source_biz_key, @BatchId, @Now
  335. FROM mdp_stg_s4_return
  336. WHERE source_table='srm_polist_ds'
  337. AND COALESCE(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.returnqty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.returnqty')) AS DECIMAL(18,6)) END, 0) > 0
  338. ON DUPLICATE KEY UPDATE return_qty=VALUES(return_qty), return_reason=VALUES(return_reason), return_status=VALUES(return_status),
  339. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  340. """, batchId, now);
  341. yield return Cmd(
  342. """
  343. INSERT INTO mdp_std_s4_shortage
  344. (tenant_id, factory_id, source_system, work_order, supplier_code, item_code, shortage_qty, risk_level, need_date, source_biz_key, sync_batch_id, sync_time)
  345. SELECT tenant_id, COALESCE(NULLIF(factory_id,0),1), COALESCE((SELECT r.source_system FROM mdp_source_table_registry r WHERE r.source_table = mdp_stg_s4_shortage.source_table LIMIT 1), NULLIF(mdp_stg_s4_shortage.source_system,''), ''),
  346. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.work_order')),
  347. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.supplier_code')), ''),
  348. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.component_item_code')), ''),
  349. COALESCE(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.shortage_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.shortage_qty')) AS DECIMAL(18,6)) END, 0),
  350. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.risk_level')), 'MEDIUM'),
  351. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.expected_supply_date')), 'null'), ''),
  352. source_biz_key, @BatchId, @Now
  353. FROM mdp_stg_s4_shortage
  354. WHERE source_table='dwd_material_shortage'
  355. AND COALESCE(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.shortage_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.shortage_qty')) AS DECIMAL(18,6)) END, 0) > 0
  356. ON DUPLICATE KEY UPDATE shortage_qty=VALUES(shortage_qty), risk_level=VALUES(risk_level), need_date=VALUES(need_date),
  357. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  358. """, batchId, now);
  359. }
  360. private async Task<int> BuildDwdAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  361. {
  362. cancellationToken.ThrowIfCancellationRequested();
  363. var scope = RequireScope();
  364. string Scoped(string sql) => MdpSqlScope.InjectTenantFactory(sql);
  365. SugarParameter[] ScopedParameters(params SugarParameter[] parameters) =>
  366. [
  367. .. parameters,
  368. new SugarParameter("@TenantId", scope.TenantId),
  369. new SugarParameter("@FactoryId", scope.FactoryId)
  370. ];
  371. var statDate = now.Date;
  372. var total = 0;
  373. try
  374. {
  375. await MdpSchemaAligner.ExecuteAsync(_db,
  376. Scoped("DELETE FROM dwd_s4_purchase_execution WHERE stat_date=@StatDate"),
  377. ScopedParameters(new SugarParameter("@StatDate", statDate)));
  378. total += await MdpSchemaAligner.ExecuteAsync(_db,
  379. Scoped("""
  380. INSERT INTO dwd_s4_purchase_execution
  381. (tenant_id, factory_id, stat_date, po_no, po_line, supplier_code, item_code,
  382. order_qty, delivery_qty, received_qty, returned_qty, shortage_qty,
  383. due_date, actual_arrival_date, risk_level, source_system, sync_batch_id, sync_time)
  384. SELECT d.tenant_id, d.factory_id, @StatDate,
  385. d.po_no, d.po_line, IFNULL(d.supplier_code,''), IFNULL(d.item_code,''),
  386. IFNULL(d.order_qty,0), IFNULL(d.delivery_qty,0), IFNULL(d.receipt_qty,0),
  387. IFNULL(ret.return_qty,0), IFNULL(sh.shortage_qty,0),
  388. DATE(d.due_date), DATE(d.last_receipt_date),
  389. CASE WHEN d.risk_level IN ('HIGH','MEDIUM','LOW') THEN LOWER(d.risk_level)
  390. WHEN IFNULL(d.remaining_qty,0) > 0 THEN 'high'
  391. WHEN IFNULL(d.receipt_qty,0) >= IFNULL(d.order_qty,0) AND IFNULL(d.order_qty,0) > 0 THEN 'low'
  392. ELSE 'medium' END,
  393. 'AIDOP_NATIVE', @BatchId, @Now
  394. FROM dwd_supplier_delivery d
  395. LEFT JOIN (
  396. SELECT tenant_id, factory_id, po_no, po_line, SUM(IFNULL(return_qty,0)) AS return_qty
  397. FROM mdp_std_s4_return
  398. WHERE IFNULL(po_no,'') <> ''
  399. GROUP BY tenant_id, factory_id, po_no, po_line
  400. ) ret ON d.tenant_id=ret.tenant_id AND d.factory_id=ret.factory_id AND d.po_no=ret.po_no AND d.po_line=ret.po_line
  401. LEFT JOIN (
  402. SELECT tenant_id, factory_id, supplier_code, item_code, SUM(IFNULL(shortage_qty,0)) AS shortage_qty
  403. FROM mdp_std_s4_shortage
  404. GROUP BY tenant_id, factory_id, supplier_code, item_code
  405. ) sh ON d.tenant_id=sh.tenant_id AND d.factory_id=sh.factory_id AND IFNULL(d.supplier_code,'')=IFNULL(sh.supplier_code,'') AND IFNULL(d.item_code,'')=IFNULL(sh.item_code,'')
  406. WHERE d.stat_date=@StatDate AND IFNULL(d.po_no,'') <> ''
  407. ON DUPLICATE KEY UPDATE
  408. order_qty=VALUES(order_qty), delivery_qty=VALUES(delivery_qty), received_qty=VALUES(received_qty),
  409. returned_qty=VALUES(returned_qty), shortage_qty=VALUES(shortage_qty),
  410. due_date=VALUES(due_date), actual_arrival_date=VALUES(actual_arrival_date),
  411. risk_level=VALUES(risk_level), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  412. """),
  413. ScopedParameters(
  414. new SugarParameter("@StatDate", statDate),
  415. new SugarParameter("@BatchId", batchId),
  416. new SugarParameter("@Now", now)));
  417. }
  418. catch
  419. {
  420. // dwd_supplier_delivery 可能尚未由 S3 生成
  421. }
  422. try
  423. {
  424. await MdpSchemaAligner.ExecuteAsync(_db,
  425. Scoped("DELETE FROM dwd_po_trans WHERE trans_date=@StatDate"),
  426. ScopedParameters(new SugarParameter("@StatDate", statDate)));
  427. total += await MdpSchemaAligner.ExecuteAsync(_db,
  428. Scoped("""
  429. INSERT INTO dwd_po_trans
  430. (tenant_id, factory_id, po_no, po_line, supplier_code, item_code, order_qty, received_qty,
  431. returned_qty, shortage_qty, due_date, actual_arrival_date, risk_level,
  432. trans_date, source_system, sync_batch_id, sync_time)
  433. SELECT tenant_id, factory_id, po_no, po_line, supplier_code, item_code,
  434. order_qty, received_qty, returned_qty, shortage_qty,
  435. due_date, actual_arrival_date, risk_level,
  436. stat_date, source_system, sync_batch_id, sync_time
  437. FROM dwd_s4_purchase_execution
  438. WHERE stat_date=@StatDate
  439. """),
  440. ScopedParameters(new SugarParameter("@StatDate", statDate)));
  441. }
  442. catch
  443. {
  444. try
  445. {
  446. total += await MdpSchemaAligner.ExecuteAsync(_db,
  447. Scoped("""
  448. INSERT INTO dwd_po_trans
  449. (tenant_id, factory_id, po_no, supplier_code, item_code, order_qty, received_qty, trans_date, source_system, sync_time)
  450. SELECT po.tenant_id, po.factory_id, po.po_no, IFNULL(po.supplier_code,''), IFNULL(po.item_code,''),
  451. IFNULL(po.order_qty,0), IFNULL(po.received_qty,0), @StatDate, 'AIDOP_NATIVE', @Now
  452. FROM mdp_std_purchase_order po
  453. WHERE IFNULL(po.po_no,'') <> ''
  454. """),
  455. ScopedParameters(
  456. new SugarParameter("@StatDate", statDate),
  457. new SugarParameter("@Now", now)));
  458. }
  459. catch
  460. {
  461. // ignore
  462. }
  463. }
  464. try
  465. {
  466. total += await MdpSchemaAligner.ExecuteAsync(_db,
  467. Scoped("""
  468. INSERT INTO dwd_qc_trans
  469. (tenant_id, factory_id, item_code, supplier_code, batch_no, sample_qty, defect_qty, result, trans_date, source_system, sync_time)
  470. SELECT tenant_id, factory_id, IFNULL(item_code,''), IFNULL(supplier_code,''), source_biz_key,
  471. CAST(IFNULL(sample_qty,0) AS SIGNED), CAST(IFNULL(defect_qty,0) AS SIGNED),
  472. CASE WHEN qc_result='FAIL' THEN 'FAIL' WHEN qc_result='CONCESSION' THEN 'CONCESSION' ELSE 'PASS' END,
  473. @StatDate, 'AIDOP_NATIVE', @Now
  474. FROM mdp_std_s4_iqc
  475. WHERE IFNULL(item_code,'') <> ''
  476. ON DUPLICATE KEY UPDATE sample_qty=VALUES(sample_qty), defect_qty=VALUES(defect_qty), result=VALUES(result), sync_time=VALUES(sync_time)
  477. """),
  478. ScopedParameters(
  479. new SugarParameter("@StatDate", statDate),
  480. new SugarParameter("@Now", now)));
  481. }
  482. catch
  483. {
  484. // dwd_qc_trans 可能无唯一键,忽略
  485. }
  486. return total;
  487. }
  488. private async Task<int> BuildS4KpiValuesAsync(DateTime now, CancellationToken cancellationToken)
  489. {
  490. var affected = 0;
  491. for (var dayOffset = 13; dayOffset >= 0; dayOffset--)
  492. {
  493. var statDate = now.Date.AddDays(-dayOffset);
  494. var rows = await CalculateS4KpiValuesAsync(statDate);
  495. foreach (var row in rows)
  496. {
  497. cancellationToken.ThrowIfCancellationRequested();
  498. affected += await UpsertS4KpiValueAsync(row, statDate, now);
  499. }
  500. }
  501. return affected;
  502. }
  503. private async Task<List<S4KpiCalcRow>> CalculateS4KpiValuesAsync(DateTime statDate)
  504. {
  505. var scope = RequireScope();
  506. try
  507. {
  508. return await _db.Ado.SqlQueryAsync<S4KpiCalcRow>(
  509. MdpSqlScope.InjectTenantFactory(
  510. """
  511. SELECT ds.tenant_id AS TenantId, ds.factory_id AS FactoryId, 'S4_L1_001' AS MetricCode,
  512. ROUND(AVG(TIMESTAMPDIFF(DAY, ds.request_date, r.receipt_date)), 4) AS MetricValue
  513. FROM mdp_std_delivery_schedule ds
  514. JOIN (
  515. SELECT tenant_id, factory_id, po_no, po_line, item_code, MAX(receipt_date) AS receipt_date
  516. FROM (
  517. SELECT tenant_id, factory_id, po_no, po_line, item_code, receipt_date
  518. FROM mdp_std_delivery_result
  519. WHERE receipt_date IS NOT NULL
  520. UNION ALL
  521. SELECT tenant_id, factory_id, po_no, po_line, item_code, receipt_date
  522. FROM mdp_std_s4_iqc
  523. WHERE receipt_date IS NOT NULL
  524. ) receipt_events
  525. GROUP BY tenant_id, factory_id, po_no, po_line, item_code
  526. ) r ON r.tenant_id=ds.tenant_id AND r.factory_id=ds.factory_id
  527. AND r.po_no=ds.po_no AND r.po_line=ds.po_line AND r.item_code=ds.item_code
  528. WHERE ds.request_date IS NOT NULL AND r.receipt_date >= ds.request_date
  529. GROUP BY ds.tenant_id, ds.factory_id
  530. HAVING MetricValue IS NOT NULL
  531. UNION ALL
  532. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S4_L1_002' AS MetricCode,
  533. ROUND(100 * SUM(CASE WHEN IFNULL(d.receipt_qty,0) >= IFNULL(d.order_qty,0) AND IFNULL(d.order_qty,0) > 0 THEN 1
  534. WHEN d.delivery_status='COMPLETED' THEN 1 ELSE 0 END) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  535. FROM dwd_supplier_delivery d
  536. WHERE d.stat_date=@StatDate AND IFNULL(d.order_qty,0) > 0
  537. GROUP BY tenant_id, factory_id
  538. UNION ALL
  539. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S4_L1_003' AS MetricCode,
  540. ROUND(SUM(IFNULL(d.receipt_qty,0)) / GREATEST(COUNT(DISTINCT NULLIF(d.supplier_code,'')), 1), 4) AS MetricValue
  541. FROM dwd_supplier_delivery d
  542. WHERE d.stat_date=@StatDate
  543. GROUP BY tenant_id, factory_id
  544. UNION ALL
  545. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S4_L1_004' AS MetricCode,
  546. ROUND(AVG(CASE WHEN IFNULL(d.order_qty,0) > 0
  547. THEN (IFNULL(d.remaining_qty,0) / d.order_qty) * 30 END), 4) AS MetricValue
  548. FROM dwd_supplier_delivery d
  549. WHERE d.stat_date=@StatDate AND IFNULL(d.order_qty,0) > 0
  550. GROUP BY tenant_id, factory_id
  551. UNION ALL
  552. SELECT ds.tenant_id AS TenantId, ds.factory_id AS FactoryId, 'S4_L2_001' AS MetricCode,
  553. ROUND(AVG(TIMESTAMPDIFF(DAY, ds.request_date, r.receipt_date)), 4) AS MetricValue
  554. FROM mdp_std_delivery_schedule ds
  555. JOIN (
  556. SELECT tenant_id, factory_id, po_no, po_line, item_code, MAX(receipt_date) AS receipt_date
  557. FROM (
  558. SELECT tenant_id, factory_id, po_no, po_line, item_code, receipt_date
  559. FROM mdp_std_delivery_result
  560. WHERE receipt_date IS NOT NULL
  561. UNION ALL
  562. SELECT tenant_id, factory_id, po_no, po_line, item_code, receipt_date
  563. FROM mdp_std_s4_iqc
  564. WHERE receipt_date IS NOT NULL
  565. ) receipt_events
  566. GROUP BY tenant_id, factory_id, po_no, po_line, item_code
  567. ) r ON r.tenant_id=ds.tenant_id AND r.factory_id=ds.factory_id
  568. AND r.po_no=ds.po_no AND r.po_line=ds.po_line AND r.item_code=ds.item_code
  569. WHERE ds.request_date IS NOT NULL AND r.receipt_date >= ds.request_date
  570. GROUP BY ds.tenant_id, ds.factory_id
  571. HAVING MetricValue IS NOT NULL
  572. UNION ALL
  573. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S4_L2_002' AS MetricCode,
  574. ROUND(100 * SUM(CASE WHEN IFNULL(d.receipt_qty,0) >= IFNULL(d.order_qty,0) AND IFNULL(d.order_qty,0) > 0 THEN 1
  575. WHEN d.delivery_status='COMPLETED' THEN 1 ELSE 0 END) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  576. FROM dwd_supplier_delivery d
  577. WHERE d.stat_date=@StatDate AND IFNULL(d.order_qty,0) > 0
  578. GROUP BY tenant_id, factory_id
  579. UNION ALL
  580. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S4_L2_003' AS MetricCode,
  581. ROUND(SUM(IFNULL(pe.received_qty,0)) / GREATEST(COUNT(DISTINCT NULLIF(pe.item_code,'')), 1), 4) AS MetricValue
  582. FROM dwd_s4_purchase_execution pe
  583. WHERE pe.stat_date=@StatDate
  584. GROUP BY tenant_id, factory_id
  585. UNION ALL
  586. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S4_L2_004' AS MetricCode,
  587. ROUND(AVG(CASE WHEN IFNULL(pe.order_qty,0) > 0
  588. THEN ((IFNULL(pe.order_qty,0) - IFNULL(pe.received_qty,0)) / pe.order_qty) * 30 END), 4) AS MetricValue
  589. FROM dwd_s4_purchase_execution pe
  590. WHERE pe.stat_date=@StatDate AND IFNULL(pe.order_qty,0) > 0
  591. GROUP BY tenant_id, factory_id
  592. UNION ALL
  593. SELECT ds.tenant_id AS TenantId, ds.factory_id AS FactoryId, 'S4_L3_001' AS MetricCode,
  594. ROUND(AVG(TIMESTAMPDIFF(DAY, ds.request_date, r.receipt_date)), 4) AS MetricValue
  595. FROM mdp_std_delivery_schedule ds
  596. JOIN (
  597. SELECT tenant_id, factory_id, po_no, po_line, item_code, MAX(receipt_date) AS receipt_date
  598. FROM (
  599. SELECT tenant_id, factory_id, po_no, po_line, item_code, receipt_date
  600. FROM mdp_std_delivery_result
  601. WHERE receipt_date IS NOT NULL
  602. UNION ALL
  603. SELECT tenant_id, factory_id, po_no, po_line, item_code, receipt_date
  604. FROM mdp_std_s4_iqc
  605. WHERE receipt_date IS NOT NULL
  606. ) receipt_events
  607. GROUP BY tenant_id, factory_id, po_no, po_line, item_code
  608. ) r ON r.tenant_id=ds.tenant_id AND r.factory_id=ds.factory_id
  609. AND r.po_no=ds.po_no AND r.po_line=ds.po_line AND r.item_code=ds.item_code
  610. WHERE ds.request_date IS NOT NULL AND r.receipt_date >= ds.request_date
  611. AND IFNULL(ds.supplier_code,'') <> ''
  612. GROUP BY ds.tenant_id, ds.factory_id
  613. HAVING MetricValue IS NOT NULL
  614. UNION ALL
  615. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S4_L3_002' AS MetricCode,
  616. ROUND(100 * SUM(CASE WHEN IFNULL(d.receipt_qty,0) >= IFNULL(d.order_qty,0) AND IFNULL(d.order_qty,0) > 0 THEN 1
  617. WHEN d.delivery_status='COMPLETED' THEN 1 ELSE 0 END) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  618. FROM dwd_supplier_delivery d
  619. WHERE d.stat_date=@StatDate AND IFNULL(d.supplier_code,'') <> '' AND IFNULL(d.order_qty,0) > 0
  620. GROUP BY tenant_id, factory_id
  621. UNION ALL
  622. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S4_L3_003' AS MetricCode,
  623. ROUND(COUNT(DISTINCT NULLIF(pe.supplier_code,'')) / GREATEST(COUNT(DISTINCT NULLIF(pe.po_no,'')), 1), 4) AS MetricValue
  624. FROM dwd_s4_purchase_execution pe
  625. WHERE pe.stat_date=@StatDate AND IFNULL(pe.supplier_code,'') <> ''
  626. GROUP BY tenant_id, factory_id
  627. UNION ALL
  628. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S4_L3_004' AS MetricCode,
  629. ROUND(AVG(CASE WHEN IFNULL(pe.order_qty,0) > 0 AND IFNULL(pe.supplier_code,'') <> ''
  630. THEN ((IFNULL(pe.order_qty,0) - IFNULL(pe.received_qty,0)) / pe.order_qty) * 30 END), 4) AS MetricValue
  631. FROM dwd_s4_purchase_execution pe
  632. WHERE pe.stat_date=@StatDate AND IFNULL(pe.order_qty,0) > 0 AND IFNULL(pe.supplier_code,'') <> ''
  633. GROUP BY tenant_id, factory_id
  634. """),
  635. new SugarParameter("@StatDate", statDate),
  636. new SugarParameter("@TenantId", scope.TenantId),
  637. new SugarParameter("@FactoryId", scope.FactoryId));
  638. }
  639. catch
  640. {
  641. return new List<S4KpiCalcRow>();
  642. }
  643. }
  644. private async Task<int> UpsertS4KpiValueAsync(S4KpiCalcRow row, DateTime statDate, DateTime now)
  645. {
  646. var meta = await _db.Ado.SqlQuerySingleAsync<S4KpiMetaRow>(
  647. """
  648. SELECT MetricLevel, Direction, YellowThreshold, RedThreshold
  649. FROM ado_smart_ops_kpi_master
  650. WHERE TenantId=@TenantId AND ModuleCode='S4' AND MetricCode=@MetricCode AND IsEnabled=1
  651. LIMIT 1
  652. """,
  653. new SugarParameter("@TenantId", row.TenantId),
  654. new SugarParameter("@MetricCode", row.MetricCode));
  655. if (meta == null || row.MetricValue == null)
  656. return 0;
  657. var table = ResolveKpiValueTable(meta.MetricLevel);
  658. var current = await _db.Ado.SqlQuerySingleAsync<S4KpiValueRow>(
  659. $"""
  660. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  661. FROM {table}
  662. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S4'
  663. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  664. ORDER BY id
  665. LIMIT 1
  666. """,
  667. new SugarParameter("@TenantId", row.TenantId),
  668. new SugarParameter("@FactoryId", row.FactoryId),
  669. new SugarParameter("@MetricCode", row.MetricCode),
  670. new SugarParameter("@BizDate", statDate));
  671. var prior = await _db.Ado.SqlQuerySingleAsync<S4KpiValueRow>(
  672. $"""
  673. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  674. FROM {table}
  675. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S4'
  676. AND metric_code=@MetricCode AND biz_date<@BizDate AND is_deleted=0
  677. ORDER BY biz_date DESC, id DESC
  678. LIMIT 1
  679. """,
  680. new SugarParameter("@TenantId", row.TenantId),
  681. new SugarParameter("@FactoryId", row.FactoryId),
  682. new SugarParameter("@MetricCode", row.MetricCode),
  683. new SugarParameter("@BizDate", statDate));
  684. var actual = Math.Round(row.MetricValue.Value, 4);
  685. var snap = await _kpiTargetResolver.ResolveAsync(row.TenantId, row.FactoryId, row.MetricCode, "S4", statDate);
  686. var target = snap.TargetValue;
  687. var status = AidopS4KpiMerge.AchievementLevel(actual, target, meta.Direction, meta.YellowThreshold, meta.RedThreshold);
  688. var trend = ResolveTrendFlag(actual, prior?.MetricValue);
  689. if (current != null)
  690. {
  691. return await MdpSchemaAligner.ExecuteAsync(_db,
  692. $"""
  693. UPDATE {table}
  694. SET {Admin.NET.Plugin.AiDOP.SmartOps.KpiResultWrite.FromMetricSet}metric_value=@MetricValue, target_value=@TargetValue, status_color=@StatusColor, trend_flag=@TrendFlag,
  695. target_config_id=@TargetConfigId, target_source=@TargetSource, target_resolved_at=@TargetResolvedAt,
  696. is_active=1, status='ACTIVE', calc_time=@CalcTime, update_time=@CalcTime
  697. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S4'
  698. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  699. """,
  700. new SugarParameter("@MetricValue", actual),
  701. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  702. new SugarParameter("@StatusColor", status),
  703. new SugarParameter("@TrendFlag", trend),
  704. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  705. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  706. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt),
  707. new SugarParameter("@CalcTime", now),
  708. new SugarParameter("@TenantId", row.TenantId),
  709. new SugarParameter("@FactoryId", row.FactoryId),
  710. new SugarParameter("@MetricCode", row.MetricCode),
  711. new SugarParameter("@BizDate", statDate));
  712. }
  713. var nextId = Yitter.IdGenerator.YitIdHelper.NextId();
  714. return await MdpSchemaAligner.ExecuteAsync(_db,
  715. $"""
  716. INSERT INTO {table}
  717. (id, tenant_id, factory_id, status, biz_date, create_time, update_time, is_deleted, is_active,
  718. module_code, metric_code, result_status, result_reason, metric_value, target_value, status_color, trend_flag, calc_time,
  719. target_config_id, target_source, target_resolved_at)
  720. VALUES
  721. (@Id, @TenantId, @FactoryId, 'ACTIVE', @BizDate, @CalcTime, @CalcTime, 0, 1,
  722. 'S4', @MetricCode, {Admin.NET.Plugin.AiDOP.SmartOps.KpiResultWrite.FromMetricVals}@MetricValue, @TargetValue, @StatusColor, @TrendFlag, @CalcTime,
  723. @TargetConfigId, @TargetSource, @TargetResolvedAt)
  724. """,
  725. new SugarParameter("@Id", nextId),
  726. new SugarParameter("@TenantId", row.TenantId),
  727. new SugarParameter("@FactoryId", row.FactoryId),
  728. new SugarParameter("@BizDate", statDate),
  729. new SugarParameter("@CalcTime", now),
  730. new SugarParameter("@MetricCode", row.MetricCode),
  731. new SugarParameter("@MetricValue", actual),
  732. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  733. new SugarParameter("@StatusColor", status),
  734. new SugarParameter("@TrendFlag", trend),
  735. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  736. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  737. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt));
  738. }
  739. private async Task<long> InsertSyncLogAsync(long entityId, string entityName, string batchId, int rowsRead)
  740. {
  741. EnsureRunScope();
  742. await MdpSchemaAligner.ExecuteAsync(_db,
  743. """
  744. INSERT INTO mdp_sync_log
  745. (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type, sync_start, rows_read, status)
  746. VALUES (@TenantId, @EntityId, 'AIDOPDEV_MYSQL', @EntityName, @BatchId, 'FULL', 'AUTO', NOW(), @RowsRead, 'RUNNING')
  747. """,
  748. new SugarParameter("@TenantId", _runScope.TenantId),
  749. new SugarParameter("@EntityId", entityId),
  750. new SugarParameter("@EntityName", entityName),
  751. new SugarParameter("@BatchId", batchId),
  752. new SugarParameter("@RowsRead", rowsRead));
  753. return await _db.Ado.GetLongAsync(
  754. "SELECT id FROM mdp_sync_log WHERE sync_batch_id=@BatchId AND entity_id=@EntityId ORDER BY id DESC LIMIT 1",
  755. new List<SugarParameter> { new("@BatchId", batchId), new("@EntityId", entityId) });
  756. }
  757. private async Task MarkSyncLogSuccessAsync(long logId, DateTime started, int affected)
  758. {
  759. await MdpSchemaAligner.ExecuteAsync(_db,
  760. """
  761. UPDATE mdp_sync_log
  762. SET status='SUCCESS', sync_end=NOW(), duration_ms=@DurationMs, rows_written=@RowsWritten
  763. WHERE id=@Id
  764. """,
  765. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  766. new SugarParameter("@RowsWritten", affected),
  767. new SugarParameter("@Id", logId));
  768. }
  769. private async Task MarkSyncLogFailedAsync(long logId, DateTime started, string message)
  770. {
  771. await MdpSchemaAligner.ExecuteAsync(_db,
  772. """
  773. UPDATE mdp_sync_log
  774. SET status='FAILED', sync_end=NOW(), duration_ms=@DurationMs, error_message=@ErrorMessage
  775. WHERE id=@Id
  776. """,
  777. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  778. new SugarParameter("@ErrorMessage", Truncate(message, 2000)),
  779. new SugarParameter("@Id", logId));
  780. }
  781. private async Task<long> InsertTransformRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  782. {
  783. await MdpSchemaAligner.ExecuteAsync(_db,
  784. """
  785. INSERT INTO mdp_transform_run_log
  786. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  787. VALUES (@TenantId, @JobCode, 'S4 MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  788. """,
  789. new SugarParameter("@TenantId", RequireScope().TenantId),
  790. new SugarParameter("@JobCode", JobCode),
  791. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  792. new SugarParameter("@BatchId", batchId),
  793. new SugarParameter("@StartTime", startedAt));
  794. return await _db.Ado.GetLongAsync(
  795. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  796. new List<SugarParameter> { new("@BatchId", batchId) });
  797. }
  798. private async Task MarkTransformRunSuccessAsync(long runLogId, DateTime startedAt, S4MdpSyncTransformResult result)
  799. {
  800. var finishedAt = DateTime.Now;
  801. await MdpSchemaAligner.ExecuteAsync(_db,
  802. """
  803. UPDATE mdp_transform_run_log
  804. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  805. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=@DwdRows,
  806. summary_json=@SummaryJson, update_time=CURRENT_TIMESTAMP
  807. WHERE id=@Id
  808. """,
  809. new SugarParameter("@EndTime", finishedAt),
  810. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  811. new SugarParameter("@StageRows", result.StageRows),
  812. new SugarParameter("@StandardRows", result.StandardRows),
  813. new SugarParameter("@DwdRows", result.DwdRows),
  814. new SugarParameter("@SummaryJson", BuildRunSummaryJson(result)),
  815. new SugarParameter("@Id", runLogId));
  816. }
  817. private async Task MarkTransformRunFailedAsync(long runLogId, DateTime startedAt, string message)
  818. {
  819. try
  820. {
  821. var finishedAt = DateTime.Now;
  822. await MdpSchemaAligner.ExecuteAsync(_db,
  823. """
  824. UPDATE mdp_transform_run_log
  825. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  826. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  827. WHERE id=@Id
  828. """,
  829. new SugarParameter("@EndTime", finishedAt),
  830. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  831. new SugarParameter("@ErrorMessage", Truncate(message, 2000)),
  832. new SugarParameter("@Id", runLogId));
  833. }
  834. catch (Exception ex)
  835. {
  836. Console.Error.WriteLine($"[S4MdpSyncTransform] MarkTransformRunFailed write failed (runLogId={runLogId}): {ex.Message}");
  837. }
  838. }
  839. private S4MdpSqlCommand Cmd(string sql, string batchId, DateTime now)
  840. {
  841. var scope = RequireScope();
  842. return new S4MdpSqlCommand(MdpSqlScope.InjectTenantFactory(sql), new[]
  843. {
  844. new SugarParameter("@BatchId", batchId),
  845. new SugarParameter("@Now", now),
  846. new SugarParameter("@StatDate", now.Date),
  847. new SugarParameter("@TenantId", scope.TenantId),
  848. new SugarParameter("@FactoryId", scope.FactoryId)
  849. });
  850. }
  851. private MdpRebuildScope RequireScope()
  852. {
  853. EnsureRunScope();
  854. return _runScope;
  855. }
  856. private void EnsureRunScope()
  857. {
  858. if (_runScope == null)
  859. throw new InvalidOperationException("S4 全量重算必须指定有效 tenantId/factoryId");
  860. }
  861. private static string BuildJsonObjectExpression(IEnumerable<string> columns)
  862. {
  863. var parts = columns.SelectMany(c => new[] { $"'{c.Replace("'", "''")}'", $"s.`{c}`" });
  864. return $"JSON_OBJECT({string.Join(",", parts)})";
  865. }
  866. private static string FindColumn(IEnumerable<string> columns, string expected) =>
  867. columns.First(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase));
  868. private static string BuildRunSummaryJson(S4MdpSyncTransformResult result) =>
  869. $$"""{"batchId":"{{result.BatchId}}","stageRows":{{result.StageRows}},"standardRows":{{result.StandardRows}},"dwdRows":{{result.DwdRows}},"kpiRows":{{result.KpiRows}}}""";
  870. private static string ResolveKpiValueTable(int metricLevel) => metricLevel switch
  871. {
  872. 1 => "ado_s9_kpi_value_l1_day",
  873. 2 => "ado_s9_kpi_value_l2_day",
  874. 3 => "ado_s9_kpi_value_l3_day",
  875. 4 => "ado_s9_kpi_value_l4_day",
  876. _ => "ado_s9_kpi_value_l2_day"
  877. };
  878. private static decimal DefaultS4Target(string metricCode) => LegacyKpiCodeTargets.GetOrZero(metricCode);
  879. private static string ResolveKpiStatus(decimal actual, decimal target, string? direction, decimal? yellowThreshold, decimal? redThreshold)
  880. {
  881. if (target <= 0) return "gray";
  882. var ratio = actual / target * 100m;
  883. if (string.Equals(direction, "lower_is_better", StringComparison.OrdinalIgnoreCase))
  884. {
  885. if (actual <= target) return "green";
  886. if (ratio <= (yellowThreshold ?? 110m)) return "yellow";
  887. return ratio >= (redThreshold ?? 120m) ? "red" : "yellow";
  888. }
  889. if (actual >= target) return "green";
  890. if (ratio >= (yellowThreshold ?? 95m)) return "yellow";
  891. return ratio <= (redThreshold ?? 80m) ? "red" : "yellow";
  892. }
  893. private static string ResolveTrendFlag(decimal actual, decimal? previous)
  894. {
  895. if (previous == null) return "flat";
  896. if (actual > previous.Value) return "up";
  897. if (actual < previous.Value) return "down";
  898. return "flat";
  899. }
  900. private static string NormalizeTriggerType(string? triggerType)
  901. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  902. private static string Truncate(string? raw, int maxLength)
  903. {
  904. if (string.IsNullOrEmpty(raw)) return string.Empty;
  905. return raw.Length <= maxLength ? raw : raw[..maxLength];
  906. }
  907. private sealed class S4ColumnRow
  908. {
  909. public string ColumnName { get; set; } = string.Empty;
  910. }
  911. private sealed class S4MdpEntityRow
  912. {
  913. public long Id { get; set; }
  914. public string EntityName { get; set; } = string.Empty;
  915. }
  916. private sealed class S4KpiCalcRow
  917. {
  918. public long TenantId { get; set; }
  919. public long FactoryId { get; set; }
  920. public string MetricCode { get; set; } = string.Empty;
  921. public decimal? MetricValue { get; set; }
  922. }
  923. private sealed class S4KpiMetaRow
  924. {
  925. public int MetricLevel { get; set; }
  926. public string Direction { get; set; } = "higher_is_better";
  927. public decimal? YellowThreshold { get; set; }
  928. public decimal? RedThreshold { get; set; }
  929. }
  930. private sealed class S4KpiValueRow
  931. {
  932. public long Id { get; set; }
  933. public decimal? MetricValue { get; set; }
  934. public decimal? TargetValue { get; set; }
  935. }
  936. }
  937. public sealed class S4MdpSyncTransformResult
  938. {
  939. public long RunLogId { get; set; }
  940. public string BatchId { get; set; } = string.Empty;
  941. public int StageRows { get; set; }
  942. public int StandardRows { get; set; }
  943. public int DwdRows { get; set; }
  944. public int KpiRows { get; set; }
  945. }
  946. internal sealed record S4MdpSqlCommand(string Sql, SugarParameter[] Parameters);
  947. internal sealed record S4MdpEntityConfig(
  948. string EntityCode,
  949. string SourceTable,
  950. string TargetTable,
  951. string SourceRowIdExpression,
  952. string SourceBizKeyExpression,
  953. bool Optional = false)
  954. {
  955. public static readonly IReadOnlyList<S4MdpEntityConfig> All = new List<S4MdpEntityConfig>
  956. {
  957. new("S4_IQC_RECEIPT", "PurOrdRctDetail", "mdp_stg_s4_iqc", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`Receiver`,''), ':', IFNULL(s.`Line`,''))"),
  958. new("S4_SHIPMENT_EXEC", "scm_shdzb", "mdp_stg_s4_shipment", "id", "CONCAT(IFNULL(s.`glid`,''), ':', IFNULL(s.`id`,''))"),
  959. new("S4_RETURN_EXEC", "srm_polist_ds", "mdp_stg_s4_return", "Id", "s.`dsnum`"),
  960. new("S4_SHORTAGE_EXEC", "dwd_material_shortage", "mdp_stg_s4_shortage", "id", "CONCAT(IFNULL(s.`work_order`,''), ':', IFNULL(s.`component_item_code`,''))", Optional: true)
  961. };
  962. }