S7MdpSyncTransformService.cs 47 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963
  1. using Admin.NET.Core.Service;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  3. using Admin.NET.Plugin.AiDOP.Infrastructure;
  4. using Admin.NET.Plugin.AiDOP.MaterialWarehouse;
  5. using Microsoft.Extensions.Logging;
  6. using System.Text.Json;
  7. namespace Admin.NET.Plugin.AiDOP.FinishedWarehouse;
  8. /// <summary>
  9. /// S7 成品仓储 — KPI 计算与刷新转换服务。
  10. /// 双模式:读本地中立标准层(由适配器从贴源→标准),不再直连源库。
  11. /// 计算逻辑沿用方老师 v5.4 KPI J 列口径(等价改写为 MySQL 读 std)。
  12. /// 包含 KPI:S7_L1_001 订单发货周期 / S7_L1_002 订单发货满足率 / S7_L1_003 成品仓储人效。
  13. /// </summary>
  14. public class S7MdpSyncTransformService : ITransient
  15. {
  16. private readonly ISqlSugarClient _db;
  17. private readonly TransformRunLogFinalizer _runLogFinalizer;
  18. private readonly SysNoticeService _sysNoticeService;
  19. private readonly ILogger<S7MdpSyncTransformService> _logger;
  20. private readonly SmartOps.KpiCalcDispatcher _kpiCalcDispatcher;
  21. private readonly SmartOps.KpiDimensionRunService _dimensionRun;
  22. private readonly SmartOps.IKpiTargetResolver _kpiTargetResolver;
  23. private readonly InventoryMdpSyncService _inventoryMdpSync;
  24. private readonly SmartOps.S9CompositeKpiWriter _s9Composite;
  25. private readonly SmartOps.KpiResultStatusService _kpiStatus;
  26. private readonly DataPlatform.MdpTenantStdSourceResolver _stdSource;
  27. private const string JobCode = "S7_MDP_SYNC_TRANSFORM";
  28. private const string JobName = "S7 成品仓储 MDP 同步与转换";
  29. private const string ModuleCode = "S7";
  30. private const string L2ValueTable = "ado_s9_kpi_value_l2_day";
  31. // FAILURE-NOTIFICATION-1:超级管理员 superAdmin.NET(AccountType=999)
  32. private const long NoticeReceiverUserId = 1300000000101L;
  33. private const string NoticeReceiverUserName = "超级管理员";
  34. public S7MdpSyncTransformService(
  35. ISqlSugarClient db,
  36. SysNoticeService sysNoticeService,
  37. ILogger<S7MdpSyncTransformService> logger,
  38. SmartOps.KpiCalcDispatcher kpiCalcDispatcher,
  39. SmartOps.KpiDimensionRunService dimensionRun,
  40. SmartOps.IKpiTargetResolver kpiTargetResolver,
  41. InventoryMdpSyncService inventoryMdpSync,
  42. TransformRunLogFinalizer runLogFinalizer,
  43. SmartOps.S9CompositeKpiWriter s9Composite,
  44. SmartOps.KpiResultStatusService kpiStatus,
  45. DataPlatform.MdpTenantStdSourceResolver stdSource)
  46. {
  47. _db = db;
  48. _runLogFinalizer = runLogFinalizer;
  49. _sysNoticeService = sysNoticeService;
  50. _logger = logger;
  51. _kpiCalcDispatcher = kpiCalcDispatcher;
  52. _dimensionRun = dimensionRun;
  53. _kpiTargetResolver = kpiTargetResolver;
  54. _inventoryMdpSync = inventoryMdpSync;
  55. _s9Composite = s9Composite;
  56. _kpiStatus = kpiStatus;
  57. _stdSource = stdSource;
  58. }
  59. private async Task<string> SourceOfAsync(long tenantId, string stdObject, CancellationToken ct)
  60. {
  61. var row = await _stdSource.GetSourceAsync(tenantId, stdObject, ct);
  62. return row?.SourceSystem ?? "";
  63. }
  64. public async Task<S7MdpSyncTransformResult> RunFullAsync(
  65. CancellationToken cancellationToken = default,
  66. string triggerType = "AUTO",
  67. S7MdpRefreshOption? option = null)
  68. {
  69. cancellationToken.ThrowIfCancellationRequested();
  70. option ??= S7MdpRefreshOption.Default();
  71. NormalizeOption(option);
  72. var now = DateTime.Now;
  73. var batchId = $"S7_MDP_FULL_{now:yyyyMMddHHmmss}";
  74. var normalizedTrigger = NormalizeTriggerType(triggerType);
  75. var runLogId = await InsertTransformRunLogAsync(batchId, now, normalizedTrigger, option);
  76. var result = new S7MdpSyncTransformResult
  77. {
  78. BatchId = batchId,
  79. RunLogId = runLogId,
  80. TriggerType = normalizedTrigger,
  81. SourceZtid = option.SourceZtid,
  82. TargetTenantId = option.TargetTenantId,
  83. TargetFactoryId = option.TargetFactoryId,
  84. BizDate = option.BizDate,
  85. BizMonth = option.BizMonth,
  86. MonthlyPeriodStart = option.MonthlyPeriodStart,
  87. MonthlyPeriodEnd = option.MonthlyPeriodEnd
  88. };
  89. try
  90. {
  91. result.StageRows = 0;
  92. result.StandardRows = await _inventoryMdpSync.TransformTransStdFromStgAsync(
  93. option.TargetTenantId, cancellationToken);
  94. var sub25 = await BuildS7L1001OrderShipmentCycleAsync(batchId, now, option, normalizedTrigger, cancellationToken);
  95. result.MergeSub("S7_L1_001", sub25);
  96. var sub26 = await BuildS7L1002OrderShipmentFulfillmentAsync(batchId, now, option, normalizedTrigger, cancellationToken);
  97. result.MergeSub("S7_L1_002", sub26);
  98. var sub27 = await BuildS7L1003FinishedWarehouseEfficiencyAsync(batchId, now, option, normalizedTrigger, cancellationToken);
  99. result.MergeSub("S7_L1_003", sub27);
  100. var currentBizDate = option.BizDate;
  101. const int backfillDays = 14;
  102. for (var dayOffset = backfillDays - 1; dayOffset >= 0; dayOffset--)
  103. {
  104. option.BizDate = currentBizDate.AddDays(-dayOffset);
  105. foreach (var metricCode in Enumerable.Range(1, 16).Select(x => $"S7_L2_{x:000}"))
  106. {
  107. var sub = await BuildS7StageKpiAsync(
  108. metricCode, batchId, now, option, normalizedTrigger, cancellationToken);
  109. result.MergeSub(metricCode, sub);
  110. }
  111. result.MergeSub("S9_L1_001", await BuildS9L1001QualityReturnRateAsync(
  112. batchId, now, option, normalizedTrigger, cancellationToken));
  113. }
  114. option.BizDate = currentBizDate;
  115. result.KpiRows += await _s9Composite.WriteRecentAsync(
  116. option.TargetTenantId, option.TargetFactoryId, option.BizDate, cancellationToken);
  117. await MarkTransformRunSuccessAsync(runLogId, now, result);
  118. return result;
  119. }
  120. catch (Exception ex)
  121. {
  122. if (ex is ModuleRebuildCancelledException)
  123. throw;
  124. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  125. if (!_runLogFinalizer.IsHostStopping)
  126. await MarkTransformRunFailedAsync(runLogId, now, ex.Message, batchId);
  127. throw;
  128. }
  129. finally
  130. {
  131. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  132. }
  133. }
  134. // ─────────────────────────────────────────────────────────────────────────
  135. /// <summary>S7_L1_001 订单发货周期 = 最晚发货日期 - 最早 FQC 报检日期(5 表 JOIN)。</summary>
  136. private async Task<KpiBuildSubResult> BuildS7L1001OrderShipmentCycleAsync(
  137. string batchId, DateTime now, S7MdpRefreshOption option, string triggerType, CancellationToken ct)
  138. {
  139. var sub = new KpiBuildSubResult();
  140. const string sql = @"
  141. select order_no as order_no,
  142. datediff(max(ship_time), min(start_time)) as cycle_days
  143. from (
  144. select b.order_no as order_no, b.item_code as item_code,
  145. ifnull(c.apply_time, b.release_time) as start_time,
  146. (case when b.closed_flag=1 then b.closed_time else d.approved_time end) as ship_time,
  147. (case when b.closed_flag=1 or b.qty_planned<=b.qty_completed then 1 else 0 end) as completed
  148. from mdp_std_work_order_line b
  149. left join (
  150. select tenant_id as tenant_id, ref_task_no, item_num,
  151. max(approved_time) as approved_time, sum(abs(qty_change)) as qty_change
  152. from mdp_std_inv_trans
  153. where tenant_id=@tenantId and biz_doc_type='SALES_SHIP'
  154. and summary_flag=0 and void_flag=0 and approved_flag=1 AND source_system=@sourceSystem
  155. group by tenant_id, ref_task_no, item_num
  156. ) d on b.task_no=d.ref_task_no and b.item_code=d.item_num and d.tenant_id=b.tenant_id
  157. left join mdp_std_fqc_task c on c.sales_order_line=b.source_row_id and c.apply_flag=1 and c.tenant_id=b.tenant_id
  158. where b.tenant_id=@tenantId and b.doc_type='SALES_ORDER' and b.void_flag=0 and b.approved_flag=1
  159. ) n
  160. group by order_no
  161. having min(completed)=1";
  162. var rows = await _db.Ado.SqlQueryAsync<S7CycleRow>(sql, new[]
  163. {
  164. new SugarParameter("@tenantId", option.TargetTenantId),
  165. new SugarParameter("@sourceDomain", option.SourceZtid),
  166. new SugarParameter("@sourceSystem", await SourceOfAsync(option.TargetTenantId, "INV_TRANS", ct))
  167. });
  168. sub.T8Rows = rows.Count;
  169. var dwdAffected = 0;
  170. var cycleList = new List<int>();
  171. foreach (var r in rows)
  172. {
  173. ct.ThrowIfCancellationRequested();
  174. if (string.IsNullOrEmpty(r.order_no)) continue;
  175. if (r.cycle_days.HasValue) cycleList.Add(r.cycle_days.Value);
  176. dwdAffected += await _db.Ado.ExecuteCommandAsync(@"
  177. INSERT INTO dwd_order_shipment_cycle
  178. (tenant_id, factory_id, biz_date, source_ztid, order_no, cycle_days, batch_id, create_time)
  179. VALUES
  180. (@tenantId, @factoryId, @bizDate, @sourceDomain, @orderNo, @cycleDays, @batchId, @now)
  181. ON DUPLICATE KEY UPDATE
  182. cycle_days=VALUES(cycle_days),
  183. batch_id=VALUES(batch_id), update_time=@now",
  184. new SugarParameter("@tenantId", option.TargetTenantId),
  185. new SugarParameter("@factoryId", option.TargetFactoryId),
  186. new SugarParameter("@bizDate", option.BizDate),
  187. new SugarParameter("@sourceDomain", option.SourceZtid),
  188. new SugarParameter("@orderNo", r.order_no),
  189. new SugarParameter("@cycleDays", r.cycle_days),
  190. new SugarParameter("@batchId", batchId),
  191. new SugarParameter("@now", now));
  192. }
  193. sub.DwdRows = dwdAffected;
  194. // 数据准备(dwd 逐单 cycle_days)已完成。最终聚合交分发器(日 KPI:period 复用 BizDate,SQL 按 biz_date 圈选)。
  195. decimal? legacyValue = cycleList.Count > 0
  196. ? Math.Round((decimal)cycleList.Average(), 4)
  197. : null;
  198. var legacyDenom = cycleList.Count > 0 ? "OK" : "NO_COMPLETED_ORDER";
  199. var dispatch = await _kpiCalcDispatcher.DispatchAsync(
  200. "S7_L1_001", ModuleCode, option.TargetTenantId, option.TargetFactoryId,
  201. option.BizDate, option.BizDate, option.BizDate, option.SourceZtid,
  202. batchId, triggerType, legacyValue, legacyDenom, ct);
  203. sub.KpiRows = dispatch.ShouldUpsert
  204. ? await UpsertKpiValueAsync("S7_L1_001", option.BizDate, dispatch.MetricValue, now, option)
  205. : 0;
  206. sub.DenominatorStatus = dispatch.DenominatorStatus;
  207. // 调度:SUMMARY 成功/NO_DATA 后触发对应 DIMENSION 跑批(共享 BatchId;SUMMARY FAILED 不触发)。
  208. // 维度失败不影响汇总链路(内部已落 dimension_run_log)。
  209. if (dispatch.ShouldUpsert)
  210. {
  211. try
  212. {
  213. await _dimensionRun.RunDimensionAsync(
  214. "S7_L1_001", ModuleCode, option.TargetTenantId, option.BizDate, batchId, triggerType, ct);
  215. }
  216. catch (Exception ex)
  217. {
  218. _logger.LogWarning(ex, "S7_L1_001 维度跑批异常(不影响汇总链路)");
  219. }
  220. }
  221. return sub;
  222. }
  223. /// <summary>S7_L1_002 订单发货满足率 = (交期前发货行数 / 该订单总行数) × 100%。</summary>
  224. private async Task<KpiBuildSubResult> BuildS7L1002OrderShipmentFulfillmentAsync(
  225. string batchId, DateTime now, S7MdpRefreshOption option, string triggerType, CancellationToken ct)
  226. {
  227. var sub = new KpiBuildSubResult();
  228. const string sql = @"
  229. select order_no as order_no,
  230. count(order_no) as total_rows,
  231. sum(completed) as in_window_rows
  232. from (
  233. select b.order_no as order_no, b.task_no as task_no, b.item_code as item_code,
  234. (case when sum(abs(d.qty_change))>=b.qty_planned then 1 else 0 end) as completed
  235. from mdp_std_work_order_line b
  236. left join (
  237. select tenant_id as tenant_id, ref_task_no, item_num,
  238. date(approved_time) as approved_date, qty_change
  239. from mdp_std_inv_trans
  240. where tenant_id=@tenantId and biz_doc_type='SALES_SHIP'
  241. and summary_flag=0 and void_flag=0 and approved_flag=1 AND source_system=@sourceSystem
  242. ) d on b.task_no=d.ref_task_no and b.item_code=d.item_num
  243. and d.approved_date<=b.plan_finish_date and d.tenant_id=b.tenant_id
  244. where b.tenant_id=@tenantId and b.doc_type='SALES_ORDER' and b.void_flag=0 and b.approved_flag=1
  245. group by b.order_no, b.task_no, b.item_code, b.qty_planned
  246. ) n
  247. group by order_no";
  248. var rows = await _db.Ado.SqlQueryAsync<S7FulfillmentRow>(sql, new[]
  249. {
  250. new SugarParameter("@tenantId", option.TargetTenantId),
  251. new SugarParameter("@sourceDomain", option.SourceZtid),
  252. new SugarParameter("@sourceSystem", await SourceOfAsync(option.TargetTenantId, "INV_TRANS", ct))
  253. });
  254. sub.T8Rows = rows.Count;
  255. var dwdAffected = 0;
  256. var rateList = new List<decimal>();
  257. foreach (var r in rows)
  258. {
  259. ct.ThrowIfCancellationRequested();
  260. if (string.IsNullOrEmpty(r.order_no)) continue;
  261. decimal? rate = (r.total_rows > 0)
  262. ? Math.Round((decimal)r.in_window_rows / r.total_rows * 100m, 4)
  263. : null;
  264. if (rate.HasValue) rateList.Add(rate.Value);
  265. dwdAffected += await _db.Ado.ExecuteCommandAsync(@"
  266. INSERT INTO dwd_order_shipment_fulfillment
  267. (tenant_id, factory_id, biz_date, source_ztid, order_no,
  268. total_rows, in_window_rows, fulfillment_rate, batch_id, create_time)
  269. VALUES
  270. (@tenantId, @factoryId, @bizDate, @sourceDomain, @orderNo,
  271. @total, @inWindow, @rate, @batchId, @now)
  272. ON DUPLICATE KEY UPDATE
  273. total_rows=VALUES(total_rows), in_window_rows=VALUES(in_window_rows),
  274. fulfillment_rate=VALUES(fulfillment_rate),
  275. batch_id=VALUES(batch_id), update_time=@now",
  276. new SugarParameter("@tenantId", option.TargetTenantId),
  277. new SugarParameter("@factoryId", option.TargetFactoryId),
  278. new SugarParameter("@bizDate", option.BizDate),
  279. new SugarParameter("@sourceDomain", option.SourceZtid),
  280. new SugarParameter("@orderNo", r.order_no),
  281. new SugarParameter("@total", r.total_rows),
  282. new SugarParameter("@inWindow", r.in_window_rows),
  283. new SugarParameter("@rate", rate),
  284. new SugarParameter("@batchId", batchId),
  285. new SugarParameter("@now", now));
  286. }
  287. sub.DwdRows = dwdAffected;
  288. // 数据准备(dwd 逐单 fulfillment_rate,已 ×100)已完成。最终聚合交分发器(CONFIG_SQL 直接 AVG 该列不再 ×100)。
  289. decimal? legacyValue = rateList.Count > 0
  290. ? Math.Round(rateList.Average(), 4)
  291. : null;
  292. var legacyDenom = rateList.Count > 0 ? "OK" : "NO_VALID_ORDER";
  293. var dispatch = await _kpiCalcDispatcher.DispatchAsync(
  294. "S7_L1_002", ModuleCode, option.TargetTenantId, option.TargetFactoryId,
  295. option.BizDate, option.BizDate, option.BizDate, option.SourceZtid,
  296. batchId, triggerType, legacyValue, legacyDenom, ct);
  297. sub.KpiRows = dispatch.ShouldUpsert
  298. ? await UpsertKpiValueAsync("S7_L1_002", option.BizDate, dispatch.MetricValue, now, option)
  299. : 0;
  300. sub.DenominatorStatus = dispatch.DenominatorStatus;
  301. // 调度:SUMMARY 成功/NO_DATA 后触发对应 DIMENSION 跑批(共享 BatchId;SUMMARY FAILED 不触发)。
  302. if (dispatch.ShouldUpsert)
  303. {
  304. try
  305. {
  306. await _dimensionRun.RunDimensionAsync(
  307. "S7_L1_002", ModuleCode, option.TargetTenantId, option.BizDate, batchId, triggerType, ct);
  308. }
  309. catch (Exception ex)
  310. {
  311. _logger.LogWarning(ex, "S7_L1_002 维度跑批异常(不影响汇总链路)");
  312. }
  313. }
  314. return sub;
  315. }
  316. /// <summary>S7_L1_003 成品仓储人效 = SALES_SHIP 数量 / WAREHOUSE ACTIVE 人数。</summary>
  317. private async Task<KpiBuildSubResult> BuildS7L1003FinishedWarehouseEfficiencyAsync(
  318. string batchId, DateTime now, S7MdpRefreshOption option, string triggerType, CancellationToken ct)
  319. {
  320. var sub = new KpiBuildSubResult();
  321. const string sqlNumer = @"
  322. select ref_task_no as task_no, item_num as item_code,
  323. date(approved_time) as approved_date, abs(qty_change) as qty_change
  324. from mdp_std_inv_trans
  325. where tenant_id=@tenantId and biz_doc_type='SALES_SHIP'
  326. and summary_flag=0 and void_flag=0 and approved_flag=1 AND source_system=@sourceSystem
  327. and date(approved_time) between @startDateText and @endDateText";
  328. const string sqlDenom = @"
  329. select count(*) as penum
  330. from mdp_std_employee
  331. where tenant_id=@tenantId and employment_status='ACTIVE' and position_code='WAREHOUSE' AND source_system=@empSource";
  332. var pNumer = new[]
  333. {
  334. new SugarParameter("@tenantId", option.TargetTenantId),
  335. new SugarParameter("@startDateText", option.MonthlyPeriodStart.ToString("yyyy-MM-dd")),
  336. new SugarParameter("@endDateText", option.MonthlyPeriodEnd.ToString("yyyy-MM-dd")),
  337. new SugarParameter("@sourceDomain", option.SourceZtid),
  338. new SugarParameter("@sourceSystem", await SourceOfAsync(option.TargetTenantId, "INV_TRANS", ct))
  339. };
  340. var pDenom = new[]
  341. {
  342. new SugarParameter("@tenantId", option.TargetTenantId),
  343. new SugarParameter("@sourceDomain", option.SourceZtid),
  344. new SugarParameter("@empSource", await SourceOfAsync(option.TargetTenantId, "EMPLOYEE", ct))
  345. };
  346. var numerRows = await _db.Ado.SqlQueryAsync<S7ShipmentDetailRow>(sqlNumer, pNumer);
  347. var denomRows = await _db.Ado.SqlQueryAsync<S7PeNumRow>(sqlDenom, pDenom);
  348. sub.T8Rows = numerRows.Count + denomRows.Count;
  349. decimal? shipmentQty = numerRows.Sum(r => r.qty_change ?? 0m);
  350. if (numerRows.Count == 0) shipmentQty = null;
  351. int? headcount = denomRows.FirstOrDefault()?.penum;
  352. decimal? efficiency = null;
  353. string denomStatus;
  354. if (!headcount.HasValue || headcount.Value <= 0)
  355. denomStatus = "NO_HEADCOUNT";
  356. else if (!shipmentQty.HasValue)
  357. denomStatus = "NO_NUMERATOR";
  358. else
  359. {
  360. efficiency = Math.Round(shipmentQty.Value / headcount.Value, 4);
  361. denomStatus = "OK";
  362. }
  363. sub.DenominatorStatus = denomStatus;
  364. var dwdAffected = await _db.Ado.ExecuteCommandAsync(@"
  365. INSERT INTO dwd_finished_warehouse_efficiency
  366. (tenant_id, factory_id, biz_month, source_ztid, period_start, period_end,
  367. shipment_qty, warehouse_headcount, efficiency, denominator_status, batch_id, create_time)
  368. VALUES
  369. (@tenantId, @factoryId, @bizMonth, @sourceDomain, @periodStart, @periodEnd,
  370. @shipmentQty, @headcount, @efficiency, @denomStatus, @batchId, @now)
  371. ON DUPLICATE KEY UPDATE
  372. period_start=VALUES(period_start), period_end=VALUES(period_end),
  373. shipment_qty=VALUES(shipment_qty), warehouse_headcount=VALUES(warehouse_headcount),
  374. efficiency=VALUES(efficiency), denominator_status=VALUES(denominator_status),
  375. batch_id=VALUES(batch_id), update_time=@now",
  376. new SugarParameter("@tenantId", option.TargetTenantId),
  377. new SugarParameter("@factoryId", option.TargetFactoryId),
  378. new SugarParameter("@bizMonth", option.BizMonth),
  379. new SugarParameter("@sourceDomain", option.SourceZtid),
  380. new SugarParameter("@periodStart", option.MonthlyPeriodStart),
  381. new SugarParameter("@periodEnd", option.MonthlyPeriodEnd),
  382. new SugarParameter("@shipmentQty", shipmentQty),
  383. new SugarParameter("@headcount", headcount),
  384. new SugarParameter("@efficiency", efficiency),
  385. new SugarParameter("@denomStatus", denomStatus),
  386. new SugarParameter("@batchId", batchId),
  387. new SugarParameter("@now", now));
  388. sub.DwdRows = dwdAffected;
  389. // 月度 KPI 最终聚合交分发器;bizDate=月末、period=当月窗口;
  390. // legacyValue=efficiency、legacyDenom 保留 NO_HEADCOUNT/NO_NUMERATOR(CONFIG_SQL 下塌缩为 NO_DATA)。
  391. var dispatch = await _kpiCalcDispatcher.DispatchAsync(
  392. "S7_L1_003", ModuleCode, option.TargetTenantId, option.TargetFactoryId,
  393. option.MonthlyPeriodEnd, option.MonthlyPeriodStart, option.MonthlyPeriodEnd, option.SourceZtid,
  394. batchId, triggerType, efficiency, denomStatus, ct);
  395. sub.KpiRows = dispatch.ShouldUpsert
  396. ? await UpsertKpiValueAsync("S7_L1_003", option.MonthlyPeriodEnd, dispatch.MetricValue, now, option)
  397. : 0;
  398. sub.DenominatorStatus = dispatch.DenominatorStatus;
  399. // 调度:SUMMARY 成功/NO_DATA 后触发人效月度 DIMENSION 跑批(月度:@biz_date=月末派生 biz_month)。
  400. if (dispatch.ShouldUpsert)
  401. {
  402. try
  403. {
  404. await _dimensionRun.RunDimensionAsync(
  405. "S7_L1_003", ModuleCode, option.TargetTenantId, option.MonthlyPeriodEnd, batchId, triggerType, ct);
  406. }
  407. catch (Exception ex)
  408. {
  409. _logger.LogWarning(ex, "S7_L1_003 维度跑批异常(不影响汇总链路)");
  410. }
  411. }
  412. return sub;
  413. }
  414. private Task<KpiBuildSubResult> BuildS7StageKpiAsync(
  415. string metricCode, string batchId, DateTime now, S7MdpRefreshOption option,
  416. string triggerType, CancellationToken ct)
  417. {
  418. var (sql, emptyDenom) = metricCode switch
  419. {
  420. "S7_L2_001" => (
  421. """
  422. SELECT ROUND(AVG(TIMESTAMPDIFF(MINUTE,q.FAPPLYTIME,q.jywcsj)/1440),4) AS MetricValue,
  423. COUNT(*) AS RowCount
  424. FROM qms_fqcbj q
  425. WHERE q.tenant_id=@TenantId
  426. AND q.FAPPLYTIME IS NOT NULL AND q.jywcsj>=q.FAPPLYTIME
  427. AND UPPER(IFNULL(q.FINSPECTSTATUS,'')) IN ('检验完成','COMPLETED','CLOSED')
  428. """,
  429. "NO_COMPLETED_FQC"),
  430. "S7_L2_002" => (
  431. """
  432. SELECT ROUND(
  433. 100 * SUM(CASE WHEN q.jywcsj IS NOT NULL AND q.jywcsj<=w.DueDate THEN 1 ELSE 0 END)
  434. / NULLIF(COUNT(*),0),
  435. 4) AS MetricValue,
  436. COUNT(*) AS RowCount
  437. FROM qms_fqcbj q
  438. INNER JOIN WorkOrdMaster w
  439. ON w.tenant_id=q.tenant_id AND w.WorkOrd=q.sczld
  440. WHERE q.tenant_id=@TenantId AND w.DueDate IS NOT NULL
  441. """,
  442. "NO_FQC_REQUIRED_DATE"),
  443. "S7_L2_003" => (
  444. """
  445. SELECT ROUND(
  446. SUM(CASE WHEN q.jywcsj IS NOT NULL
  447. AND UPPER(IFNULL(q.FINSPECTSTATUS,'')) IN ('检验完成','COMPLETED','CLOSED')
  448. THEN 1 ELSE 0 END)
  449. / NULLIF(COUNT(DISTINCT NULLIF(TRIM(q.jyfzr),'')),0),
  450. 4) AS MetricValue,
  451. SUM(CASE WHEN q.jywcsj IS NOT NULL THEN 1 ELSE 0 END) AS RowCount
  452. FROM qms_fqcbj q
  453. WHERE q.tenant_id=@TenantId
  454. """,
  455. "NO_FQC_INSPECTOR"),
  456. "S7_L2_004" => (WarehouseTurnoverSql("FG_FQC_RELEASE", "FG_PROD_RECEIPT"), "NO_FINISHED_RECEIPT_COST"),
  457. "S7_L2_005" => (WarehouseCycleSql("FG_FQC_RELEASE", "FG_PUTAWAY"), "NO_FINISHED_PUTAWAY_CYCLE"),
  458. "S7_L2_006" => (WarehouseSatisfactionSql("FG_PUTAWAY"), "NO_FINISHED_PUTAWAY_REQUIRED_DATE"),
  459. "S7_L2_007" => (WarehouseEfficiencySql("FG_PUTAWAY"), "NO_FINISHED_PUTAWAY_OPERATOR"),
  460. "S7_L2_008" => (WarehouseTurnoverSql("FG_PUTAWAY", "FG_PROD_RECEIPT"), "NO_FINISHED_RECEIPT_COST"),
  461. "S7_L2_009" => (WarehouseCycleSql("FG_PUTAWAY", "FG_PICK"), "NO_FINISHED_PICK_CYCLE"),
  462. "S7_L2_010" => (WarehouseSatisfactionSql("FG_PICK"), "NO_FINISHED_PICK_REQUIRED_DATE"),
  463. "S7_L2_011" => (WarehouseEfficiencySql("FG_PICK"), "NO_FINISHED_PICK_OPERATOR"),
  464. "S7_L2_012" => (WarehouseTurnoverSql("FG_PICK", "FG_SHIP"), "NO_FINISHED_SHIPMENT_COST"),
  465. "S7_L2_013" => (WarehouseCycleSql("FG_SHIP", "FG_RECEIPT"), "NO_FINISHED_DELIVERY_CYCLE"),
  466. "S7_L2_014" => (
  467. """
  468. SELECT ROUND(100 * SUM(completed) / NULLIF(COUNT(*),0),4) AS MetricValue,
  469. COUNT(*) AS RowCount
  470. FROM (
  471. SELECT b.order_no, b.task_no, b.item_code,
  472. CASE WHEN SUM(ABS(IFNULL(d.qty_change,0)))>=b.qty_planned THEN 1 ELSE 0 END AS completed
  473. FROM mdp_std_work_order_line b
  474. LEFT JOIN (
  475. SELECT tenant_id, ref_task_no, item_num, qty_change, DATE(approved_time) AS ship_date
  476. FROM mdp_std_inv_trans
  477. WHERE tenant_id=@TenantId
  478. AND biz_doc_type='SALES_SHIP' AND summary_flag=0 AND void_flag=0 AND approved_flag=1 AND source_system=@sourceSystem
  479. ) d ON d.tenant_id=b.tenant_id AND d.ref_task_no=b.task_no
  480. AND d.item_num=b.item_code AND d.ship_date<=b.plan_finish_date
  481. WHERE b.tenant_id=@TenantId
  482. AND b.doc_type='SALES_ORDER' AND b.void_flag=0 AND b.approved_flag=1
  483. GROUP BY b.order_no, b.task_no, b.item_code, b.qty_planned
  484. ) x
  485. """,
  486. "NO_SHIPMENT_NOTICE"),
  487. "S7_L2_015" => (WarehouseEfficiencySql("FG_SHIP"), "NO_FINISHED_SHIPMENT_OPERATOR"),
  488. "S7_L2_016" => (WarehouseTurnoverSql("FG_SHIP", "FG_RECEIPT"), "NO_SIGNED_RECEIPT_COST"),
  489. _ => throw new ArgumentOutOfRangeException(nameof(metricCode), metricCode, "不支持的 S7 阶段指标")
  490. };
  491. return DispatchS7StageKpiAsync(
  492. metricCode, sql, emptyDenom, batchId, now, option, triggerType, ct);
  493. }
  494. private static string WarehouseCycleSql(string fromStage, string toStage) =>
  495. $"""
  496. SELECT ROUND(AVG(TIMESTAMPDIFF(MINUTE,a.trans_time,b.trans_time)/1440),4) AS MetricValue,
  497. COUNT(*) AS RowCount
  498. FROM mdp_std_inv_trans a
  499. INNER JOIN mdp_std_inv_trans b
  500. ON b.tenant_id=a.tenant_id AND b.source_system=a.source_system
  501. AND b.item_num=a.item_num AND b.lot_serial=a.lot_serial
  502. AND b.trans_type='{toStage}'
  503. WHERE a.tenant_id=@TenantId AND a.trans_type='{fromStage}'
  504. AND a.source_system=@sourceSystem
  505. AND a.trans_time IS NOT NULL AND b.trans_time>=a.trans_time
  506. """;
  507. private static string WarehouseSatisfactionSql(string stage) =>
  508. $"""
  509. SELECT ROUND(100 * SUM(CASE WHEN trans_time<=eff_date THEN 1 ELSE 0 END)
  510. / NULLIF(COUNT(*),0),4) AS MetricValue,
  511. COUNT(*) AS RowCount
  512. FROM mdp_std_inv_trans
  513. WHERE tenant_id=@TenantId AND trans_type='{stage}'
  514. AND source_system=@sourceSystem
  515. AND trans_time IS NOT NULL AND eff_date IS NOT NULL
  516. """;
  517. private static string WarehouseEfficiencySql(string stage) =>
  518. $"""
  519. SELECT ROUND(COUNT(DISTINCT NULLIF(lot_serial,''))
  520. / NULLIF(COUNT(DISTINCT NULLIF(TRIM(create_user),'')),0),4) AS MetricValue,
  521. COUNT(*) AS RowCount
  522. FROM mdp_std_inv_trans
  523. WHERE tenant_id=@TenantId AND trans_type='{stage}'
  524. AND source_system=@sourceSystem
  525. """;
  526. private static string WarehouseTurnoverSql(string inventoryStage, string flowStage) =>
  527. $"""
  528. SELECT ROUND(
  529. 30 * SUM(CASE WHEN trans_type='{inventoryStage}'
  530. THEN IFNULL(end_balance,0)
  531. * IFNULL(CAST(NULLIF(dimension1,'') AS DECIMAL(18,4)),1)
  532. ELSE 0 END)
  533. / NULLIF(SUM(CASE WHEN trans_type='{flowStage}'
  534. THEN ABS(IFNULL(qty_change,0))
  535. * IFNULL(CAST(NULLIF(dimension1,'') AS DECIMAL(18,4)),1)
  536. ELSE 0 END),0),
  537. 4) AS MetricValue,
  538. SUM(CASE WHEN trans_type='{inventoryStage}' THEN 1 ELSE 0 END) AS RowCount
  539. FROM mdp_std_inv_trans
  540. WHERE tenant_id=@TenantId AND trans_type IN ('{inventoryStage}','{flowStage}')
  541. AND source_system=@sourceSystem
  542. """;
  543. private async Task<KpiBuildSubResult> DispatchS7StageKpiAsync(
  544. string metricCode, string sql, string emptyDenom, string batchId, DateTime now,
  545. S7MdpRefreshOption option, string triggerType, CancellationToken ct)
  546. {
  547. var row = await _db.Ado.SqlQuerySingleAsync<S7StageKpiRow>(
  548. sql,
  549. new SugarParameter("@TenantId", option.TargetTenantId),
  550. new SugarParameter("@sourceDomain", option.SourceZtid),
  551. new SugarParameter("@sourceSystem", await SourceOfAsync(option.TargetTenantId, "INV_TRANS", ct)));
  552. var value = row?.RowCount > 0 ? row.MetricValue : null;
  553. var dispatch = await _kpiCalcDispatcher.DispatchAsync(
  554. metricCode, ModuleCode, option.TargetTenantId, option.TargetFactoryId,
  555. option.BizDate, option.BizDate, option.BizDate, option.SourceZtid,
  556. batchId, triggerType, value, value.HasValue ? "OK" : emptyDenom, ct);
  557. var sub = new KpiBuildSubResult
  558. {
  559. T8Rows = row?.RowCount ?? 0,
  560. KpiRows = dispatch.ShouldUpsert
  561. ? await UpsertKpiValueAsync(
  562. metricCode, option.BizDate, dispatch.MetricValue, now, option,
  563. valueTable: L2ValueTable)
  564. : 0,
  565. DenominatorStatus = dispatch.DenominatorStatus
  566. };
  567. if (dispatch.ShouldUpsert)
  568. {
  569. try
  570. {
  571. await _dimensionRun.RunDimensionAsync(
  572. metricCode, ModuleCode, option.TargetTenantId, option.BizDate, batchId, triggerType, ct);
  573. }
  574. catch (Exception ex)
  575. {
  576. _logger.LogWarning(ex, "{MetricCode} 维度跑批异常(不影响汇总链路)", metricCode);
  577. }
  578. }
  579. return sub;
  580. }
  581. // ─────────────────────────────────────────────────────────────────────────
  582. /// <summary>S9_L1_001 质量退货率 = SALES_RETURN 数量 / SALES_SHIP 数量 × 1,000,000 PPM。</summary>
  583. private async Task<KpiBuildSubResult> BuildS9L1001QualityReturnRateAsync(
  584. string batchId, DateTime now, S7MdpRefreshOption option, string triggerType, CancellationToken ct)
  585. {
  586. var sub = new KpiBuildSubResult();
  587. var summary = await _db.Ado.SqlQuerySingleAsync<S7QualityReturnSummaryRow>(
  588. """
  589. SELECT
  590. SUM(CASE WHEN biz_doc_type='SALES_RETURN' THEN ABS(IFNULL(qty_change,0)) ELSE 0 END) AS return_qty,
  591. SUM(CASE WHEN biz_doc_type='SALES_SHIP' THEN ABS(IFNULL(qty_change,0)) ELSE 0 END) AS shipped_qty
  592. FROM mdp_std_inv_trans
  593. WHERE tenant_id=@TenantId
  594. AND summary_flag=0 AND void_flag=0 AND approved_flag=1 AND source_system=@sourceSystem
  595. AND biz_doc_type IN ('SALES_RETURN','SALES_SHIP')
  596. AND approved_time BETWEEN @PeriodStart AND @PeriodEnd
  597. """,
  598. new SugarParameter("@TenantId", option.TargetTenantId),
  599. new SugarParameter("@PeriodStart", option.MonthlyPeriodStart),
  600. new SugarParameter("@PeriodEnd", option.MonthlyPeriodEnd),
  601. new SugarParameter("@sourceDomain", option.SourceZtid),
  602. new SugarParameter("@sourceSystem", await SourceOfAsync(option.TargetTenantId, "INV_TRANS", ct)));
  603. sub.T8Rows = summary == null ? 0 : 1;
  604. decimal? ppm = summary?.shipped_qty > 0
  605. ? Math.Round((summary.return_qty ?? 0m) / summary.shipped_qty.Value * 1_000_000m, 4)
  606. : null;
  607. var denominatorStatus = summary?.shipped_qty > 0 ? "OK" : "NO_SHIPMENT";
  608. var dispatch = await _kpiCalcDispatcher.DispatchAsync(
  609. "S9_L1_001", "S9", option.TargetTenantId, option.TargetFactoryId,
  610. option.MonthlyPeriodEnd, option.MonthlyPeriodStart, option.MonthlyPeriodEnd, option.SourceZtid,
  611. batchId, triggerType, ppm, denominatorStatus, ct);
  612. sub.KpiRows = dispatch.ShouldUpsert
  613. ? await UpsertKpiValueAsync("S9_L1_001", option.MonthlyPeriodEnd, dispatch.MetricValue, now, option, "S9")
  614. : 0;
  615. sub.DenominatorStatus = dispatch.DenominatorStatus;
  616. return sub;
  617. }
  618. private async Task<int> UpsertKpiValueAsync(
  619. string metricCode,
  620. DateTime bizDate,
  621. decimal? metricValue,
  622. DateTime now,
  623. S7MdpRefreshOption option,
  624. string moduleCode = ModuleCode,
  625. string valueTable = "ado_s9_kpi_value_l1_day")
  626. {
  627. // 沿用 S3 UpsertS3KpiValueAsync 范式:先查现存行 → UPDATE;不存在 → SELECT MAX(id)+1 显式生成 id 后 INSERT。
  628. // ado_s9_kpi_value_l1_day.id 为手工分配主键(无 AUTO_INCREMENT),必须显式 set;
  629. // metric_value 允许 NULL(分母缺失不得伪装真实 0)。
  630. // FIX-2:截断时分秒(月度 KPI 入参可能为 YYYY-MM-DD 23:59:59),保证 SELECT WHERE biz_date=@BizDate 与 DB date 列匹配,避免重复 INSERT。
  631. // FIX-1:tenant_id/factory_id 取自 option,默认仍为 1300000000001/1,不破坏 Demo。
  632. bizDate = bizDate.Date;
  633. var decision = await _kpiStatus.ResolveAsync(option.TargetTenantId, metricCode, metricValue);
  634. metricValue = decision.Value;
  635. var snap = await _kpiTargetResolver.ResolveAsync(option.TargetTenantId, option.TargetFactoryId, metricCode, moduleCode, bizDate);
  636. var existingId = await _db.Ado.GetLongAsync(
  637. $"SELECT IFNULL((SELECT id FROM {valueTable} WHERE tenant_id=@TenantId AND factory_id=@FactoryId " +
  638. "AND module_code=@ModuleCode AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0 " +
  639. "ORDER BY id LIMIT 1), 0)",
  640. new List<SugarParameter>
  641. {
  642. new("@TenantId", option.TargetTenantId),
  643. new("@FactoryId", option.TargetFactoryId),
  644. new("@ModuleCode", moduleCode),
  645. new("@MetricCode", metricCode),
  646. new("@BizDate", bizDate)
  647. });
  648. if (existingId > 0)
  649. {
  650. return await _db.Ado.ExecuteCommandAsync(
  651. $"UPDATE {valueTable} SET {SmartOps.KpiResultWrite.SetClause}metric_value=@MetricValue, target_value=@TargetValue, " +
  652. "target_config_id=@TargetConfigId, target_source=@TargetSource, target_resolved_at=@TargetResolvedAt, " +
  653. "calc_time=@Now, update_time=@Now, is_deleted=0, is_active=1 WHERE id=@Id",
  654. new SugarParameter("@MetricValue", metricValue),
  655. SmartOps.KpiResultWrite.Status(decision.Status),
  656. SmartOps.KpiResultWrite.Reason(decision.Reason),
  657. new SugarParameter("@TargetValue", SmartOps.KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  658. new SugarParameter("@TargetConfigId", SmartOps.KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  659. new SugarParameter("@TargetSource", SmartOps.KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  660. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt),
  661. new SugarParameter("@Now", now),
  662. new SugarParameter("@Id", existingId));
  663. }
  664. var nextId = Yitter.IdGenerator.YitIdHelper.NextId();
  665. return await _db.Ado.ExecuteCommandAsync($@"
  666. INSERT INTO {valueTable}
  667. (id, tenant_id, org_id, company_id, factory_id, status, biz_date,
  668. create_time, update_time, is_deleted, is_active,
  669. module_code, metric_code, result_status, result_reason, metric_value, target_value, calc_time,
  670. target_config_id, target_source, target_resolved_at)
  671. VALUES
  672. (@Id, @TenantId, NULL, NULL, @FactoryId, NULL, @BizDate,
  673. @Now, @Now, 0, 1,
  674. @ModuleCode, @MetricCode, @ResultStatus, @ResultReason, @MetricValue, @TargetValue, @Now,
  675. @TargetConfigId, @TargetSource, @TargetResolvedAt)",
  676. new SugarParameter("@Id", nextId),
  677. new SugarParameter("@TenantId", option.TargetTenantId),
  678. new SugarParameter("@FactoryId", option.TargetFactoryId),
  679. new SugarParameter("@BizDate", bizDate),
  680. new SugarParameter("@Now", now),
  681. new SugarParameter("@ModuleCode", moduleCode),
  682. new SugarParameter("@MetricCode", metricCode),
  683. SmartOps.KpiResultWrite.Status(decision.Status),
  684. SmartOps.KpiResultWrite.Reason(decision.Reason),
  685. new SugarParameter("@MetricValue", metricValue),
  686. new SugarParameter("@TargetValue", SmartOps.KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  687. new SugarParameter("@TargetConfigId", SmartOps.KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  688. new SugarParameter("@TargetSource", SmartOps.KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  689. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt));
  690. }
  691. private async Task<long> InsertTransformRunLogAsync(string batchId, DateTime startedAt, string triggerType, S7MdpRefreshOption option)
  692. {
  693. await _db.Ado.ExecuteCommandAsync(@"
  694. INSERT INTO mdp_transform_run_log
  695. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time, stage_rows, standard_rows, dwd_rows, create_time, update_time)
  696. VALUES
  697. (@TenantId, @JobCode, @JobName, @TriggerType, @BatchId, 'RUNNING', @StartTime, 0, 0, 0, @StartTime, @StartTime)",
  698. new SugarParameter("@TenantId", option.TargetTenantId),
  699. new SugarParameter("@JobCode", JobCode),
  700. new SugarParameter("@JobName", JobName),
  701. new SugarParameter("@TriggerType", triggerType),
  702. new SugarParameter("@BatchId", batchId),
  703. new SugarParameter("@StartTime", startedAt));
  704. return await _db.Ado.GetLongAsync(
  705. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  706. new List<SugarParameter> { new("@BatchId", batchId) });
  707. }
  708. private async Task MarkTransformRunSuccessAsync(long runLogId, DateTime startedAt, S7MdpSyncTransformResult result)
  709. {
  710. var finishedAt = DateTime.Now;
  711. await _db.Ado.ExecuteCommandAsync(@"
  712. UPDATE mdp_transform_run_log
  713. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  714. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=@DwdRows,
  715. summary_json=@SummaryJson, update_time=CURRENT_TIMESTAMP
  716. WHERE id=@Id",
  717. new SugarParameter("@EndTime", finishedAt),
  718. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  719. new SugarParameter("@StageRows", result.StageRows),
  720. new SugarParameter("@StandardRows", result.StandardRows),
  721. new SugarParameter("@DwdRows", result.DwdRows),
  722. new SugarParameter("@SummaryJson", JsonSerializer.Serialize(new
  723. {
  724. batchId = result.BatchId,
  725. sourceZtid = result.SourceZtid,
  726. bizDate = result.BizDate.ToString("yyyy-MM-dd"),
  727. bizMonth = result.BizMonth,
  728. dwdRows = result.DwdRows,
  729. kpiRows = result.KpiRows,
  730. perKpiDwdRows = result.PerKpiDwdRows,
  731. perKpiKpiRows = result.PerKpiKpiRows,
  732. denominatorStatus = result.KpiDenominatorStatus
  733. })),
  734. new SugarParameter("@Id", runLogId));
  735. }
  736. private async Task MarkTransformRunFailedAsync(long runLogId, DateTime startedAt, string message, string batchId)
  737. {
  738. bool runLogUpdated = false;
  739. try
  740. {
  741. var finishedAt = DateTime.Now;
  742. await _db.Ado.ExecuteCommandAsync(@"
  743. UPDATE mdp_transform_run_log
  744. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  745. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  746. WHERE id=@Id",
  747. new SugarParameter("@EndTime", finishedAt),
  748. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  749. new SugarParameter("@ErrorMessage", Truncate(message, 2000)),
  750. new SugarParameter("@Id", runLogId));
  751. runLogUpdated = true;
  752. }
  753. catch (Exception ex)
  754. {
  755. Console.Error.WriteLine($"[S7MdpSyncTransform] MarkTransformRunFailed write failed (runLogId={runLogId}): {ex.Message}");
  756. }
  757. // FAILURE-NOTIFICATION-1:写库 FAILED 成功后发通知给超级管理员;通知失败不影响主流程
  758. if (!runLogUpdated) return;
  759. try
  760. {
  761. await _sysNoticeService.AddNotice(new AddNoticeInput
  762. {
  763. Title = "S7 成品仓储 T8 KPI 跑批失败",
  764. Content = $"模块:S7 成品仓储\n批次ID:{batchId}\n失败时间:{DateTime.Now:yyyy-MM-dd HH:mm:ss}\n错误信息:{Truncate(message, 1000)}\n\n请查看 mdp_transform_run_log 获取完整错误与重试记录。",
  765. Type = NoticeTypeEnum.NOTICE,
  766. PublicTime = DateTime.Now,
  767. Status = NoticeStatusEnum.PUBLIC,
  768. PublicUserId = NoticeReceiverUserId,
  769. PublicUserName = NoticeReceiverUserName
  770. });
  771. }
  772. catch (Exception notifyEx)
  773. {
  774. _logger.LogError(notifyEx, "[S7MdpSyncTransform] SysNotice 发送失败 (runLogId={RunLogId}, batchId={BatchId})", runLogId, batchId);
  775. }
  776. }
  777. private static string NormalizeTriggerType(string s) =>
  778. string.IsNullOrWhiteSpace(s) ? "AUTO" : s.Trim().ToUpperInvariant();
  779. private static void NormalizeOption(S7MdpRefreshOption option)
  780. {
  781. var d = S7MdpRefreshOption.Default();
  782. if (option.TargetFactoryId <= 0) option.TargetFactoryId = d.TargetFactoryId;
  783. if (string.IsNullOrWhiteSpace(option.SourceZtid)) option.SourceZtid = d.SourceZtid;
  784. option.TargetTenantId = AidopSourceTenantMap.ResolveTenantId(option.SourceZtid, option.TargetTenantId);
  785. if (option.BizDate == default) option.BizDate = d.BizDate;
  786. if (string.IsNullOrWhiteSpace(option.BizMonth)) option.BizMonth = d.BizMonth;
  787. if (option.MonthlyPeriodStart == default) option.MonthlyPeriodStart = d.MonthlyPeriodStart;
  788. if (option.MonthlyPeriodEnd == default) option.MonthlyPeriodEnd = d.MonthlyPeriodEnd;
  789. }
  790. private static string Truncate(string s, int max) =>
  791. string.IsNullOrEmpty(s) ? "" : (s.Length <= max ? s : s.Substring(0, max));
  792. }
  793. // DTO ────────────────────────────────────────────────────────────────────────
  794. public sealed class S7MdpRefreshOption
  795. {
  796. public string SourceZtid { get; set; } = "pbxfxp";
  797. /// <summary>KPI/DWD 落库目标租户;≤0 时由 <see cref="AidopSourceTenantMap"/> 按 SourceZtid 解析。</summary>
  798. public long TargetTenantId { get; set; }
  799. /// <summary>KPI/DWD 落库目标工厂;默认 1。</summary>
  800. public long TargetFactoryId { get; set; } = 1L;
  801. public DateTime BizDate { get; set; }
  802. public string BizMonth { get; set; } = "";
  803. public DateTime MonthlyPeriodStart { get; set; }
  804. public DateTime MonthlyPeriodEnd { get; set; }
  805. public static S7MdpRefreshOption Default()
  806. {
  807. var today = DateTime.Today;
  808. var yesterday = today.AddDays(-1);
  809. var lastMonth = today.AddMonths(-1);
  810. var monthStart = new DateTime(lastMonth.Year, lastMonth.Month, 1);
  811. var monthEnd = monthStart.AddMonths(1).AddDays(-1);
  812. return new S7MdpRefreshOption
  813. {
  814. SourceZtid = "pbxfxp",
  815. TargetTenantId = 0,
  816. TargetFactoryId = 1L,
  817. BizDate = yesterday,
  818. BizMonth = lastMonth.ToString("yyyy-MM"),
  819. MonthlyPeriodStart = monthStart,
  820. MonthlyPeriodEnd = monthEnd.AddDays(1).AddSeconds(-1)
  821. };
  822. }
  823. }
  824. public sealed class S7MdpSyncTransformResult
  825. {
  826. public string BatchId { get; set; } = "";
  827. public long RunLogId { get; set; }
  828. public string TriggerType { get; set; } = "AUTO";
  829. public string SourceZtid { get; set; } = "";
  830. public long TargetTenantId { get; set; }
  831. public long TargetFactoryId { get; set; }
  832. public DateTime BizDate { get; set; }
  833. public string BizMonth { get; set; } = "";
  834. public DateTime MonthlyPeriodStart { get; set; }
  835. public DateTime MonthlyPeriodEnd { get; set; }
  836. public int StageRows { get; set; }
  837. public int StandardRows { get; set; }
  838. public int DwdRows { get; set; }
  839. public int KpiRows { get; set; }
  840. public Dictionary<string, int> PerKpiDwdRows { get; } = new();
  841. public Dictionary<string, int> PerKpiKpiRows { get; } = new();
  842. public List<string> KpiDenominatorStatus { get; } = new();
  843. public void MergeSub(string kpiCode, KpiBuildSubResult sub)
  844. {
  845. PerKpiDwdRows[kpiCode] = sub.DwdRows;
  846. PerKpiKpiRows[kpiCode] = sub.KpiRows;
  847. DwdRows += sub.DwdRows;
  848. KpiRows += sub.KpiRows;
  849. KpiDenominatorStatus.Add($"{kpiCode}:{sub.DenominatorStatus}");
  850. }
  851. }
  852. public sealed class KpiBuildSubResult
  853. {
  854. public int T8Rows { get; set; }
  855. public int DwdRows { get; set; }
  856. public int KpiRows { get; set; }
  857. public string DenominatorStatus { get; set; } = "OK";
  858. }
  859. // T8 result set 投影类型 ──────────────────────────────────────────────────────
  860. internal sealed class S7CycleRow
  861. {
  862. public string? order_no { get; set; }
  863. public int? cycle_days { get; set; }
  864. }
  865. internal sealed class S7FulfillmentRow
  866. {
  867. public string? order_no { get; set; }
  868. public int total_rows { get; set; }
  869. public int in_window_rows { get; set; }
  870. }
  871. internal sealed class S7ShipmentDetailRow
  872. {
  873. public string? task_no { get; set; }
  874. public string? item_code { get; set; }
  875. public string? approved_date { get; set; }
  876. public decimal? qty_change { get; set; }
  877. }
  878. internal sealed class S7PeNumRow
  879. {
  880. public int penum { get; set; }
  881. }
  882. internal sealed class S7QualityReturnSummaryRow
  883. {
  884. public decimal? return_qty { get; set; }
  885. public decimal? shipped_qty { get; set; }
  886. }
  887. internal sealed class S7StageKpiRow
  888. {
  889. public decimal? MetricValue { get; set; }
  890. public int RowCount { get; set; }
  891. }