ActionEventMdpSyncService.cs 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  2. using Microsoft.Extensions.Logging;
  3. using SqlSugar;
  4. namespace Admin.NET.Plugin.AiDOP.Order;
  5. /// <summary>
  6. /// S8 Stage-1「订单评审」ActualHours 数据链:
  7. /// <c>aidop_action_run_log → mdp_stg_action_run → mdp_std_action_event → dwd_order_review_kpi</c>。
  8. ///
  9. /// <para><b>为什么必须经中台</b>:S8 正式运行时禁止直读 <c>aidop_action_run_log</c> /
  10. /// <c>crm_seorder</c> / <c>crm_seorderentry</c>。配对、FAILED 过滤、重复动作处理全部在本服务完成,
  11. /// S8 只消费 <c>dwd_order_review_kpi</c>。</para>
  12. ///
  13. /// <para><b>租户是唯一隔离边界</b>:每一层的读、写、JOIN、唯一键首列都带 <c>tenant_id</c>;
  14. /// 源侧 <c>tenant_id</c> 为 NULL 或 0 一律拒绝入站,绝不补默认租户。</para>
  15. /// </summary>
  16. public sealed class ActionEventMdpSyncService : ITransient
  17. {
  18. /// <summary>MDP 实体码(注册见 UpdateScripts/1.0.527.sql)。</summary>
  19. public const string EntityCode = "S1_ACTION_RUN_LOG";
  20. private const string SourceTable = "aidop_action_run_log";
  21. private const string ReviewActionCode = "S1_ORDER_REVIEW";
  22. private const string ConfirmActionCode = "S1_DELIVERY_CONFIRM";
  23. private const string OrderBizType = "crm_seorder";
  24. private readonly ISqlSugarClient _db;
  25. private readonly MdpModuleStagingPuller _stagingPuller;
  26. private readonly ILogger<ActionEventMdpSyncService> _logger;
  27. public ActionEventMdpSyncService(
  28. ISqlSugarClient db,
  29. MdpModuleStagingPuller stagingPuller,
  30. ILogger<ActionEventMdpSyncService> logger)
  31. {
  32. _db = db;
  33. _stagingPuller = stagingPuller;
  34. _logger = logger;
  35. }
  36. /// <summary>
  37. /// 全链路跑批:源→STG(增量 + RUNNING 回捞)→ STD → DWD KPI。
  38. /// 幂等:同一份源数据重复执行,STD 与 DWD 结果完全一致。
  39. /// </summary>
  40. public async Task<ActionEventSyncResult> RunAsync(long tenantId, CancellationToken cancellationToken = default)
  41. {
  42. if (tenantId <= 0)
  43. throw new InvalidOperationException("ActionEventMdpSyncService 必须显式传入有效 tenantId,禁止默认租户。");
  44. var now = DateTime.Now;
  45. var batchId = $"S1_ACTION_EVENT_{tenantId}_{now:yyyyMMddHHmmss}";
  46. var result = new ActionEventSyncResult { BatchId = batchId, TenantId = tenantId };
  47. result.StageRows = await PullStagingAsync(tenantId, batchId, cancellationToken);
  48. result.StageRepulledRunning = await RepullRunningAsync(tenantId, batchId, now, cancellationToken);
  49. result.StandardRows = await TransformStdFromStgAsync(tenantId, batchId, cancellationToken);
  50. result.KpiRows = await BuildOrderReviewKpiAsync(tenantId, batchId, now, cancellationToken);
  51. // Stage-2 必须在 Stage-1 KPI 之后跑:它直接消费 dwd_order_review_kpi.review_end_time
  52. // 作为 Start,而不是自己再实现一遍 Stage-1 的配对逻辑。
  53. result.ProductDesignKpiRows = await BuildProductDesignKpiAsync(tenantId, batchId, now, cancellationToken);
  54. _logger.LogInformation(
  55. "S1 动作事件链路完成 tenant={TenantId} batch={BatchId} stg={Stg} repull={Repull} std={Std} kpi={Kpi} pdKpi={PdKpi}",
  56. tenantId, batchId, result.StageRows, result.StageRepulledRunning,
  57. result.StandardRows, result.KpiRows, result.ProductDesignKpiRows);
  58. return result;
  59. }
  60. // ══════════════════════════════════════════════════════════════════════
  61. // PHASE 2 · 源 → STG
  62. // ══════════════════════════════════════════════════════════════════════
  63. /// <summary>
  64. /// 主增量:走统一执行器(<c>sync_mode=INCR</c>,主游标 <c>start_time</c>,水位存
  65. /// <c>mdp_entity.last_cursor</c>)。写入经 <see cref="MdpStagingWriter"/> 的 canonical 契约,
  66. /// 其中已内置「tenant 解析不出即跳过 + 跨租户行丢弃」双重守卫。
  67. /// </summary>
  68. private Task<int> PullStagingAsync(long tenantId, string batchId, CancellationToken cancellationToken) =>
  69. _stagingPuller.PullEntitiesAsync(
  70. [EntityCode], batchId, tenantId,
  71. fullRefresh: false,
  72. taskCode: "S1_ACTION_EVENT_INBOUND",
  73. cancellationToken,
  74. factoryId: 0,
  75. requireMatchingSourceTenant: true);
  76. /// <summary>
  77. /// <b>RUNNING 回捞补偿(本批的核心风险控制点)。</b>
  78. ///
  79. /// <para>源行不是内容级 append-only:<c>StartAsync</c> 先 INSERT 一行
  80. /// <c>status='RUNNING', end_time=NULL</c>,之后 <c>FinishAsync</c> /
  81. /// <c>FailStaleRunningAsync</c> 会<b>原地 UPDATE</b> 成 SUCCESS/FAILED 并补 <c>end_time</c>。
  82. /// 而源表<b>没有 update_time 列</b>,且 <c>create_time ≡ start_time</c>、不随 UPDATE 推进
  83. /// —— 所以主游标 <c>start_time &gt; @cursor</c> 永远不会第二次捞到这一行,
  84. /// 若无补偿,STG/STD 会永久停留在 RUNNING。</para>
  85. ///
  86. /// <para>补偿判据:以<b>标准层当前仍为 RUNNING 的 action_run_id 集合</b>反向驱动重拉。
  87. /// 该集合天然极小(实测全库仅 1 行),因此:
  88. /// ① tenant-safe —— 读写两侧都带 <c>tenant_id=@TenantId</c>;
  89. /// ② 幂等 —— <c>ON DUPLICATE KEY UPDATE</c>,重复执行不产生新行;
  90. /// ③ 不全表扫描 —— 由 RUNNING 集合驱动,走 <c>uk_std_action_event</c> 与源主键;
  91. /// ④ 不依赖 id 单调性 —— 按 id 精确等值匹配,与雪花乱序无关;
  92. /// ⑤ 不遗漏迟到终态 —— 只要某行在 STD 还是 RUNNING,每轮都会被重新检查。</para>
  93. ///
  94. /// <para>此处不经 <see cref="MdpStagingWriter"/> 而用定向 <c>INSERT…SELECT</c>,
  95. /// 原因是框架执行器只支持「按水位翻页」、不支持「按 id 集合定向重拉」;
  96. /// 因此租户守卫在本 SQL 内显式重建(<c>tenant_id=@TenantId</c> +
  97. /// <c>tenant_id IS NOT NULL</c> + <c>tenant_id&lt;&gt;0</c>),语义与框架守卫一致。</para>
  98. /// </summary>
  99. private async Task<int> RepullRunningAsync(
  100. long tenantId, string batchId, DateTime now, CancellationToken cancellationToken)
  101. {
  102. cancellationToken.ThrowIfCancellationRequested();
  103. var ps = new SugarParameter[]
  104. {
  105. new("@TenantId", tenantId),
  106. new("@BatchId", batchId),
  107. new("@Now", now)
  108. };
  109. // 逻辑计数:本轮被回捞的 distinct action_run_id 数。
  110. // 必须先于 upsert 求值(upsert 不改变该集合,但顺序固定更易解释),
  111. // 且与下方 INSERT 共用同一份 WHERE(RepullRunningFrom),杜绝谓词漂移。
  112. var logicalRows = await _db.Ado.GetIntAsync(
  113. $"SELECT COUNT(DISTINCT s.`id`) {RepullRunningFrom}", ps);
  114. const string sql = $"""
  115. INSERT INTO `mdp_stg_action_run`
  116. (`tenant_id`, `source_system`, `source_table`, `source_row_id`, `source_biz_key`,
  117. `raw_data`, `sync_batch_id`, `sync_time`, `process_status`, `create_time`, `update_time`)
  118. SELECT
  119. s.`tenant_id`,
  120. 'AIDOPDEV_MYSQL',
  121. 'aidop_action_run_log',
  122. CAST(s.`id` AS CHAR),
  123. CAST(s.`id` AS CHAR),
  124. JSON_OBJECT(
  125. 'id', CAST(s.`id` AS CHAR),
  126. 'tenant_id', s.`tenant_id`,
  127. 'action_code', s.`action_code`,
  128. 'biz_type', s.`biz_type`,
  129. 'biz_id', s.`biz_id`,
  130. 'biz_no', s.`biz_no`,
  131. 'status', s.`status`,
  132. 'message', s.`message`,
  133. 'detail_json', s.`detail_json`,
  134. 'start_time', DATE_FORMAT(s.`start_time`, '%Y-%m-%d %H:%i:%s'),
  135. 'end_time', DATE_FORMAT(s.`end_time`, '%Y-%m-%d %H:%i:%s'),
  136. 'create_time', DATE_FORMAT(s.`create_time`,'%Y-%m-%d %H:%i:%s')
  137. ),
  138. @BatchId, @Now, 'PENDING', @Now, @Now
  139. {RepullRunningFrom}
  140. ON DUPLICATE KEY UPDATE
  141. `source_row_id` = VALUES(`source_row_id`),
  142. `raw_data` = VALUES(`raw_data`),
  143. `sync_batch_id` = VALUES(`sync_batch_id`),
  144. `sync_time` = VALUES(`sync_time`),
  145. `process_status` = 'PENDING',
  146. `update_time` = VALUES(`update_time`)
  147. """;
  148. // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
  149. await _db.Ado.ExecuteCommandAsync(sql, ps);
  150. if (logicalRows > 0)
  151. _logger.LogInformation("RUNNING 回捞补偿命中 tenant={TenantId} rows={Rows}", tenantId, logicalRows);
  152. return logicalRows;
  153. }
  154. /// <summary>
  155. /// RUNNING 回捞的取数范围(FROM + WHERE)。被计数查询与 upsert 共用,保证两者永远同一集合。
  156. /// </summary>
  157. private const string RepullRunningFrom = """
  158. FROM `aidop_action_run_log` s
  159. WHERE s.`tenant_id` = @TenantId
  160. AND s.`tenant_id` IS NOT NULL
  161. AND s.`tenant_id` <> 0
  162. AND EXISTS (
  163. SELECT 1 FROM `mdp_std_action_event` e
  164. WHERE e.`tenant_id` = @TenantId
  165. AND e.`action_status` = 'RUNNING'
  166. AND e.`action_run_id` = s.`id`
  167. )
  168. """;
  169. // ══════════════════════════════════════════════════════════════════════
  170. // PHASE 3 · STG → STD(审计事实层,全状态保留)
  171. // ══════════════════════════════════════════════════════════════════════
  172. /// <summary>
  173. /// 贴源 → 标准事件层。Grain = <c>tenant_id + action_run_id</c>。
  174. /// <para>RUNNING / SUCCESS / FAILED <b>全部保留</b> —— STD 是审计事实层,
  175. /// 不因为 KPI 只需要 SUCCESS 就在此丢弃其它状态。</para>
  176. /// <para>同一 <c>action_run_id</c> 重拉时,<c>ON DUPLICATE KEY UPDATE</c> 允许
  177. /// RUNNING 被 SUCCESS/FAILED 更新(终态推进)。</para>
  178. /// </summary>
  179. public async Task<int> TransformStdFromStgAsync(
  180. long tenantId, string batchId, CancellationToken cancellationToken = default)
  181. {
  182. if (tenantId <= 0) throw new InvalidOperationException("TransformStdFromStgAsync 需要有效 tenantId。");
  183. cancellationToken.ThrowIfCancellationRequested();
  184. var ps = new SugarParameter[]
  185. {
  186. new("@TenantId", tenantId),
  187. new("@BatchId", batchId)
  188. };
  189. // 逻辑计数:本轮参与 STD transform 的 distinct action_run_id 数。
  190. // 与下方 INSERT 共用同一份 FROM/WHERE(StdEligibleFrom),保证「计了什么」= 「写了什么」。
  191. var logicalRows = await _db.Ado.GetIntAsync(
  192. $"SELECT COUNT(DISTINCT JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id'))) {StdEligibleFrom}", ps);
  193. const string sql = $"""
  194. INSERT INTO `mdp_std_action_event`
  195. (`tenant_id`, `action_run_id`, `action_code`, `biz_type`, `biz_id`, `biz_no`,
  196. `action_status`, `action_start_time`, `action_end_time`, `message`, `detail_json`,
  197. `source_system`, `source_table`, `source_biz_key`, `sync_batch_id`)
  198. SELECT
  199. g.`tenant_id`,
  200. CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id')) AS SIGNED),
  201. JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.action_code')),
  202. JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_type')),
  203. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_id')) REGEXP '^-?[0-9]+$'
  204. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_id')) AS SIGNED) END,
  205. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_no')), 'null'),
  206. JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.status')),
  207. CAST(REPLACE(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.start_time')), 'T', ' ') AS DATETIME),
  208. CASE WHEN NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.end_time')), 'null') IS NULL
  209. THEN NULL
  210. ELSE CAST(REPLACE(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.end_time')), 'T', ' ') AS DATETIME) END,
  211. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.message')), 'null'),
  212. CASE
  213. WHEN JSON_EXTRACT(g.`raw_data`, '$.detail_json') IS NULL THEN NULL
  214. WHEN JSON_VALID(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.detail_json')))
  215. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.detail_json')) AS JSON)
  216. ELSE JSON_EXTRACT(g.`raw_data`, '$.detail_json')
  217. END,
  218. g.`source_system`,
  219. g.`source_table`,
  220. g.`source_biz_key`,
  221. @BatchId
  222. {StdEligibleFrom}
  223. ON DUPLICATE KEY UPDATE
  224. `action_code` = VALUES(`action_code`),
  225. `biz_type` = VALUES(`biz_type`),
  226. `biz_id` = VALUES(`biz_id`),
  227. `biz_no` = VALUES(`biz_no`),
  228. `action_status` = VALUES(`action_status`),
  229. `action_start_time` = VALUES(`action_start_time`),
  230. `action_end_time` = VALUES(`action_end_time`),
  231. `message` = VALUES(`message`),
  232. `detail_json` = VALUES(`detail_json`),
  233. `sync_batch_id` = VALUES(`sync_batch_id`)
  234. """;
  235. // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
  236. await _db.Ado.ExecuteCommandAsync(sql, ps);
  237. return logicalRows;
  238. }
  239. /// <summary>
  240. /// STD transform 的取数范围(FROM + WHERE)。被计数查询与 upsert 共用,保证两者永远同一集合。
  241. /// </summary>
  242. private const string StdEligibleFrom = """
  243. FROM `mdp_stg_action_run` g
  244. WHERE g.`tenant_id` = @TenantId
  245. AND g.`tenant_id` <> 0
  246. AND g.`source_table` = 'aidop_action_run_log'
  247. AND JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id')) REGEXP '^[0-9]+$'
  248. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.action_code')), 'null') IS NOT NULL
  249. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_type')), 'null') IS NOT NULL
  250. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.status')), 'null') IS NOT NULL
  251. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.start_time')), 'null') IS NOT NULL
  252. """;
  253. // ══════════════════════════════════════════════════════════════════════
  254. // PHASE 4 · STD → DWD KPI(FIRST VALID PAIR)
  255. // ══════════════════════════════════════════════════════════════════════
  256. /// <summary>
  257. /// 标准事件层 → 订单评审 KPI。Grain = <c>tenant_id + sales_order_id</c>。
  258. ///
  259. /// <para><b>算法(定版,不得改用 LAST PAIR / MAX)</b>:
  260. /// ① 只取 <c>status='SUCCESS'</c>,FAILED 保留在 STD 但不参与配对;
  261. /// ② ReviewStart = 评审 SUCCESS 的 <b>FIRST</b>(<c>ORDER BY action_start_time, action_run_id</c>);
  262. /// ③ ReviewEnd = <b>首个 <c>action_end_time &gt; ReviewStart</c></b> 的确认 SUCCESS 的
  263. /// <c>action_end_time</c>(= 确认动作<b>完成</b>时刻,是 <c>progress='3'</c> 提交后的最近上界);
  264. /// ④ ActualHours = 自然 elapsed 小时,2 位小数,<b>不扣夜间/周末/节假日</b>(本项目无工作日历 Authority)。</para>
  265. ///
  266. /// <para><b>冻结语义</b>:算法只依赖「第一条评审」与「首个晚于它的确认」,两者一旦产生即恒定,
  267. /// 后续任意次重复点击都不会改写已形成的 KPI;因此重复执行结果完全一致(deterministic / idempotent)。</para>
  268. ///
  269. /// <para>本批<b>不计算</b> TargetHours / kpi_status —— 那依赖 S0 PI 匹配,属另一批次。</para>
  270. /// </summary>
  271. public async Task<int> BuildOrderReviewKpiAsync(
  272. long tenantId, string batchId, DateTime now, CancellationToken cancellationToken = default)
  273. {
  274. if (tenantId <= 0) throw new InvalidOperationException("BuildOrderReviewKpiAsync 需要有效 tenantId。");
  275. cancellationToken.ThrowIfCancellationRequested();
  276. var ps = new SugarParameter[]
  277. {
  278. new("@TenantId", tenantId),
  279. new("@BatchId", batchId),
  280. new("@Now", now)
  281. };
  282. // 逻辑计数:本轮参与 KPI transform 的 distinct (tenant_id, sales_order_id) 数,
  283. // 即下方 CTE `agg` 的基数 —— 与写入集合一一对应(agg 是最终 SELECT 的驱动表)。
  284. var logicalRows = await _db.Ado.GetIntAsync($"""
  285. SELECT COUNT(*) FROM (
  286. SELECT `tenant_id`, `biz_id` {KpiEventFrom} GROUP BY `tenant_id`, `biz_id`
  287. ) agg
  288. """, ps);
  289. const string sql = $"""
  290. INSERT INTO `dwd_order_review_kpi`
  291. (`tenant_id`, `sales_order_id`, `sales_order_no`,
  292. `review_action_run_id`, `confirm_action_run_id`,
  293. `review_start_time`, `review_end_time`, `actual_hours`, `calculation_status`,
  294. `review_success_count`, `confirm_success_count`,
  295. `review_failed_count`, `confirm_failed_count`,
  296. `calc_batch_id`, `calc_time`)
  297. WITH ev AS (
  298. SELECT `tenant_id`, `biz_id`, `biz_no`, `action_run_id`, `action_code`,
  299. `action_status`, `action_start_time`, `action_end_time`
  300. {KpiEventFrom}
  301. ),
  302. agg AS (
  303. SELECT `tenant_id`, `biz_id`,
  304. MAX(`biz_no`) AS biz_no,
  305. SUM(`action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'SUCCESS') AS rs,
  306. SUM(`action_code` = 'S1_DELIVERY_CONFIRM' AND `action_status` = 'SUCCESS') AS cs,
  307. SUM(`action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'FAILED') AS rf,
  308. SUM(`action_code` = 'S1_DELIVERY_CONFIRM' AND `action_status` = 'FAILED') AS cf
  309. FROM ev GROUP BY `tenant_id`, `biz_id`
  310. ),
  311. rv AS (
  312. SELECT `tenant_id`, `biz_id`, `action_run_id`, `action_start_time`,
  313. ROW_NUMBER() OVER (PARTITION BY `tenant_id`, `biz_id`
  314. ORDER BY `action_start_time` ASC, `action_run_id` ASC) AS rn
  315. FROM ev
  316. WHERE `action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'SUCCESS'
  317. ),
  318. r1 AS (SELECT * FROM rv WHERE rn = 1),
  319. cv AS (
  320. SELECT c.`tenant_id`, c.`biz_id`, c.`action_run_id`, c.`action_end_time`,
  321. ROW_NUMBER() OVER (PARTITION BY c.`tenant_id`, c.`biz_id`
  322. ORDER BY c.`action_end_time` ASC, c.`action_run_id` ASC) AS rn
  323. FROM ev c
  324. JOIN r1 ON r1.`tenant_id` = c.`tenant_id` AND r1.`biz_id` = c.`biz_id`
  325. WHERE c.`action_code` = 'S1_DELIVERY_CONFIRM'
  326. AND c.`action_status` = 'SUCCESS'
  327. AND c.`action_end_time` IS NOT NULL
  328. AND c.`action_end_time` > r1.`action_start_time`
  329. ),
  330. c1 AS (SELECT * FROM cv WHERE rn = 1)
  331. SELECT
  332. a.`tenant_id`,
  333. a.`biz_id`,
  334. a.biz_no,
  335. r1.`action_run_id`,
  336. c1.`action_run_id`,
  337. r1.`action_start_time`,
  338. c1.`action_end_time`,
  339. CASE WHEN r1.`action_start_time` IS NOT NULL AND c1.`action_end_time` IS NOT NULL
  340. THEN ROUND(TIMESTAMPDIFF(SECOND, r1.`action_start_time`, c1.`action_end_time`) / 3600.0, 2)
  341. END,
  342. CASE
  343. WHEN r1.`action_start_time` IS NULL AND a.cs = 0 THEN 'NO_START'
  344. WHEN r1.`action_start_time` IS NULL THEN 'INVALID'
  345. WHEN c1.`action_end_time` IS NOT NULL THEN 'OK'
  346. WHEN a.cs = 0 THEN 'IN_PROGRESS'
  347. ELSE 'INVALID_TIME_SEQUENCE'
  348. END,
  349. a.rs, a.cs, a.rf, a.cf,
  350. @BatchId, @Now
  351. FROM agg a
  352. LEFT JOIN r1 ON r1.`tenant_id` = a.`tenant_id` AND r1.`biz_id` = a.`biz_id`
  353. LEFT JOIN c1 ON c1.`tenant_id` = a.`tenant_id` AND c1.`biz_id` = a.`biz_id`
  354. ON DUPLICATE KEY UPDATE
  355. `sales_order_no` = VALUES(`sales_order_no`),
  356. `review_action_run_id` = VALUES(`review_action_run_id`),
  357. `confirm_action_run_id` = VALUES(`confirm_action_run_id`),
  358. `review_start_time` = VALUES(`review_start_time`),
  359. `review_end_time` = VALUES(`review_end_time`),
  360. `actual_hours` = VALUES(`actual_hours`),
  361. `calculation_status` = VALUES(`calculation_status`),
  362. `review_success_count` = VALUES(`review_success_count`),
  363. `confirm_success_count` = VALUES(`confirm_success_count`),
  364. `review_failed_count` = VALUES(`review_failed_count`),
  365. `confirm_failed_count` = VALUES(`confirm_failed_count`),
  366. `calc_batch_id` = VALUES(`calc_batch_id`),
  367. `calc_time` = VALUES(`calc_time`)
  368. """;
  369. // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
  370. await _db.Ado.ExecuteCommandAsync(sql, ps);
  371. return logicalRows;
  372. }
  373. /// <summary>
  374. /// KPI transform 的事件取数范围(FROM + WHERE)。被计数查询与 CTE <c>ev</c> 共用,
  375. /// 保证「计的订单数」与「写入的订单集合」同源。
  376. /// </summary>
  377. private const string KpiEventFrom = """
  378. FROM `mdp_std_action_event`
  379. WHERE `tenant_id` = @TenantId
  380. AND `biz_type` = 'crm_seorder'
  381. AND `biz_id` IS NOT NULL
  382. AND `action_code` IN ('S1_ORDER_REVIEW', 'S1_DELIVERY_CONFIRM')
  383. """;
  384. // ══════════════════════════════════════════════════════════════════════
  385. // Stage-2 · 产品设计 KPI(dwd_product_design_kpi)
  386. // ══════════════════════════════════════════════════════════════════════
  387. /// <summary>
  388. /// Stage-2「产品设计」KPI。Grain = <c>tenant_id + sales_order_id</c>。
  389. ///
  390. /// <para><b>Start</b> 直接取 <c>dwd_order_review_kpi.review_end_time</c> —— Stage-1 已有正式
  391. /// DWD Authority,这里只消费,<b>不重新实现 Stage-1 的配对逻辑</b>。</para>
  392. ///
  393. /// <para><b>End</b> 取该订单 <c>S2_DAILY_PLAN_RELEASE</c> 的 <b>FIRST SUCCESS</b>
  394. /// (<c>ORDER BY action_end_time, action_run_id</c>)的 <c>action_end_time</c>。
  395. /// 将来同一订单再次下达也<b>不会改写</b>已成型的 Stage-2 End —— 首个里程碑一旦产生即恒定,
  396. /// 故本 transform 是 deterministic / idempotent。</para>
  397. ///
  398. /// <para><b>END_UNKNOWN 是本表存在的核心理由</b>:1.0.532 之前没有释放事件,历史上已下达的订单
  399. /// 只能证明「已下达」、无法证明「何时下达」。这类订单必须与「确实还没下达」的 IN_PROGRESS
  400. /// 区分开,否则会把历史数据缺口误判成业务未完成。区分依据是中台里是否存在该订单
  401. /// <c>status='Y'</c> 的日计划行 —— <b>只取这个布尔事实,绝不取它的时间</b>。</para>
  402. ///
  403. /// <para><b>MDP ONLY</b>:全部取数来自 <c>dwd_order_review_kpi</c> / <c>mdp_std_action_event</c> /
  404. /// <c>mdp_std_operation_schedule</c> / <c>mdp_std_work_order_schedule</c> / <c>mdp_std_so</c>,
  405. /// 不 JOIN 任何源业务表。</para>
  406. ///
  407. /// <para><b>订单全集以 <c>mdp_std_so</c> 为准</b>,而不是由释放事件驱动 ——
  408. /// 否则「尚未开始评审」的订单会从档案里整个消失。</para>
  409. /// </summary>
  410. public async Task<int> BuildProductDesignKpiAsync(
  411. long tenantId, string batchId, DateTime now, CancellationToken cancellationToken = default)
  412. {
  413. if (tenantId <= 0) throw new InvalidOperationException("BuildProductDesignKpiAsync 需要有效 tenantId。");
  414. cancellationToken.ThrowIfCancellationRequested();
  415. var ps = new SugarParameter[]
  416. {
  417. new("@TenantId", tenantId),
  418. new("@BatchId", batchId),
  419. new("@Now", now)
  420. };
  421. // 逻辑计数:本轮参与 Stage-2 KPI transform 的 distinct 销售订单数(非 DB affected rows)
  422. var logicalRows = await _db.Ado.GetIntAsync($"""
  423. SELECT COUNT(*) FROM ({ProductDesignOrderScope}) o
  424. """, ps);
  425. const string sql = $"""
  426. INSERT INTO `dwd_product_design_kpi`
  427. (`tenant_id`, `sales_order_id`, `sales_order_no`, `review_end_time`, `release_action_run_id`,
  428. `product_design_start_time`, `product_design_end_time`, `actual_hours`, `calculation_status`,
  429. `release_success_count`, `release_failed_count`, `has_released_schedule`,
  430. `calc_batch_id`, `calc_time`)
  431. WITH orders AS ({ProductDesignOrderScope}),
  432. -- 该订单的释放事件全集(用于计数;SUCCESS/FAILED 都留在事件层)
  433. rel_all AS (
  434. SELECT `tenant_id`, `biz_id`,
  435. SUM(`action_status` = 'SUCCESS') AS succ_cnt,
  436. SUM(`action_status` = 'FAILED') AS fail_cnt
  437. FROM `mdp_std_action_event`
  438. WHERE `tenant_id` = @TenantId
  439. AND `action_code` = 'S2_DAILY_PLAN_RELEASE'
  440. AND `biz_type` = 'crm_seorder'
  441. AND `biz_id` IS NOT NULL
  442. GROUP BY `tenant_id`, `biz_id`
  443. ),
  444. -- FIRST SUCCESS:Stage-2 End 永远取首个成功释放,后续释放不改写历史
  445. rel_rank AS (
  446. SELECT `tenant_id`, `biz_id`, `action_run_id`, `action_end_time`,
  447. ROW_NUMBER() OVER (PARTITION BY `tenant_id`, `biz_id`
  448. ORDER BY `action_end_time` ASC, `action_run_id` ASC) AS rn
  449. FROM `mdp_std_action_event`
  450. WHERE `tenant_id` = @TenantId
  451. AND `action_code` = 'S2_DAILY_PLAN_RELEASE'
  452. AND `biz_type` = 'crm_seorder'
  453. AND `biz_id` IS NOT NULL
  454. AND `action_status` = 'SUCCESS'
  455. AND `action_end_time` IS NOT NULL
  456. ),
  457. rel1 AS (SELECT * FROM rel_rank WHERE rn = 1),
  458. -- 历史「已下达」状态佐证:MDP-only ID 桥
  459. -- operation_schedule → work_order_schedule.sales_order_entry_id → std_so.order_entry_id → order_id
  460. -- 只产出布尔,不取任何时间(该层的 update_time 会被覆盖,不是下达时刻)
  461. rel_sched AS (
  462. SELECT DISTINCT so.`tenant_id`, so.`order_id`
  463. FROM `mdp_std_operation_schedule` op
  464. JOIN `mdp_std_work_order_schedule` wo
  465. ON wo.`tenant_id` = op.`tenant_id` AND wo.`work_order` = op.`work_order`
  466. JOIN `mdp_std_so` so
  467. ON so.`tenant_id` = wo.`tenant_id` AND so.`order_entry_id` = wo.`sales_order_entry_id`
  468. WHERE op.`tenant_id` = @TenantId
  469. AND op.`status` = 'Y'
  470. AND wo.`sales_order_entry_id` IS NOT NULL
  471. AND so.`order_id` IS NOT NULL
  472. )
  473. SELECT
  474. o.`tenant_id`,
  475. o.`order_id`,
  476. o.order_no,
  477. k.`review_end_time`,
  478. r.`action_run_id`,
  479. k.`review_end_time` AS product_design_start_time,
  480. r.`action_end_time` AS product_design_end_time,
  481. CASE WHEN k.`review_end_time` IS NOT NULL
  482. AND r.`action_end_time` IS NOT NULL
  483. AND r.`action_end_time` > k.`review_end_time`
  484. THEN ROUND(TIMESTAMPDIFF(SECOND, k.`review_end_time`, r.`action_end_time`) / 3600.0, 2)
  485. END AS actual_hours,
  486. CASE
  487. -- 无 Start 无 End → 还没走到这一步
  488. WHEN k.`review_end_time` IS NULL AND r.`action_end_time` IS NULL THEN 'NOT_STARTED'
  489. -- 有 End 却没有 Start → 数据异常,不得倒推 Start
  490. WHEN k.`review_end_time` IS NULL THEN 'INVALID_NO_START'
  491. -- End <= Start → 时序异常,绝不 abs()/swap() 硬算正数
  492. WHEN r.`action_end_time` IS NOT NULL
  493. AND r.`action_end_time` <= k.`review_end_time` THEN 'INVALID_TIME_SEQUENCE'
  494. WHEN r.`action_end_time` IS NOT NULL THEN 'OK'
  495. -- 无释放事件但中台证明已下达 → 历史数据,时间 Authority 缺失
  496. WHEN s.`order_id` IS NOT NULL THEN 'END_UNKNOWN'
  497. ELSE 'IN_PROGRESS'
  498. END AS calculation_status,
  499. IFNULL(ra.succ_cnt, 0),
  500. IFNULL(ra.fail_cnt, 0),
  501. CASE WHEN s.`order_id` IS NOT NULL THEN 1 ELSE 0 END,
  502. @BatchId, @Now
  503. FROM orders o
  504. LEFT JOIN `dwd_order_review_kpi` k
  505. ON k.`tenant_id` = o.`tenant_id` AND k.`sales_order_id` = o.`order_id`
  506. LEFT JOIN rel1 r ON r.`tenant_id` = o.`tenant_id` AND r.`biz_id` = o.`order_id`
  507. LEFT JOIN rel_all ra ON ra.`tenant_id` = o.`tenant_id` AND ra.`biz_id` = o.`order_id`
  508. LEFT JOIN rel_sched s ON s.`tenant_id` = o.`tenant_id` AND s.`order_id` = o.`order_id`
  509. ON DUPLICATE KEY UPDATE
  510. `sales_order_no` = VALUES(`sales_order_no`),
  511. `review_end_time` = VALUES(`review_end_time`),
  512. `release_action_run_id` = VALUES(`release_action_run_id`),
  513. `product_design_start_time` = VALUES(`product_design_start_time`),
  514. `product_design_end_time` = VALUES(`product_design_end_time`),
  515. `actual_hours` = VALUES(`actual_hours`),
  516. `calculation_status` = VALUES(`calculation_status`),
  517. `release_success_count` = VALUES(`release_success_count`),
  518. `release_failed_count` = VALUES(`release_failed_count`),
  519. `has_released_schedule` = VALUES(`has_released_schedule`),
  520. `calc_batch_id` = VALUES(`calc_batch_id`),
  521. `calc_time` = VALUES(`calc_time`)
  522. """;
  523. // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
  524. await _db.Ado.ExecuteCommandAsync(sql, ps);
  525. return logicalRows;
  526. }
  527. /// <summary>
  528. /// Stage-2 KPI 的订单全集:以 <c>mdp_std_so</c> 的订单头为准(一订单一行)。
  529. /// <para>被计数查询与 CTE <c>orders</c> 共用。刻意<b>不</b>由释放事件驱动 ——
  530. /// 那样「尚未开始评审」的订单会从档案里消失。</para>
  531. /// </summary>
  532. private const string ProductDesignOrderScope = """
  533. SELECT `tenant_id`, `order_id`, MAX(`order_no`) AS order_no
  534. FROM `mdp_std_so`
  535. WHERE `tenant_id` = @TenantId
  536. AND `order_id` IS NOT NULL
  537. AND `order_id` <> 0
  538. AND `deleted_flag` = 0
  539. GROUP BY `tenant_id`, `order_id`
  540. """;
  541. }
  542. /// <summary>
  543. /// 动作事件链路跑批结果。
  544. ///
  545. /// <para><b>计数语义统一为「本轮参与处理的逻辑对象数」</b>,不是数据库 affected rows。
  546. /// 原因:<c>INSERT … ON DUPLICATE KEY UPDATE</c> 在 MySQL 下「插入计 1、更新计 2」,
  547. /// 直接返回会让「什么都没变」的第二次跑批显示成处理量翻倍,属于虚报。
  548. /// 因此每层都改为对该层的取数集合求 <c>COUNT(DISTINCT 业务键)</c>:
  549. /// 一个 <c>action_run_id</c> 最多计 1、一张销售订单最多计 1,重复 upsert 不翻倍。</para>
  550. /// </summary>
  551. public sealed class ActionEventSyncResult
  552. {
  553. public string BatchId { get; set; } = "";
  554. public long TenantId { get; set; }
  555. /// <summary>本轮实际进入贴源处理链的 source event 数(增量拉取到的源行数)。</summary>
  556. public int StageRows { get; set; }
  557. /// <summary>本轮被 RUNNING 回捞补偿重拉的 distinct action_run_id 数。</summary>
  558. public int StageRepulledRunning { get; set; }
  559. /// <summary>本轮参与 STD transform 的 distinct action_run_id 数。</summary>
  560. public int StandardRows { get; set; }
  561. /// <summary>本轮参与 KPI transform 的 distinct 销售订单数。</summary>
  562. public int KpiRows { get; set; }
  563. /// <summary>本轮参与 Stage-2 产品设计 KPI transform 的 distinct 销售订单数。</summary>
  564. public int ProductDesignKpiRows { get; set; }
  565. }