|
|
@@ -0,0 +1,441 @@
|
|
|
+using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
|
|
|
+using Microsoft.Extensions.Logging;
|
|
|
+using SqlSugar;
|
|
|
+
|
|
|
+namespace Admin.NET.Plugin.AiDOP.Order;
|
|
|
+
|
|
|
+/// <summary>
|
|
|
+/// S8 Stage-1「订单评审」ActualHours 数据链:
|
|
|
+/// <c>aidop_action_run_log → mdp_stg_action_run → mdp_std_action_event → dwd_order_review_kpi</c>。
|
|
|
+///
|
|
|
+/// <para><b>为什么必须经中台</b>:S8 正式运行时禁止直读 <c>aidop_action_run_log</c> /
|
|
|
+/// <c>crm_seorder</c> / <c>crm_seorderentry</c>。配对、FAILED 过滤、重复动作处理全部在本服务完成,
|
|
|
+/// S8 只消费 <c>dwd_order_review_kpi</c>。</para>
|
|
|
+///
|
|
|
+/// <para><b>租户是唯一隔离边界</b>:每一层的读、写、JOIN、唯一键首列都带 <c>tenant_id</c>;
|
|
|
+/// 源侧 <c>tenant_id</c> 为 NULL 或 0 一律拒绝入站,绝不补默认租户。</para>
|
|
|
+/// </summary>
|
|
|
+public sealed class ActionEventMdpSyncService : ITransient
|
|
|
+{
|
|
|
+ /// <summary>MDP 实体码(注册见 UpdateScripts/1.0.527.sql)。</summary>
|
|
|
+ public const string EntityCode = "S1_ACTION_RUN_LOG";
|
|
|
+
|
|
|
+ private const string SourceTable = "aidop_action_run_log";
|
|
|
+ private const string ReviewActionCode = "S1_ORDER_REVIEW";
|
|
|
+ private const string ConfirmActionCode = "S1_DELIVERY_CONFIRM";
|
|
|
+ private const string OrderBizType = "crm_seorder";
|
|
|
+
|
|
|
+ private readonly ISqlSugarClient _db;
|
|
|
+ private readonly MdpModuleStagingPuller _stagingPuller;
|
|
|
+ private readonly ILogger<ActionEventMdpSyncService> _logger;
|
|
|
+
|
|
|
+ public ActionEventMdpSyncService(
|
|
|
+ ISqlSugarClient db,
|
|
|
+ MdpModuleStagingPuller stagingPuller,
|
|
|
+ ILogger<ActionEventMdpSyncService> logger)
|
|
|
+ {
|
|
|
+ _db = db;
|
|
|
+ _stagingPuller = stagingPuller;
|
|
|
+ _logger = logger;
|
|
|
+ }
|
|
|
+
|
|
|
+ /// <summary>
|
|
|
+ /// 全链路跑批:源→STG(增量 + RUNNING 回捞)→ STD → DWD KPI。
|
|
|
+ /// 幂等:同一份源数据重复执行,STD 与 DWD 结果完全一致。
|
|
|
+ /// </summary>
|
|
|
+ public async Task<ActionEventSyncResult> RunAsync(long tenantId, CancellationToken cancellationToken = default)
|
|
|
+ {
|
|
|
+ if (tenantId <= 0)
|
|
|
+ throw new InvalidOperationException("ActionEventMdpSyncService 必须显式传入有效 tenantId,禁止默认租户。");
|
|
|
+
|
|
|
+ var now = DateTime.Now;
|
|
|
+ var batchId = $"S1_ACTION_EVENT_{tenantId}_{now:yyyyMMddHHmmss}";
|
|
|
+ var result = new ActionEventSyncResult { BatchId = batchId, TenantId = tenantId };
|
|
|
+
|
|
|
+ result.StageRows = await PullStagingAsync(tenantId, batchId, cancellationToken);
|
|
|
+ result.StageRepulledRunning = await RepullRunningAsync(tenantId, batchId, now, cancellationToken);
|
|
|
+ result.StandardRows = await TransformStdFromStgAsync(tenantId, batchId, cancellationToken);
|
|
|
+ result.KpiRows = await BuildOrderReviewKpiAsync(tenantId, batchId, now, cancellationToken);
|
|
|
+
|
|
|
+ _logger.LogInformation(
|
|
|
+ "S1 动作事件链路完成 tenant={TenantId} batch={BatchId} stg={Stg} repull={Repull} std={Std} kpi={Kpi}",
|
|
|
+ tenantId, batchId, result.StageRows, result.StageRepulledRunning, result.StandardRows, result.KpiRows);
|
|
|
+
|
|
|
+ return result;
|
|
|
+ }
|
|
|
+
|
|
|
+ // ══════════════════════════════════════════════════════════════════════
|
|
|
+ // PHASE 2 · 源 → STG
|
|
|
+ // ══════════════════════════════════════════════════════════════════════
|
|
|
+
|
|
|
+ /// <summary>
|
|
|
+ /// 主增量:走统一执行器(<c>sync_mode=INCR</c>,主游标 <c>start_time</c>,水位存
|
|
|
+ /// <c>mdp_entity.last_cursor</c>)。写入经 <see cref="MdpStagingWriter"/> 的 canonical 契约,
|
|
|
+ /// 其中已内置「tenant 解析不出即跳过 + 跨租户行丢弃」双重守卫。
|
|
|
+ /// </summary>
|
|
|
+ private Task<int> PullStagingAsync(long tenantId, string batchId, CancellationToken cancellationToken) =>
|
|
|
+ _stagingPuller.PullEntitiesAsync(
|
|
|
+ [EntityCode], batchId, tenantId,
|
|
|
+ fullRefresh: false,
|
|
|
+ taskCode: "S1_ACTION_EVENT_INBOUND",
|
|
|
+ cancellationToken,
|
|
|
+ factoryId: 0,
|
|
|
+ requireMatchingSourceTenant: true);
|
|
|
+
|
|
|
+ /// <summary>
|
|
|
+ /// <b>RUNNING 回捞补偿(本批的核心风险控制点)。</b>
|
|
|
+ ///
|
|
|
+ /// <para>源行不是内容级 append-only:<c>StartAsync</c> 先 INSERT 一行
|
|
|
+ /// <c>status='RUNNING', end_time=NULL</c>,之后 <c>FinishAsync</c> /
|
|
|
+ /// <c>FailStaleRunningAsync</c> 会<b>原地 UPDATE</b> 成 SUCCESS/FAILED 并补 <c>end_time</c>。
|
|
|
+ /// 而源表<b>没有 update_time 列</b>,且 <c>create_time ≡ start_time</c>、不随 UPDATE 推进
|
|
|
+ /// —— 所以主游标 <c>start_time > @cursor</c> 永远不会第二次捞到这一行,
|
|
|
+ /// 若无补偿,STG/STD 会永久停留在 RUNNING。</para>
|
|
|
+ ///
|
|
|
+ /// <para>补偿判据:以<b>标准层当前仍为 RUNNING 的 action_run_id 集合</b>反向驱动重拉。
|
|
|
+ /// 该集合天然极小(实测全库仅 1 行),因此:
|
|
|
+ /// ① tenant-safe —— 读写两侧都带 <c>tenant_id=@TenantId</c>;
|
|
|
+ /// ② 幂等 —— <c>ON DUPLICATE KEY UPDATE</c>,重复执行不产生新行;
|
|
|
+ /// ③ 不全表扫描 —— 由 RUNNING 集合驱动,走 <c>uk_std_action_event</c> 与源主键;
|
|
|
+ /// ④ 不依赖 id 单调性 —— 按 id 精确等值匹配,与雪花乱序无关;
|
|
|
+ /// ⑤ 不遗漏迟到终态 —— 只要某行在 STD 还是 RUNNING,每轮都会被重新检查。</para>
|
|
|
+ ///
|
|
|
+ /// <para>此处不经 <see cref="MdpStagingWriter"/> 而用定向 <c>INSERT…SELECT</c>,
|
|
|
+ /// 原因是框架执行器只支持「按水位翻页」、不支持「按 id 集合定向重拉」;
|
|
|
+ /// 因此租户守卫在本 SQL 内显式重建(<c>tenant_id=@TenantId</c> +
|
|
|
+ /// <c>tenant_id IS NOT NULL</c> + <c>tenant_id<>0</c>),语义与框架守卫一致。</para>
|
|
|
+ /// </summary>
|
|
|
+ private async Task<int> RepullRunningAsync(
|
|
|
+ long tenantId, string batchId, DateTime now, CancellationToken cancellationToken)
|
|
|
+ {
|
|
|
+ cancellationToken.ThrowIfCancellationRequested();
|
|
|
+
|
|
|
+ var ps = new SugarParameter[]
|
|
|
+ {
|
|
|
+ new("@TenantId", tenantId),
|
|
|
+ new("@BatchId", batchId),
|
|
|
+ new("@Now", now)
|
|
|
+ };
|
|
|
+
|
|
|
+ // 逻辑计数:本轮被回捞的 distinct action_run_id 数。
|
|
|
+ // 必须先于 upsert 求值(upsert 不改变该集合,但顺序固定更易解释),
|
|
|
+ // 且与下方 INSERT 共用同一份 WHERE(RepullRunningFrom),杜绝谓词漂移。
|
|
|
+ var logicalRows = await _db.Ado.GetIntAsync(
|
|
|
+ $"SELECT COUNT(DISTINCT s.`id`) {RepullRunningFrom}", ps);
|
|
|
+
|
|
|
+ const string sql = $"""
|
|
|
+ INSERT INTO `mdp_stg_action_run`
|
|
|
+ (`tenant_id`, `source_system`, `source_table`, `source_row_id`, `source_biz_key`,
|
|
|
+ `raw_data`, `sync_batch_id`, `sync_time`, `process_status`, `create_time`, `update_time`)
|
|
|
+ SELECT
|
|
|
+ s.`tenant_id`,
|
|
|
+ 'AIDOPDEV_MYSQL',
|
|
|
+ 'aidop_action_run_log',
|
|
|
+ CAST(s.`id` AS CHAR),
|
|
|
+ CAST(s.`id` AS CHAR),
|
|
|
+ JSON_OBJECT(
|
|
|
+ 'id', CAST(s.`id` AS CHAR),
|
|
|
+ 'tenant_id', s.`tenant_id`,
|
|
|
+ 'action_code', s.`action_code`,
|
|
|
+ 'biz_type', s.`biz_type`,
|
|
|
+ 'biz_id', s.`biz_id`,
|
|
|
+ 'biz_no', s.`biz_no`,
|
|
|
+ 'status', s.`status`,
|
|
|
+ 'message', s.`message`,
|
|
|
+ 'detail_json', s.`detail_json`,
|
|
|
+ 'start_time', DATE_FORMAT(s.`start_time`, '%Y-%m-%d %H:%i:%s'),
|
|
|
+ 'end_time', DATE_FORMAT(s.`end_time`, '%Y-%m-%d %H:%i:%s'),
|
|
|
+ 'create_time', DATE_FORMAT(s.`create_time`,'%Y-%m-%d %H:%i:%s')
|
|
|
+ ),
|
|
|
+ @BatchId, @Now, 'PENDING', @Now, @Now
|
|
|
+ {RepullRunningFrom}
|
|
|
+ ON DUPLICATE KEY UPDATE
|
|
|
+ `source_row_id` = VALUES(`source_row_id`),
|
|
|
+ `raw_data` = VALUES(`raw_data`),
|
|
|
+ `sync_batch_id` = VALUES(`sync_batch_id`),
|
|
|
+ `sync_time` = VALUES(`sync_time`),
|
|
|
+ `process_status` = 'PENDING',
|
|
|
+ `update_time` = VALUES(`update_time`)
|
|
|
+ """;
|
|
|
+
|
|
|
+ // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
|
|
|
+ await _db.Ado.ExecuteCommandAsync(sql, ps);
|
|
|
+
|
|
|
+ if (logicalRows > 0)
|
|
|
+ _logger.LogInformation("RUNNING 回捞补偿命中 tenant={TenantId} rows={Rows}", tenantId, logicalRows);
|
|
|
+
|
|
|
+ return logicalRows;
|
|
|
+ }
|
|
|
+
|
|
|
+ /// <summary>
|
|
|
+ /// RUNNING 回捞的取数范围(FROM + WHERE)。被计数查询与 upsert 共用,保证两者永远同一集合。
|
|
|
+ /// </summary>
|
|
|
+ private const string RepullRunningFrom = """
|
|
|
+ FROM `aidop_action_run_log` s
|
|
|
+ WHERE s.`tenant_id` = @TenantId
|
|
|
+ AND s.`tenant_id` IS NOT NULL
|
|
|
+ AND s.`tenant_id` <> 0
|
|
|
+ AND EXISTS (
|
|
|
+ SELECT 1 FROM `mdp_std_action_event` e
|
|
|
+ WHERE e.`tenant_id` = @TenantId
|
|
|
+ AND e.`action_status` = 'RUNNING'
|
|
|
+ AND e.`action_run_id` = s.`id`
|
|
|
+ )
|
|
|
+ """;
|
|
|
+
|
|
|
+ // ══════════════════════════════════════════════════════════════════════
|
|
|
+ // PHASE 3 · STG → STD(审计事实层,全状态保留)
|
|
|
+ // ══════════════════════════════════════════════════════════════════════
|
|
|
+
|
|
|
+ /// <summary>
|
|
|
+ /// 贴源 → 标准事件层。Grain = <c>tenant_id + action_run_id</c>。
|
|
|
+ /// <para>RUNNING / SUCCESS / FAILED <b>全部保留</b> —— STD 是审计事实层,
|
|
|
+ /// 不因为 KPI 只需要 SUCCESS 就在此丢弃其它状态。</para>
|
|
|
+ /// <para>同一 <c>action_run_id</c> 重拉时,<c>ON DUPLICATE KEY UPDATE</c> 允许
|
|
|
+ /// RUNNING 被 SUCCESS/FAILED 更新(终态推进)。</para>
|
|
|
+ /// </summary>
|
|
|
+ public async Task<int> TransformStdFromStgAsync(
|
|
|
+ long tenantId, string batchId, CancellationToken cancellationToken = default)
|
|
|
+ {
|
|
|
+ if (tenantId <= 0) throw new InvalidOperationException("TransformStdFromStgAsync 需要有效 tenantId。");
|
|
|
+ cancellationToken.ThrowIfCancellationRequested();
|
|
|
+
|
|
|
+ var ps = new SugarParameter[]
|
|
|
+ {
|
|
|
+ new("@TenantId", tenantId),
|
|
|
+ new("@BatchId", batchId)
|
|
|
+ };
|
|
|
+
|
|
|
+ // 逻辑计数:本轮参与 STD transform 的 distinct action_run_id 数。
|
|
|
+ // 与下方 INSERT 共用同一份 FROM/WHERE(StdEligibleFrom),保证「计了什么」= 「写了什么」。
|
|
|
+ var logicalRows = await _db.Ado.GetIntAsync(
|
|
|
+ $"SELECT COUNT(DISTINCT JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id'))) {StdEligibleFrom}", ps);
|
|
|
+
|
|
|
+ const string sql = $"""
|
|
|
+ INSERT INTO `mdp_std_action_event`
|
|
|
+ (`tenant_id`, `action_run_id`, `action_code`, `biz_type`, `biz_id`, `biz_no`,
|
|
|
+ `action_status`, `action_start_time`, `action_end_time`, `message`, `detail_json`,
|
|
|
+ `source_system`, `source_table`, `source_biz_key`, `sync_batch_id`)
|
|
|
+ SELECT
|
|
|
+ g.`tenant_id`,
|
|
|
+ CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id')) AS SIGNED),
|
|
|
+ JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.action_code')),
|
|
|
+ JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_type')),
|
|
|
+ CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_id')) REGEXP '^-?[0-9]+$'
|
|
|
+ THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_id')) AS SIGNED) END,
|
|
|
+ NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_no')), 'null'),
|
|
|
+ JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.status')),
|
|
|
+ CAST(REPLACE(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.start_time')), 'T', ' ') AS DATETIME),
|
|
|
+ CASE WHEN NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.end_time')), 'null') IS NULL
|
|
|
+ THEN NULL
|
|
|
+ ELSE CAST(REPLACE(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.end_time')), 'T', ' ') AS DATETIME) END,
|
|
|
+ NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.message')), 'null'),
|
|
|
+ CASE
|
|
|
+ WHEN JSON_EXTRACT(g.`raw_data`, '$.detail_json') IS NULL THEN NULL
|
|
|
+ WHEN JSON_VALID(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.detail_json')))
|
|
|
+ THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.detail_json')) AS JSON)
|
|
|
+ ELSE JSON_EXTRACT(g.`raw_data`, '$.detail_json')
|
|
|
+ END,
|
|
|
+ g.`source_system`,
|
|
|
+ g.`source_table`,
|
|
|
+ g.`source_biz_key`,
|
|
|
+ @BatchId
|
|
|
+ {StdEligibleFrom}
|
|
|
+ ON DUPLICATE KEY UPDATE
|
|
|
+ `action_code` = VALUES(`action_code`),
|
|
|
+ `biz_type` = VALUES(`biz_type`),
|
|
|
+ `biz_id` = VALUES(`biz_id`),
|
|
|
+ `biz_no` = VALUES(`biz_no`),
|
|
|
+ `action_status` = VALUES(`action_status`),
|
|
|
+ `action_start_time` = VALUES(`action_start_time`),
|
|
|
+ `action_end_time` = VALUES(`action_end_time`),
|
|
|
+ `message` = VALUES(`message`),
|
|
|
+ `detail_json` = VALUES(`detail_json`),
|
|
|
+ `sync_batch_id` = VALUES(`sync_batch_id`)
|
|
|
+ """;
|
|
|
+
|
|
|
+ // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
|
|
|
+ await _db.Ado.ExecuteCommandAsync(sql, ps);
|
|
|
+ return logicalRows;
|
|
|
+ }
|
|
|
+
|
|
|
+ /// <summary>
|
|
|
+ /// STD transform 的取数范围(FROM + WHERE)。被计数查询与 upsert 共用,保证两者永远同一集合。
|
|
|
+ /// </summary>
|
|
|
+ private const string StdEligibleFrom = """
|
|
|
+ FROM `mdp_stg_action_run` g
|
|
|
+ WHERE g.`tenant_id` = @TenantId
|
|
|
+ AND g.`tenant_id` <> 0
|
|
|
+ AND g.`source_table` = 'aidop_action_run_log'
|
|
|
+ AND JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.id')) REGEXP '^[0-9]+$'
|
|
|
+ AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.action_code')), 'null') IS NOT NULL
|
|
|
+ AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.biz_type')), 'null') IS NOT NULL
|
|
|
+ AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.status')), 'null') IS NOT NULL
|
|
|
+ AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(g.`raw_data`, '$.start_time')), 'null') IS NOT NULL
|
|
|
+ """;
|
|
|
+
|
|
|
+ // ══════════════════════════════════════════════════════════════════════
|
|
|
+ // PHASE 4 · STD → DWD KPI(FIRST VALID PAIR)
|
|
|
+ // ══════════════════════════════════════════════════════════════════════
|
|
|
+
|
|
|
+ /// <summary>
|
|
|
+ /// 标准事件层 → 订单评审 KPI。Grain = <c>tenant_id + sales_order_id</c>。
|
|
|
+ ///
|
|
|
+ /// <para><b>算法(定版,不得改用 LAST PAIR / MAX)</b>:
|
|
|
+ /// ① 只取 <c>status='SUCCESS'</c>,FAILED 保留在 STD 但不参与配对;
|
|
|
+ /// ② ReviewStart = 评审 SUCCESS 的 <b>FIRST</b>(<c>ORDER BY action_start_time, action_run_id</c>);
|
|
|
+ /// ③ ReviewEnd = <b>首个 <c>action_end_time > ReviewStart</c></b> 的确认 SUCCESS 的
|
|
|
+ /// <c>action_end_time</c>(= 确认动作<b>完成</b>时刻,是 <c>progress='3'</c> 提交后的最近上界);
|
|
|
+ /// ④ ActualHours = 自然 elapsed 小时,2 位小数,<b>不扣夜间/周末/节假日</b>(本项目无工作日历 Authority)。</para>
|
|
|
+ ///
|
|
|
+ /// <para><b>冻结语义</b>:算法只依赖「第一条评审」与「首个晚于它的确认」,两者一旦产生即恒定,
|
|
|
+ /// 后续任意次重复点击都不会改写已形成的 KPI;因此重复执行结果完全一致(deterministic / idempotent)。</para>
|
|
|
+ ///
|
|
|
+ /// <para>本批<b>不计算</b> TargetHours / kpi_status —— 那依赖 S0 PI 匹配,属另一批次。</para>
|
|
|
+ /// </summary>
|
|
|
+ public async Task<int> BuildOrderReviewKpiAsync(
|
|
|
+ long tenantId, string batchId, DateTime now, CancellationToken cancellationToken = default)
|
|
|
+ {
|
|
|
+ if (tenantId <= 0) throw new InvalidOperationException("BuildOrderReviewKpiAsync 需要有效 tenantId。");
|
|
|
+ cancellationToken.ThrowIfCancellationRequested();
|
|
|
+
|
|
|
+ var ps = new SugarParameter[]
|
|
|
+ {
|
|
|
+ new("@TenantId", tenantId),
|
|
|
+ new("@BatchId", batchId),
|
|
|
+ new("@Now", now)
|
|
|
+ };
|
|
|
+
|
|
|
+ // 逻辑计数:本轮参与 KPI transform 的 distinct (tenant_id, sales_order_id) 数,
|
|
|
+ // 即下方 CTE `agg` 的基数 —— 与写入集合一一对应(agg 是最终 SELECT 的驱动表)。
|
|
|
+ var logicalRows = await _db.Ado.GetIntAsync($"""
|
|
|
+ SELECT COUNT(*) FROM (
|
|
|
+ SELECT `tenant_id`, `biz_id` {KpiEventFrom} GROUP BY `tenant_id`, `biz_id`
|
|
|
+ ) agg
|
|
|
+ """, ps);
|
|
|
+
|
|
|
+ const string sql = $"""
|
|
|
+ INSERT INTO `dwd_order_review_kpi`
|
|
|
+ (`tenant_id`, `sales_order_id`, `sales_order_no`,
|
|
|
+ `review_action_run_id`, `confirm_action_run_id`,
|
|
|
+ `review_start_time`, `review_end_time`, `actual_hours`, `calculation_status`,
|
|
|
+ `review_success_count`, `confirm_success_count`,
|
|
|
+ `review_failed_count`, `confirm_failed_count`,
|
|
|
+ `calc_batch_id`, `calc_time`)
|
|
|
+ WITH ev AS (
|
|
|
+ SELECT `tenant_id`, `biz_id`, `biz_no`, `action_run_id`, `action_code`,
|
|
|
+ `action_status`, `action_start_time`, `action_end_time`
|
|
|
+ {KpiEventFrom}
|
|
|
+ ),
|
|
|
+ agg AS (
|
|
|
+ SELECT `tenant_id`, `biz_id`,
|
|
|
+ MAX(`biz_no`) AS biz_no,
|
|
|
+ SUM(`action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'SUCCESS') AS rs,
|
|
|
+ SUM(`action_code` = 'S1_DELIVERY_CONFIRM' AND `action_status` = 'SUCCESS') AS cs,
|
|
|
+ SUM(`action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'FAILED') AS rf,
|
|
|
+ SUM(`action_code` = 'S1_DELIVERY_CONFIRM' AND `action_status` = 'FAILED') AS cf
|
|
|
+ FROM ev GROUP BY `tenant_id`, `biz_id`
|
|
|
+ ),
|
|
|
+ rv AS (
|
|
|
+ SELECT `tenant_id`, `biz_id`, `action_run_id`, `action_start_time`,
|
|
|
+ ROW_NUMBER() OVER (PARTITION BY `tenant_id`, `biz_id`
|
|
|
+ ORDER BY `action_start_time` ASC, `action_run_id` ASC) AS rn
|
|
|
+ FROM ev
|
|
|
+ WHERE `action_code` = 'S1_ORDER_REVIEW' AND `action_status` = 'SUCCESS'
|
|
|
+ ),
|
|
|
+ r1 AS (SELECT * FROM rv WHERE rn = 1),
|
|
|
+ cv AS (
|
|
|
+ SELECT c.`tenant_id`, c.`biz_id`, c.`action_run_id`, c.`action_end_time`,
|
|
|
+ ROW_NUMBER() OVER (PARTITION BY c.`tenant_id`, c.`biz_id`
|
|
|
+ ORDER BY c.`action_end_time` ASC, c.`action_run_id` ASC) AS rn
|
|
|
+ FROM ev c
|
|
|
+ JOIN r1 ON r1.`tenant_id` = c.`tenant_id` AND r1.`biz_id` = c.`biz_id`
|
|
|
+ WHERE c.`action_code` = 'S1_DELIVERY_CONFIRM'
|
|
|
+ AND c.`action_status` = 'SUCCESS'
|
|
|
+ AND c.`action_end_time` IS NOT NULL
|
|
|
+ AND c.`action_end_time` > r1.`action_start_time`
|
|
|
+ ),
|
|
|
+ c1 AS (SELECT * FROM cv WHERE rn = 1)
|
|
|
+ SELECT
|
|
|
+ a.`tenant_id`,
|
|
|
+ a.`biz_id`,
|
|
|
+ a.biz_no,
|
|
|
+ r1.`action_run_id`,
|
|
|
+ c1.`action_run_id`,
|
|
|
+ r1.`action_start_time`,
|
|
|
+ c1.`action_end_time`,
|
|
|
+ CASE WHEN r1.`action_start_time` IS NOT NULL AND c1.`action_end_time` IS NOT NULL
|
|
|
+ THEN ROUND(TIMESTAMPDIFF(SECOND, r1.`action_start_time`, c1.`action_end_time`) / 3600.0, 2)
|
|
|
+ END,
|
|
|
+ CASE
|
|
|
+ WHEN r1.`action_start_time` IS NULL AND a.cs = 0 THEN 'NO_START'
|
|
|
+ WHEN r1.`action_start_time` IS NULL THEN 'INVALID'
|
|
|
+ WHEN c1.`action_end_time` IS NOT NULL THEN 'OK'
|
|
|
+ WHEN a.cs = 0 THEN 'IN_PROGRESS'
|
|
|
+ ELSE 'INVALID_TIME_SEQUENCE'
|
|
|
+ END,
|
|
|
+ a.rs, a.cs, a.rf, a.cf,
|
|
|
+ @BatchId, @Now
|
|
|
+ FROM agg a
|
|
|
+ LEFT JOIN r1 ON r1.`tenant_id` = a.`tenant_id` AND r1.`biz_id` = a.`biz_id`
|
|
|
+ LEFT JOIN c1 ON c1.`tenant_id` = a.`tenant_id` AND c1.`biz_id` = a.`biz_id`
|
|
|
+ ON DUPLICATE KEY UPDATE
|
|
|
+ `sales_order_no` = VALUES(`sales_order_no`),
|
|
|
+ `review_action_run_id` = VALUES(`review_action_run_id`),
|
|
|
+ `confirm_action_run_id` = VALUES(`confirm_action_run_id`),
|
|
|
+ `review_start_time` = VALUES(`review_start_time`),
|
|
|
+ `review_end_time` = VALUES(`review_end_time`),
|
|
|
+ `actual_hours` = VALUES(`actual_hours`),
|
|
|
+ `calculation_status` = VALUES(`calculation_status`),
|
|
|
+ `review_success_count` = VALUES(`review_success_count`),
|
|
|
+ `confirm_success_count` = VALUES(`confirm_success_count`),
|
|
|
+ `review_failed_count` = VALUES(`review_failed_count`),
|
|
|
+ `confirm_failed_count` = VALUES(`confirm_failed_count`),
|
|
|
+ `calc_batch_id` = VALUES(`calc_batch_id`),
|
|
|
+ `calc_time` = VALUES(`calc_time`)
|
|
|
+ """;
|
|
|
+
|
|
|
+ // 返回值刻意丢弃:ODKU 的 affected rows 把「更新」计 2,不是逻辑记录数。
|
|
|
+ await _db.Ado.ExecuteCommandAsync(sql, ps);
|
|
|
+ return logicalRows;
|
|
|
+ }
|
|
|
+
|
|
|
+ /// <summary>
|
|
|
+ /// KPI transform 的事件取数范围(FROM + WHERE)。被计数查询与 CTE <c>ev</c> 共用,
|
|
|
+ /// 保证「计的订单数」与「写入的订单集合」同源。
|
|
|
+ /// </summary>
|
|
|
+ private const string KpiEventFrom = """
|
|
|
+ FROM `mdp_std_action_event`
|
|
|
+ WHERE `tenant_id` = @TenantId
|
|
|
+ AND `biz_type` = 'crm_seorder'
|
|
|
+ AND `biz_id` IS NOT NULL
|
|
|
+ AND `action_code` IN ('S1_ORDER_REVIEW', 'S1_DELIVERY_CONFIRM')
|
|
|
+ """;
|
|
|
+}
|
|
|
+
|
|
|
+/// <summary>
|
|
|
+/// 动作事件链路跑批结果。
|
|
|
+///
|
|
|
+/// <para><b>计数语义统一为「本轮参与处理的逻辑对象数」</b>,不是数据库 affected rows。
|
|
|
+/// 原因:<c>INSERT … ON DUPLICATE KEY UPDATE</c> 在 MySQL 下「插入计 1、更新计 2」,
|
|
|
+/// 直接返回会让「什么都没变」的第二次跑批显示成处理量翻倍,属于虚报。
|
|
|
+/// 因此每层都改为对该层的取数集合求 <c>COUNT(DISTINCT 业务键)</c>:
|
|
|
+/// 一个 <c>action_run_id</c> 最多计 1、一张销售订单最多计 1,重复 upsert 不翻倍。</para>
|
|
|
+/// </summary>
|
|
|
+public sealed class ActionEventSyncResult
|
|
|
+{
|
|
|
+ public string BatchId { get; set; } = "";
|
|
|
+ public long TenantId { get; set; }
|
|
|
+
|
|
|
+ /// <summary>本轮实际进入贴源处理链的 source event 数(增量拉取到的源行数)。</summary>
|
|
|
+ public int StageRows { get; set; }
|
|
|
+
|
|
|
+ /// <summary>本轮被 RUNNING 回捞补偿重拉的 distinct action_run_id 数。</summary>
|
|
|
+ public int StageRepulledRunning { get; set; }
|
|
|
+
|
|
|
+ /// <summary>本轮参与 STD transform 的 distinct action_run_id 数。</summary>
|
|
|
+ public int StandardRows { get; set; }
|
|
|
+
|
|
|
+ /// <summary>本轮参与 KPI transform 的 distinct 销售订单数。</summary>
|
|
|
+ public int KpiRows { get; set; }
|
|
|
+}
|