S7MdpSyncTransformService.cs 44 KB

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