S7MdpSyncTransformService.cs 45 KB

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