using Admin.NET.Plugin.AiDOP.DataPlatform.Executors; namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse; /// /// S5 采购收货单 数据中台只读同步转换服务(DOP 内部,方案 1)。 /// /// 源:aidopdev PurOrdRctDetail(p) + PurOrdRctMaster(d),业务类型 RctType='rc'(采购收货); /// 头明细关联 p.Domain=d.Domain AND p.Receiver=d.Receiver; /// 维表 ItemMaster(i) / SuppMaster(s) / ConsigneeAddressMaster(a) / PurOrdDetail(pd,sd) / srm_pr_main(dr) 全 LEFT JOIN。 /// 链路(Phase 1):执行器抽主/明细 → mdp_stg_purchase_receipt → mdp_std_purchase_receipt(typed)。 /// 双模式入站:mdp_entity=S5_PURCHASE_RECEIPT_DETAIL/MASTER(DB 或 API)→ MdpStagingWriter → stg, /// 标准层读 stg.raw_data;维表 ItemMaster/SuppMaster 等仍直连本库(本地口径补全)。 /// /// 约束: /// - 只读源/贴源,仅写 mdp_stg_purchase_receipt / mdp_std_purchase_receipt;绝不写源表;不读/不改 S3 mdp_stg_receipt、S4 ado_s4_receipt。 /// - 过滤 RctType='rc';租户隔离在读 API 侧按 tenant_id。 /// - rc 为 0 行时成功完成、处理数为 0,不报错。 /// public class PurchaseReceiptMdpSyncService : ITransient { private const string JobCode = "S5_PURCHASE_RECEIPT_MDP_SYNC"; private const string InboundEntityCode = "S5_PURCHASE_RECEIPT_DETAIL"; private const string InboundMasterEntityCode = "S5_PURCHASE_RECEIPT_MASTER"; private readonly ISqlSugarClient _db; private readonly MdpSourcePullDispatcher _pullDispatcher; public PurchaseReceiptMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher) { _db = db; _pullDispatcher = pullDispatcher; } /// 全量:本地 DB 执行器灌 stg → 标准层(读 stg)。 public async Task RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO") { cancellationToken.ThrowIfCancellationRequested(); await EnsureTablesAsync(); await EnsureStgTableAsync(); var now = DateTime.Now; var batchId = $"S5_PUR_RCT_FULL_{now:yyyyMMddHHmmss}"; var runLogId = await InsertRunLogAsync(batchId, now, triggerType); var result = new PurchaseReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId }; try { var pullCtx = new MdpPullContext { TenantId = 0, FullRefresh = true, TaskCode = "S5_PURCHASE_RECEIPT_INBOUND", BatchId = $"{batchId}_PULL" }; await PopulateStgAsync(pullCtx, cancellationToken); result.StdRows = await TransformStandardAsync(batchId, now); await MarkRunSuccessAsync(runLogId, now, result); return result; } catch (Exception ex) { await MarkRunFailedAsync(runLogId, now, ex.Message); throw; } } /// /// 双模式入站:执行器抽主/明细落 stg,再跑标准层(读 stg)。 /// public async Task RunInboundAsync( long tenantId = 0, bool fullRefresh = false, CancellationToken cancellationToken = default) { cancellationToken.ThrowIfCancellationRequested(); await EnsureTablesAsync(); await EnsureStgTableAsync(); var now = DateTime.Now; var pullCtx = new MdpPullContext { TenantId = tenantId, FullRefresh = fullRefresh, TaskCode = "S5_PURCHASE_RECEIPT_INBOUND", BatchId = $"S5_PUR_RCT_IN_{now:yyyyMMddHHmmss}" }; var pull = await PopulateStgAsync(pullCtx, cancellationToken); var batchId = $"S5_PUR_RCT_STD_{now:yyyyMMddHHmmss}"; var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND"); var result = new PurchaseReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId }; try { result.StdRows = await TransformStandardAsync(batchId, now); await MarkRunSuccessAsync(runLogId, now, result); } catch (Exception ex) { await MarkRunFailedAsync(runLogId, now, ex.Message); throw; } return new PurchaseReceiptInboundResult { PullBatchId = pullCtx.BatchId, RowsPulled = pull.RowsPulled, RowsWrittenStg = pull.RowsWritten, NewCursor = pull.NewCursor, PullMessage = pull.Message, TransformBatchId = result.BatchId, StdRows = result.StdRows }; } private async Task<(int RowsPulled, int RowsWritten, string? NewCursor, string? Message)> PopulateStgAsync( MdpPullContext pullCtx, CancellationToken cancellationToken) { var master = await _pullDispatcher.PullByEntityCodeAsync(InboundMasterEntityCode, pullCtx, cancellationToken); var detail = await _pullDispatcher.PullByEntityCodeAsync(InboundEntityCode, pullCtx, cancellationToken); return ( master.RowsPulled + detail.RowsPulled, master.RowsWritten + detail.RowsWritten, detail.NewCursor ?? master.NewCursor, $"{master.Message}; {detail.Message}"); } private async Task EnsureStgTableAsync() { await _db.Ado.ExecuteCommandAsync( """ CREATE TABLE IF NOT EXISTS mdp_stg_purchase_receipt ( id BIGINT PRIMARY KEY AUTO_INCREMENT, tenant_id BIGINT NOT NULL, source_system VARCHAR(50) NULL, source_table VARCHAR(200), source_row_id VARCHAR(200), source_biz_key VARCHAR(300) NULL, raw_data JSON, sync_batch_id VARCHAR(100), sync_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP, process_status VARCHAR(20) NOT NULL DEFAULT 'PENDING', process_message VARCHAR(500) NULL, create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, update_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key), KEY idx_batch (sync_batch_id), KEY idx_src (source_table, source_row_id), KEY idx_tenant (tenant_id) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='S5采购收货贴源层' """); } /// 防御式建表(与 UpdateScripts WIP-S5PR.sql / 正式 1.0.<n>.sql 同构,幂等)。 private async Task EnsureTablesAsync() { await _db.Ado.ExecuteCommandAsync( """ CREATE TABLE IF NOT EXISTS mdp_std_purchase_receipt ( id BIGINT AUTO_INCREMENT PRIMARY KEY, tenant_id BIGINT NOT NULL DEFAULT 0, factory_id BIGINT NULL DEFAULT 1, source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP', domain VARCHAR(24) NOT NULL, receiver VARCHAR(24) NOT NULL, line SMALLINT NOT NULL DEFAULT 0, rct_date DATETIME NULL, supp VARCHAR(20) NULL, sort_name VARCHAR(255) NULL, item_num VARCHAR(60) NULL, item_name VARCHAR(200) NULL, item_spec VARCHAR(200) NULL, um VARCHAR(8) NULL, qty_ordered DECIMAL(18,6) NULL DEFAULT 0, qty_received DECIMAL(18,6) NULL DEFAULT 0, lot_serial VARCHAR(120) NULL, location VARCHAR(8) NULL, ord_nbr VARCHAR(24) NULL, ord_line SMALLINT NULL, blanket_line INT NULL, pur_ord VARCHAR(24) NULL, pur_line SMALLINT NULL, sales_job VARCHAR(200) NULL, address1 VARCHAR(200) NULL, req VARCHAR(20) NULL, req_line INT NULL, dop_req VARCHAR(255) NULL, source_biz_key VARCHAR(200) NULL, sync_batch_id VARCHAR(100) NOT NULL, sync_time DATETIME NOT NULL, create_time DATETIME DEFAULT CURRENT_TIMESTAMP, update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, UNIQUE KEY uk_mdp_std_pur_rct (tenant_id, domain, receiver, line), KEY idx_mdp_std_pur_rct_date (tenant_id, rct_date), KEY idx_mdp_std_pur_rct_item (tenant_id, item_num), KEY idx_mdp_std_pur_rct_supp (tenant_id, supp), KEY idx_mdp_std_pur_rct_purord (tenant_id, pur_ord), KEY idx_mdp_std_pur_rct_salesjob (tenant_id, sales_job), KEY idx_mdp_std_pur_rct_req (tenant_id, req), KEY idx_mdp_std_pur_rct_dopreq (tenant_id, dop_req) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S5采购收货单标准层' """); } /// /// 标准化:stg(PurOrdRctDetail/Master) + 维表(本库) → mdp_std_purchase_receipt。 /// private async Task TransformStandardAsync(string batchId, DateTime now) { var rows = await _db.Ado.GetIntAsync( """ SELECT COUNT(1) FROM mdp_stg_purchase_receipt p INNER JOIN mdp_stg_purchase_receipt d ON d.source_table='PurOrdRctMaster' AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain')) AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Receiver')) WHERE p.source_table='PurOrdRctDetail' AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctType'))='rc' """); await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_std_purchase_receipt (tenant_id, factory_id, source_system, domain, receiver, line, rct_date, supp, sort_name, item_num, item_name, item_spec, um, qty_ordered, qty_received, lot_serial, location, ord_nbr, ord_line, blanket_line, pur_ord, pur_line, sales_job, address1, req, req_line, dop_req, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(p.tenant_id, 0), 1, IFNULL(NULLIF(p.source_system,''), 'AIDOP'), JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')), JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')), CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Line')) AS SIGNED), STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RctDate')),'null'),''), '%Y-%m-%d %H:%i:%s'), JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Supp')), TRIM(CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Supp')),''), ' ', IFNULL(s.SortName, ''))), JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.ItemNum')), i.Descr, i.Descr1, JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.UM')), CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.QtyOrded')) AS DECIMAL(18,6)), CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.QtyReceived')) AS DECIMAL(18,6)), JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.LotSerial')), JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Location')), JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdNbr')), CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdLine')) AS SIGNED), pd.BlanketLine, pd.PurOrd, pd.Line, pd.SalesJob, a.Address1, (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO' THEN '' ELSE sd.Req END), sd.ReqLine, (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO' THEN sd.Req ELSE dr.pr_billno END), IFNULL(NULLIF(p.source_biz_key,''), CONCAT( JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')), '#', JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')), '#', JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Line')))), @BatchId, @Now FROM mdp_stg_purchase_receipt p INNER JOIN mdp_stg_purchase_receipt d ON d.source_table='PurOrdRctMaster' AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain')) AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Receiver')) LEFT JOIN ItemMaster i ON JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = i.Domain AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.ItemNum')) = i.ItemNum LEFT JOIN SuppMaster s ON JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain')) = s.Domain AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Supp')) = s.Supp LEFT JOIN ConsigneeAddressMaster a ON JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = a.Domain AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Supp')) = a.Address AND a.Typed = 'Supp' LEFT JOIN PurOrdDetail pd ON JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = pd.Domain AND (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctNbr')),'')<>'' THEN JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctNbr')) ELSE JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdNbr')) END) = (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctNbr')),'')<>'' THEN pd.Contract ELSE pd.PurOrd END) AND (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctNbr')),'')<>'' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.BlanketLine')) AS SIGNED) ELSE CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdLine')) AS SIGNED) END) = pd.Line LEFT JOIN PurOrdDetail sd ON JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = sd.Domain AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdNbr')) = sd.PurOrd AND CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdLine')) AS SIGNED) = sd.Line LEFT JOIN srm_pr_main dr ON CAST(dr.factory_id AS CHAR) = sd.Domain AND dr.SAP_pr_billno = sd.Req AND IFNULL(dr.SAP_pr_billno,'')<>'' WHERE p.source_table='PurOrdRctDetail' AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctType'))='rc' ON DUPLICATE KEY UPDATE factory_id=VALUES(factory_id), rct_date=VALUES(rct_date), supp=VALUES(supp), sort_name=VALUES(sort_name), item_num=VALUES(item_num), item_name=VALUES(item_name), item_spec=VALUES(item_spec), um=VALUES(um), qty_ordered=VALUES(qty_ordered), qty_received=VALUES(qty_received), lot_serial=VALUES(lot_serial), location=VALUES(location), ord_nbr=VALUES(ord_nbr), ord_line=VALUES(ord_line), blanket_line=VALUES(blanket_line), pur_ord=VALUES(pur_ord), pur_line=VALUES(pur_line), sales_job=VALUES(sales_job), address1=VALUES(address1), req=VALUES(req), req_line=VALUES(req_line), dop_req=VALUES(dop_req), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP """, new SugarParameter("@BatchId", batchId), new SugarParameter("@Now", now)); return rows; } private async Task InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType) { await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_transform_run_log (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time) VALUES (0, @JobCode, 'S5采购收货单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime) """, new SugarParameter("@JobCode", JobCode), new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)), new SugarParameter("@BatchId", batchId), new SugarParameter("@StartTime", startedAt)); return await _db.Ado.GetLongAsync( "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1", new List { new("@BatchId", batchId) }); } private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, PurchaseReceiptMdpSyncResult result) { var finishedAt = DateTime.Now; await _db.Ado.ExecuteCommandAsync( """ UPDATE mdp_transform_run_log SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs, stage_rows=0, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP WHERE id=@Id """, new SugarParameter("@EndTime", finishedAt), new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds), new SugarParameter("@StandardRows", result.StdRows), new SugarParameter("@Id", runLogId)); } private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message) { var finishedAt = DateTime.Now; await _db.Ado.ExecuteCommandAsync( """ UPDATE mdp_transform_run_log SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs, error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP WHERE id=@Id """, new SugarParameter("@EndTime", finishedAt), new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds), new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message), new SugarParameter("@Id", runLogId)); } private static string NormalizeTriggerType(string? triggerType) => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant(); } /// 采购收货单 MDP 同步转换结果。 public sealed class PurchaseReceiptMdpSyncResult { public long RunLogId { get; set; } public string BatchId { get; set; } = string.Empty; public int StdRows { get; set; } } /// 采购收货双模式入站结果(stg 抽数 + std 转换)。 public sealed class PurchaseReceiptInboundResult { public string PullBatchId { get; set; } = string.Empty; public int RowsPulled { get; set; } public int RowsWrittenStg { get; set; } public string? NewCursor { get; set; } public string? PullMessage { get; set; } public string TransformBatchId { get; set; } = string.Empty; public int StdRows { get; set; } }