ActionEventMdpSyncService.cs 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441
  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. _logger.LogInformation(
  52. "S1 动作事件链路完成 tenant={TenantId} batch={BatchId} stg={Stg} repull={Repull} std={Std} kpi={Kpi}",
  53. tenantId, batchId, result.StageRows, result.StageRepulledRunning, result.StandardRows, result.KpiRows);
  54. return result;
  55. }
  56. // ══════════════════════════════════════════════════════════════════════
  57. // PHASE 2 · 源 → STG
  58. // ══════════════════════════════════════════════════════════════════════
  59. /// <summary>
  60. /// 主增量:走统一执行器(<c>sync_mode=INCR</c>,主游标 <c>start_time</c>,水位存
  61. /// <c>mdp_entity.last_cursor</c>)。写入经 <see cref="MdpStagingWriter"/> 的 canonical 契约,
  62. /// 其中已内置「tenant 解析不出即跳过 + 跨租户行丢弃」双重守卫。
  63. /// </summary>
  64. private Task<int> PullStagingAsync(long tenantId, string batchId, CancellationToken cancellationToken) =>
  65. _stagingPuller.PullEntitiesAsync(
  66. [EntityCode], batchId, tenantId,
  67. fullRefresh: false,
  68. taskCode: "S1_ACTION_EVENT_INBOUND",
  69. cancellationToken,
  70. factoryId: 0,
  71. requireMatchingSourceTenant: true);
  72. /// <summary>
  73. /// <b>RUNNING 回捞补偿(本批的核心风险控制点)。</b>
  74. ///
  75. /// <para>源行不是内容级 append-only:<c>StartAsync</c> 先 INSERT 一行
  76. /// <c>status='RUNNING', end_time=NULL</c>,之后 <c>FinishAsync</c> /
  77. /// <c>FailStaleRunningAsync</c> 会<b>原地 UPDATE</b> 成 SUCCESS/FAILED 并补 <c>end_time</c>。
  78. /// 而源表<b>没有 update_time 列</b>,且 <c>create_time ≡ start_time</c>、不随 UPDATE 推进
  79. /// —— 所以主游标 <c>start_time &gt; @cursor</c> 永远不会第二次捞到这一行,
  80. /// 若无补偿,STG/STD 会永久停留在 RUNNING。</para>
  81. ///
  82. /// <para>补偿判据:以<b>标准层当前仍为 RUNNING 的 action_run_id 集合</b>反向驱动重拉。
  83. /// 该集合天然极小(实测全库仅 1 行),因此:
  84. /// ① tenant-safe —— 读写两侧都带 <c>tenant_id=@TenantId</c>;
  85. /// ② 幂等 —— <c>ON DUPLICATE KEY UPDATE</c>,重复执行不产生新行;
  86. /// ③ 不全表扫描 —— 由 RUNNING 集合驱动,走 <c>uk_std_action_event</c> 与源主键;
  87. /// ④ 不依赖 id 单调性 —— 按 id 精确等值匹配,与雪花乱序无关;
  88. /// ⑤ 不遗漏迟到终态 —— 只要某行在 STD 还是 RUNNING,每轮都会被重新检查。</para>
  89. ///
  90. /// <para>此处不经 <see cref="MdpStagingWriter"/> 而用定向 <c>INSERT…SELECT</c>,
  91. /// 原因是框架执行器只支持「按水位翻页」、不支持「按 id 集合定向重拉」;
  92. /// 因此租户守卫在本 SQL 内显式重建(<c>tenant_id=@TenantId</c> +
  93. /// <c>tenant_id IS NOT NULL</c> + <c>tenant_id&lt;&gt;0</c>),语义与框架守卫一致。</para>
  94. /// </summary>
  95. private async Task<int> RepullRunningAsync(
  96. long tenantId, string batchId, DateTime now, CancellationToken cancellationToken)
  97. {
  98. cancellationToken.ThrowIfCancellationRequested();
  99. var ps = new SugarParameter[]
  100. {
  101. new("@TenantId", tenantId),
  102. new("@BatchId", batchId),
  103. new("@Now", now)
  104. };
  105. // 逻辑计数:本轮被回捞的 distinct action_run_id 数。
  106. // 必须先于 upsert 求值(upsert 不改变该集合,但顺序固定更易解释),
  107. // 且与下方 INSERT 共用同一份 WHERE(RepullRunningFrom),杜绝谓词漂移。
  108. var logicalRows = await _db.Ado.GetIntAsync(
  109. $"SELECT COUNT(DISTINCT s.`id`) {RepullRunningFrom}", ps);
  110. const string sql = $"""
  111. INSERT INTO `mdp_stg_action_run`
  112. (`tenant_id`, `source_system`, `source_table`, `source_row_id`, `source_biz_key`,
  113. `raw_data`, `sync_batch_id`, `sync_time`, `process_status`, `create_time`, `update_time`)
  114. SELECT
  115. s.`tenant_id`,
  116. 'AIDOPDEV_MYSQL',
  117. 'aidop_action_run_log',
  118. CAST(s.`id` AS CHAR),
  119. CAST(s.`id` AS CHAR),
  120. JSON_OBJECT(
  121. 'id', CAST(s.`id` AS CHAR),
  122. 'tenant_id', s.`tenant_id`,
  123. 'action_code', s.`action_code`,
  124. 'biz_type', s.`biz_type`,
  125. 'biz_id', s.`biz_id`,
  126. 'biz_no', s.`biz_no`,
  127. 'status', s.`status`,
  128. 'message', s.`message`,
  129. 'detail_json', s.`detail_json`,
  130. 'start_time', DATE_FORMAT(s.`start_time`, '%Y-%m-%d %H:%i:%s'),
  131. 'end_time', DATE_FORMAT(s.`end_time`, '%Y-%m-%d %H:%i:%s'),
  132. 'create_time', DATE_FORMAT(s.`create_time`,'%Y-%m-%d %H:%i:%s')
  133. ),
  134. @BatchId, @Now, 'PENDING', @Now, @Now
  135. {RepullRunningFrom}
  136. ON DUPLICATE KEY UPDATE
  137. `source_row_id` = VALUES(`source_row_id`),
  138. `raw_data` = VALUES(`raw_data`),
  139. `sync_batch_id` = VALUES(`sync_batch_id`),
  140. `sync_time` = VALUES(`sync_time`),
  141. `process_status` = 'PENDING',
  142. `update_time` = VALUES(`update_time`)
  143. """;
  144. // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
  145. await _db.Ado.ExecuteCommandAsync(sql, ps);
  146. if (logicalRows > 0)
  147. _logger.LogInformation("RUNNING 回捞补偿命中 tenant={TenantId} rows={Rows}", tenantId, logicalRows);
  148. return logicalRows;
  149. }
  150. /// <summary>
  151. /// RUNNING 回捞的取数范围(FROM + WHERE)。被计数查询与 upsert 共用,保证两者永远同一集合。
  152. /// </summary>
  153. private const string RepullRunningFrom = """
  154. FROM `aidop_action_run_log` s
  155. WHERE s.`tenant_id` = @TenantId
  156. AND s.`tenant_id` IS NOT NULL
  157. AND s.`tenant_id` <> 0
  158. AND EXISTS (
  159. SELECT 1 FROM `mdp_std_action_event` e
  160. WHERE e.`tenant_id` = @TenantId
  161. AND e.`action_status` = 'RUNNING'
  162. AND e.`action_run_id` = s.`id`
  163. )
  164. """;
  165. // ══════════════════════════════════════════════════════════════════════
  166. // PHASE 3 · STG → STD(审计事实层,全状态保留)
  167. // ══════════════════════════════════════════════════════════════════════
  168. /// <summary>
  169. /// 贴源 → 标准事件层。Grain = <c>tenant_id + action_run_id</c>。
  170. /// <para>RUNNING / SUCCESS / FAILED <b>全部保留</b> —— STD 是审计事实层,
  171. /// 不因为 KPI 只需要 SUCCESS 就在此丢弃其它状态。</para>
  172. /// <para>同一 <c>action_run_id</c> 重拉时,<c>ON DUPLICATE KEY UPDATE</c> 允许
  173. /// RUNNING 被 SUCCESS/FAILED 更新(终态推进)。</para>
  174. /// </summary>
  175. public async Task<int> TransformStdFromStgAsync(
  176. long tenantId, string batchId, CancellationToken cancellationToken = default)
  177. {
  178. if (tenantId <= 0) throw new InvalidOperationException("TransformStdFromStgAsync 需要有效 tenantId。");
  179. cancellationToken.ThrowIfCancellationRequested();
  180. var ps = new SugarParameter[]
  181. {
  182. new("@TenantId", tenantId),
  183. new("@BatchId", batchId)
  184. };
  185. // 逻辑计数:本轮参与 STD transform 的 distinct action_run_id 数。
  186. // 与下方 INSERT 共用同一份 FROM/WHERE(StdEligibleFrom),保证「计了什么」= 「写了什么」。
  187. var logicalRows = await _db.Ado.GetIntAsync(
  188. $"SELECT COUNT(DISTINCT JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id'))) {StdEligibleFrom}", ps);
  189. const string sql = $"""
  190. INSERT INTO `mdp_std_action_event`
  191. (`tenant_id`, `action_run_id`, `action_code`, `biz_type`, `biz_id`, `biz_no`,
  192. `action_status`, `action_start_time`, `action_end_time`, `message`, `detail_json`,
  193. `source_system`, `source_table`, `source_biz_key`, `sync_batch_id`)
  194. SELECT
  195. g.`tenant_id`,
  196. CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id')) AS SIGNED),
  197. JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.action_code')),
  198. JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_type')),
  199. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_id')) REGEXP '^-?[0-9]+$'
  200. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_id')) AS SIGNED) END,
  201. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_no')), 'null'),
  202. JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.status')),
  203. CAST(REPLACE(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.start_time')), 'T', ' ') AS DATETIME),
  204. CASE WHEN NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.end_time')), 'null') IS NULL
  205. THEN NULL
  206. ELSE CAST(REPLACE(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.end_time')), 'T', ' ') AS DATETIME) END,
  207. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.message')), 'null'),
  208. CASE
  209. WHEN JSON_EXTRACT(g.`raw_data`, '$.detail_json') IS NULL THEN NULL
  210. WHEN JSON_VALID(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.detail_json')))
  211. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.detail_json')) AS JSON)
  212. ELSE JSON_EXTRACT(g.`raw_data`, '$.detail_json')
  213. END,
  214. g.`source_system`,
  215. g.`source_table`,
  216. g.`source_biz_key`,
  217. @BatchId
  218. {StdEligibleFrom}
  219. ON DUPLICATE KEY UPDATE
  220. `action_code` = VALUES(`action_code`),
  221. `biz_type` = VALUES(`biz_type`),
  222. `biz_id` = VALUES(`biz_id`),
  223. `biz_no` = VALUES(`biz_no`),
  224. `action_status` = VALUES(`action_status`),
  225. `action_start_time` = VALUES(`action_start_time`),
  226. `action_end_time` = VALUES(`action_end_time`),
  227. `message` = VALUES(`message`),
  228. `detail_json` = VALUES(`detail_json`),
  229. `sync_batch_id` = VALUES(`sync_batch_id`)
  230. """;
  231. // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
  232. await _db.Ado.ExecuteCommandAsync(sql, ps);
  233. return logicalRows;
  234. }
  235. /// <summary>
  236. /// STD transform 的取数范围(FROM + WHERE)。被计数查询与 upsert 共用,保证两者永远同一集合。
  237. /// </summary>
  238. private const string StdEligibleFrom = """
  239. FROM `mdp_stg_action_run` g
  240. WHERE g.`tenant_id` = @TenantId
  241. AND g.`tenant_id` <> 0
  242. AND g.`source_table` = 'aidop_action_run_log'
  243. AND JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id')) REGEXP '^[0-9]+$'
  244. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.action_code')), 'null') IS NOT NULL
  245. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_type')), 'null') IS NOT NULL
  246. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.status')), 'null') IS NOT NULL
  247. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.start_time')), 'null') IS NOT NULL
  248. """;
  249. // ══════════════════════════════════════════════════════════════════════
  250. // PHASE 4 · STD → DWD KPI(FIRST VALID PAIR)
  251. // ══════════════════════════════════════════════════════════════════════
  252. /// <summary>
  253. /// 标准事件层 → 订单评审 KPI。Grain = <c>tenant_id + sales_order_id</c>。
  254. ///
  255. /// <para><b>算法(定版,不得改用 LAST PAIR / MAX)</b>:
  256. /// ① 只取 <c>status='SUCCESS'</c>,FAILED 保留在 STD 但不参与配对;
  257. /// ② ReviewStart = 评审 SUCCESS 的 <b>FIRST</b>(<c>ORDER BY action_start_time, action_run_id</c>);
  258. /// ③ ReviewEnd = <b>首个 <c>action_end_time &gt; ReviewStart</c></b> 的确认 SUCCESS 的
  259. /// <c>action_end_time</c>(= 确认动作<b>完成</b>时刻,是 <c>progress='3'</c> 提交后的最近上界);
  260. /// ④ ActualHours = 自然 elapsed 小时,2 位小数,<b>不扣夜间/周末/节假日</b>(本项目无工作日历 Authority)。</para>
  261. ///
  262. /// <para><b>冻结语义</b>:算法只依赖「第一条评审」与「首个晚于它的确认」,两者一旦产生即恒定,
  263. /// 后续任意次重复点击都不会改写已形成的 KPI;因此重复执行结果完全一致(deterministic / idempotent)。</para>
  264. ///
  265. /// <para>本批<b>不计算</b> TargetHours / kpi_status —— 那依赖 S0 PI 匹配,属另一批次。</para>
  266. /// </summary>
  267. public async Task<int> BuildOrderReviewKpiAsync(
  268. long tenantId, string batchId, DateTime now, CancellationToken cancellationToken = default)
  269. {
  270. if (tenantId <= 0) throw new InvalidOperationException("BuildOrderReviewKpiAsync 需要有效 tenantId。");
  271. cancellationToken.ThrowIfCancellationRequested();
  272. var ps = new SugarParameter[]
  273. {
  274. new("@TenantId", tenantId),
  275. new("@BatchId", batchId),
  276. new("@Now", now)
  277. };
  278. // 逻辑计数:本轮参与 KPI transform 的 distinct (tenant_id, sales_order_id) 数,
  279. // 即下方 CTE `agg` 的基数 —— 与写入集合一一对应(agg 是最终 SELECT 的驱动表)。
  280. var logicalRows = await _db.Ado.GetIntAsync($"""
  281. SELECT COUNT(*) FROM (
  282. SELECT `tenant_id`, `biz_id` {KpiEventFrom} GROUP BY `tenant_id`, `biz_id`
  283. ) agg
  284. """, ps);
  285. const string sql = $"""
  286. INSERT INTO `dwd_order_review_kpi`
  287. (`tenant_id`, `sales_order_id`, `sales_order_no`,
  288. `review_action_run_id`, `confirm_action_run_id`,
  289. `review_start_time`, `review_end_time`, `actual_hours`, `calculation_status`,
  290. `review_success_count`, `confirm_success_count`,
  291. `review_failed_count`, `confirm_failed_count`,
  292. `calc_batch_id`, `calc_time`)
  293. WITH ev AS (
  294. SELECT `tenant_id`, `biz_id`, `biz_no`, `action_run_id`, `action_code`,
  295. `action_status`, `action_start_time`, `action_end_time`
  296. {KpiEventFrom}
  297. ),
  298. agg AS (
  299. SELECT `tenant_id`, `biz_id`,
  300. MAX(`biz_no`) AS biz_no,
  301. SUM(`action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'SUCCESS') AS rs,
  302. SUM(`action_code` = 'S1_DELIVERY_CONFIRM' AND `action_status` = 'SUCCESS') AS cs,
  303. SUM(`action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'FAILED') AS rf,
  304. SUM(`action_code` = 'S1_DELIVERY_CONFIRM' AND `action_status` = 'FAILED') AS cf
  305. FROM ev GROUP BY `tenant_id`, `biz_id`
  306. ),
  307. rv AS (
  308. SELECT `tenant_id`, `biz_id`, `action_run_id`, `action_start_time`,
  309. ROW_NUMBER() OVER (PARTITION BY `tenant_id`, `biz_id`
  310. ORDER BY `action_start_time` ASC, `action_run_id` ASC) AS rn
  311. FROM ev
  312. WHERE `action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'SUCCESS'
  313. ),
  314. r1 AS (SELECT * FROM rv WHERE rn = 1),
  315. cv AS (
  316. SELECT c.`tenant_id`, c.`biz_id`, c.`action_run_id`, c.`action_end_time`,
  317. ROW_NUMBER() OVER (PARTITION BY c.`tenant_id`, c.`biz_id`
  318. ORDER BY c.`action_end_time` ASC, c.`action_run_id` ASC) AS rn
  319. FROM ev c
  320. JOIN r1 ON r1.`tenant_id` = c.`tenant_id` AND r1.`biz_id` = c.`biz_id`
  321. WHERE c.`action_code` = 'S1_DELIVERY_CONFIRM'
  322. AND c.`action_status` = 'SUCCESS'
  323. AND c.`action_end_time` IS NOT NULL
  324. AND c.`action_end_time` > r1.`action_start_time`
  325. ),
  326. c1 AS (SELECT * FROM cv WHERE rn = 1)
  327. SELECT
  328. a.`tenant_id`,
  329. a.`biz_id`,
  330. a.biz_no,
  331. r1.`action_run_id`,
  332. c1.`action_run_id`,
  333. r1.`action_start_time`,
  334. c1.`action_end_time`,
  335. CASE WHEN r1.`action_start_time` IS NOT NULL AND c1.`action_end_time` IS NOT NULL
  336. THEN ROUND(TIMESTAMPDIFF(SECOND, r1.`action_start_time`, c1.`action_end_time`) / 3600.0, 2)
  337. END,
  338. CASE
  339. WHEN r1.`action_start_time` IS NULL AND a.cs = 0 THEN 'NO_START'
  340. WHEN r1.`action_start_time` IS NULL THEN 'INVALID'
  341. WHEN c1.`action_end_time` IS NOT NULL THEN 'OK'
  342. WHEN a.cs = 0 THEN 'IN_PROGRESS'
  343. ELSE 'INVALID_TIME_SEQUENCE'
  344. END,
  345. a.rs, a.cs, a.rf, a.cf,
  346. @BatchId, @Now
  347. FROM agg a
  348. LEFT JOIN r1 ON r1.`tenant_id` = a.`tenant_id` AND r1.`biz_id` = a.`biz_id`
  349. LEFT JOIN c1 ON c1.`tenant_id` = a.`tenant_id` AND c1.`biz_id` = a.`biz_id`
  350. ON DUPLICATE KEY UPDATE
  351. `sales_order_no` = VALUES(`sales_order_no`),
  352. `review_action_run_id` = VALUES(`review_action_run_id`),
  353. `confirm_action_run_id` = VALUES(`confirm_action_run_id`),
  354. `review_start_time` = VALUES(`review_start_time`),
  355. `review_end_time` = VALUES(`review_end_time`),
  356. `actual_hours` = VALUES(`actual_hours`),
  357. `calculation_status` = VALUES(`calculation_status`),
  358. `review_success_count` = VALUES(`review_success_count`),
  359. `confirm_success_count` = VALUES(`confirm_success_count`),
  360. `review_failed_count` = VALUES(`review_failed_count`),
  361. `confirm_failed_count` = VALUES(`confirm_failed_count`),
  362. `calc_batch_id` = VALUES(`calc_batch_id`),
  363. `calc_time` = VALUES(`calc_time`)
  364. """;
  365. // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
  366. await _db.Ado.ExecuteCommandAsync(sql, ps);
  367. return logicalRows;
  368. }
  369. /// <summary>
  370. /// KPI transform 的事件取数范围(FROM + WHERE)。被计数查询与 CTE <c>ev</c> 共用,
  371. /// 保证「计的订单数」与「写入的订单集合」同源。
  372. /// </summary>
  373. private const string KpiEventFrom = """
  374. FROM `mdp_std_action_event`
  375. WHERE `tenant_id` = @TenantId
  376. AND `biz_type` = 'crm_seorder'
  377. AND `biz_id` IS NOT NULL
  378. AND `action_code` IN ('S1_ORDER_REVIEW', 'S1_DELIVERY_CONFIRM')
  379. """;
  380. }
  381. /// <summary>
  382. /// 动作事件链路跑批结果。
  383. ///
  384. /// <para><b>计数语义统一为「本轮参与处理的逻辑对象数」</b>,不是数据库 affected rows。
  385. /// 原因:<c>INSERT … ON DUPLICATE KEY UPDATE</c> 在 MySQL 下「插入计 1、更新计 2」,
  386. /// 直接返回会让「什么都没变」的第二次跑批显示成处理量翻倍,属于虚报。
  387. /// 因此每层都改为对该层的取数集合求 <c>COUNT(DISTINCT 业务键)</c>:
  388. /// 一个 <c>action_run_id</c> 最多计 1、一张销售订单最多计 1,重复 upsert 不翻倍。</para>
  389. /// </summary>
  390. public sealed class ActionEventSyncResult
  391. {
  392. public string BatchId { get; set; } = "";
  393. public long TenantId { get; set; }
  394. /// <summary>本轮实际进入贴源处理链的 source event 数(增量拉取到的源行数)。</summary>
  395. public int StageRows { get; set; }
  396. /// <summary>本轮被 RUNNING 回捞补偿重拉的 distinct action_run_id 数。</summary>
  397. public int StageRepulledRunning { get; set; }
  398. /// <summary>本轮参与 STD transform 的 distinct action_run_id 数。</summary>
  399. public int StandardRows { get; set; }
  400. /// <summary>本轮参与 KPI transform 的 distinct 销售订单数。</summary>
  401. public int KpiRows { get; set; }
  402. }