PurchaseReceiptMdpSyncService.cs 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454
  1. using Admin.NET.Plugin.AiDOP.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  3. namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
  4. /// <summary>
  5. /// S5 采购收货单 数据中台只读同步转换服务(DOP 内部,方案 1)。
  6. ///
  7. /// 源:aidopdev PurOrdRctDetail(p) + PurOrdRctMaster(d),业务类型 RctType='rc'(采购收货);
  8. /// 头明细关联 p.Domain=d.Domain AND p.Receiver=d.Receiver;
  9. /// 维表 ItemMaster(i) / SuppMaster(s) / ConsigneeAddressMaster(a) / PurOrdDetail(pd,sd) / srm_pr_main(dr) 全 LEFT JOIN。
  10. /// 链路(Phase 1):执行器抽主/明细 → mdp_stg_purchase_receipt → mdp_std_purchase_receipt(typed)。
  11. /// 双模式入站:mdp_entity=S5_PURCHASE_RECEIPT_DETAIL/MASTER(DB 或 API)→ MdpStagingWriter → stg,
  12. /// 标准层读 stg.raw_data;维表 ItemMaster/SuppMaster 等仍直连本库(本地口径补全)。
  13. ///
  14. /// 约束:
  15. /// - 只读源/贴源,仅写 mdp_stg_purchase_receipt / mdp_std_purchase_receipt;绝不写源表;不读/不改 S3 mdp_stg_receipt、S4 ado_s4_receipt。
  16. /// - 过滤 RctType='rc';租户隔离在读 API 侧按 tenant_id。
  17. /// - rc 为 0 行时成功完成、处理数为 0,不报错。
  18. /// </summary>
  19. public class PurchaseReceiptMdpSyncService : ITransient
  20. {
  21. private const string JobCode = "S5_PURCHASE_RECEIPT_MDP_SYNC";
  22. private const string InboundEntityCode = "S5_PURCHASE_RECEIPT_DETAIL";
  23. private const string InboundMasterEntityCode = "S5_PURCHASE_RECEIPT_MASTER";
  24. // 双源(dopdemorq SQL Server)Phase 1:默认第二源与其入站实体。实体默认 status=0,未启用前 RunSourceSwitchFullAsync 会因找不到启用实体而抛错(就位不启用)。
  25. private const string SqlServerSourceCode = "DOPDEMORQ_SQLSERVER";
  26. private const string SqlServerMasterEntityCode = "S5_PURCHASE_RECEIPT_MASTER_SQLSERVER";
  27. private const string SqlServerDetailEntityCode = "S5_PURCHASE_RECEIPT_DETAIL_SQLSERVER";
  28. private readonly ISqlSugarClient _db;
  29. private readonly MdpSourcePullDispatcher _pullDispatcher;
  30. public PurchaseReceiptMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher)
  31. {
  32. _db = db;
  33. _pullDispatcher = pullDispatcher;
  34. }
  35. /// <summary>全量:本地 DB 执行器灌 stg → 标准层(读 stg)。</summary>
  36. public async Task<PurchaseReceiptMdpSyncResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
  37. {
  38. cancellationToken.ThrowIfCancellationRequested();
  39. await EnsureTablesAsync();
  40. await EnsureStgTableAsync();
  41. var now = DateTime.Now;
  42. var batchId = $"S5_PUR_RCT_FULL_{now:yyyyMMddHHmmss}";
  43. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  44. var result = new PurchaseReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  45. try
  46. {
  47. var pullCtx = new MdpPullContext
  48. {
  49. TenantId = 0,
  50. FullRefresh = true,
  51. TaskCode = "S5_PURCHASE_RECEIPT_INBOUND",
  52. BatchId = $"{batchId}_PULL"
  53. };
  54. await PopulateStgAsync(pullCtx, cancellationToken);
  55. result.StdRows = await TransformStandardAsync(batchId, now);
  56. await MarkRunSuccessAsync(runLogId, now, result);
  57. return result;
  58. }
  59. catch (Exception ex)
  60. {
  61. await MarkRunFailedAsync(runLogId, now, ex.Message);
  62. throw;
  63. }
  64. }
  65. /// <summary>
  66. /// 双模式入站:执行器抽主/明细落 stg,再跑标准层(读 stg)。
  67. /// </summary>
  68. public async Task<PurchaseReceiptInboundResult> RunInboundAsync(
  69. long tenantId = 0,
  70. bool fullRefresh = false,
  71. CancellationToken cancellationToken = default)
  72. {
  73. cancellationToken.ThrowIfCancellationRequested();
  74. await EnsureTablesAsync();
  75. await EnsureStgTableAsync();
  76. var now = DateTime.Now;
  77. var pullCtx = new MdpPullContext
  78. {
  79. TenantId = tenantId,
  80. FullRefresh = fullRefresh,
  81. TaskCode = "S5_PURCHASE_RECEIPT_INBOUND",
  82. BatchId = $"S5_PUR_RCT_IN_{now:yyyyMMddHHmmss}"
  83. };
  84. var pull = await PopulateStgAsync(pullCtx, cancellationToken);
  85. var batchId = $"S5_PUR_RCT_STD_{now:yyyyMMddHHmmss}";
  86. var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
  87. var result = new PurchaseReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  88. try
  89. {
  90. result.StdRows = await TransformStandardAsync(batchId, now);
  91. await MarkRunSuccessAsync(runLogId, now, result);
  92. }
  93. catch (Exception ex)
  94. {
  95. await MarkRunFailedAsync(runLogId, now, ex.Message);
  96. throw;
  97. }
  98. return new PurchaseReceiptInboundResult
  99. {
  100. PullBatchId = pullCtx.BatchId,
  101. RowsPulled = pull.RowsPulled,
  102. RowsWrittenStg = pull.RowsWritten,
  103. NewCursor = pull.NewCursor,
  104. PullMessage = pull.Message,
  105. TransformBatchId = result.BatchId,
  106. StdRows = result.StdRows
  107. };
  108. }
  109. /// <summary>
  110. /// 双源切换 FULL Replace(Phase 1):从指定 DB 源(默认 DOPDEMORQ_SQLSERVER)PullAll 全量灌 stg,
  111. /// 成功后在单事务内以「仅当前源」结果 FULL 重建 mdp_std_purchase_receipt(按 tenant 精确隔离),消除旧源独有业务键残留。
  112. ///
  113. /// 与既有 RunFull/RunInbound 的区别:① 用 PullAllByEntityCodeAsync 抽尽(非单批 PullByEntityCode);
  114. /// ② FullRefresh=true 显式全量 bootstrap;③ transform 仅当前 source_system;④ std DELETE+INSERT 同事务,Pull 成功后才进入。
  115. ///
  116. /// 切源语义为「部署级单活源」。SQLSERVER 实体默认 status=0,未启用前本方法运行时会因 PullByEntityCode 找不到启用实体抛错——符合 Phase 1「就位不启用」。
  117. /// dopdemorq 当前六表为空:本方法 FULL Replace 的理论结果是把当前 tenant 对应 std 置空;Phase 1 不实际执行 destructive 切换。
  118. /// </summary>
  119. public async Task<PurchaseReceiptInboundResult> RunSourceSwitchFullAsync(
  120. string sourceCode = SqlServerSourceCode,
  121. string masterEntityCode = SqlServerMasterEntityCode,
  122. string detailEntityCode = SqlServerDetailEntityCode,
  123. long tenantId = 0,
  124. CancellationToken cancellationToken = default)
  125. {
  126. cancellationToken.ThrowIfCancellationRequested();
  127. await EnsureTablesAsync();
  128. await EnsureStgTableAsync();
  129. var now = DateTime.Now;
  130. // 1) PullAll 抽尽 + fullRefresh=true(不用单批 PullByEntityCode)。任一 Pull 失败会抛异常,std 未动。
  131. var pullCtx = new MdpPullContext
  132. {
  133. TenantId = tenantId,
  134. FullRefresh = true,
  135. TaskCode = "S5_PURCHASE_RECEIPT_INBOUND",
  136. BatchId = $"S5_PUR_RCT_SW_{now:yyyyMMddHHmmss}"
  137. };
  138. var master = await _pullDispatcher.PullAllByEntityCodeAsync(masterEntityCode, pullCtx, cancellationToken);
  139. var detail = await _pullDispatcher.PullAllByEntityCodeAsync(detailEntityCode, pullCtx, cancellationToken);
  140. // 2) FULL Replace:事务内 DELETE 当前 tenant 的 std,再 INSERT 仅当前 source_system 的结果。
  141. var batchId = $"S5_PUR_RCT_SWSTD_{now:yyyyMMddHHmmss}";
  142. var runLogId = await InsertRunLogAsync(batchId, now, "SOURCE_SWITCH");
  143. var result = new PurchaseReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  144. try
  145. {
  146. result.StdRows = await MdpStdFullReplace.ReplaceAsync(
  147. _db, "mdp_std_purchase_receipt", tenantId, extraWhere: null,
  148. insertScopedAsync: () => TransformStandardAsync(batchId, now, sourceCode),
  149. cancellationToken);
  150. await MarkRunSuccessAsync(runLogId, now, result);
  151. }
  152. catch (Exception ex)
  153. {
  154. await MarkRunFailedAsync(runLogId, now, ex.Message);
  155. throw;
  156. }
  157. return new PurchaseReceiptInboundResult
  158. {
  159. PullBatchId = pullCtx.BatchId,
  160. RowsPulled = master.RowsPulled + detail.RowsPulled,
  161. RowsWrittenStg = master.RowsWritten + detail.RowsWritten,
  162. NewCursor = detail.NewCursor ?? master.NewCursor,
  163. PullMessage = $"{master.Message}; {detail.Message}",
  164. TransformBatchId = batchId,
  165. StdRows = result.StdRows
  166. };
  167. }
  168. private async Task<(int RowsPulled, int RowsWritten, string? NewCursor, string? Message)> PopulateStgAsync(
  169. MdpPullContext pullCtx, CancellationToken cancellationToken)
  170. {
  171. var master = await _pullDispatcher.PullByEntityCodeAsync(InboundMasterEntityCode, pullCtx, cancellationToken);
  172. var detail = await _pullDispatcher.PullByEntityCodeAsync(InboundEntityCode, pullCtx, cancellationToken);
  173. return (
  174. master.RowsPulled + detail.RowsPulled,
  175. master.RowsWritten + detail.RowsWritten,
  176. detail.NewCursor ?? master.NewCursor,
  177. $"{master.Message}; {detail.Message}");
  178. }
  179. private async Task EnsureStgTableAsync()
  180. {
  181. await _db.Ado.ExecuteCommandAsync(
  182. """
  183. CREATE TABLE IF NOT EXISTS mdp_stg_purchase_receipt (
  184. id BIGINT PRIMARY KEY AUTO_INCREMENT,
  185. tenant_id BIGINT NOT NULL,
  186. source_system VARCHAR(50) NULL,
  187. source_table VARCHAR(200),
  188. source_row_id VARCHAR(200),
  189. source_biz_key VARCHAR(300) NULL,
  190. raw_data JSON,
  191. sync_batch_id VARCHAR(100),
  192. sync_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP,
  193. process_status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
  194. process_message VARCHAR(500) NULL,
  195. create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
  196. update_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  197. UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key),
  198. KEY idx_batch (sync_batch_id),
  199. KEY idx_src (source_table, source_row_id),
  200. KEY idx_tenant (tenant_id)
  201. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='S5采购收货贴源层'
  202. """);
  203. }
  204. /// <summary>防御式建表(与 UpdateScripts WIP-S5PR.sql / 正式 1.0.&lt;n&gt;.sql 同构,幂等)。</summary>
  205. private async Task EnsureTablesAsync()
  206. {
  207. await _db.Ado.ExecuteCommandAsync(
  208. """
  209. CREATE TABLE IF NOT EXISTS mdp_std_purchase_receipt (
  210. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  211. tenant_id BIGINT NOT NULL DEFAULT 0,
  212. factory_id BIGINT NULL DEFAULT 1,
  213. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  214. domain VARCHAR(24) NOT NULL,
  215. receiver VARCHAR(24) NOT NULL,
  216. line SMALLINT NOT NULL DEFAULT 0,
  217. rct_date DATETIME NULL,
  218. supp VARCHAR(20) NULL,
  219. sort_name VARCHAR(255) NULL,
  220. item_num VARCHAR(60) NULL,
  221. item_name VARCHAR(200) NULL,
  222. item_spec VARCHAR(200) NULL,
  223. um VARCHAR(8) NULL,
  224. qty_ordered DECIMAL(18,6) NULL DEFAULT 0,
  225. qty_received DECIMAL(18,6) NULL DEFAULT 0,
  226. lot_serial VARCHAR(120) NULL,
  227. location VARCHAR(8) NULL,
  228. ord_nbr VARCHAR(24) NULL,
  229. ord_line SMALLINT NULL,
  230. blanket_line INT NULL,
  231. pur_ord VARCHAR(24) NULL,
  232. pur_line SMALLINT NULL,
  233. sales_job VARCHAR(200) NULL,
  234. address1 VARCHAR(200) NULL,
  235. req VARCHAR(20) NULL,
  236. req_line INT NULL,
  237. dop_req VARCHAR(255) NULL,
  238. source_biz_key VARCHAR(200) NULL,
  239. sync_batch_id VARCHAR(100) NOT NULL,
  240. sync_time DATETIME NOT NULL,
  241. create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
  242. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  243. UNIQUE KEY uk_mdp_std_pur_rct (tenant_id, domain, receiver, line),
  244. KEY idx_mdp_std_pur_rct_date (tenant_id, rct_date),
  245. KEY idx_mdp_std_pur_rct_item (tenant_id, item_num),
  246. KEY idx_mdp_std_pur_rct_supp (tenant_id, supp),
  247. KEY idx_mdp_std_pur_rct_purord (tenant_id, pur_ord),
  248. KEY idx_mdp_std_pur_rct_salesjob (tenant_id, sales_job),
  249. KEY idx_mdp_std_pur_rct_req (tenant_id, req),
  250. KEY idx_mdp_std_pur_rct_dopreq (tenant_id, dop_req)
  251. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S5采购收货单标准层'
  252. """);
  253. }
  254. /// <summary>
  255. /// 标准化:stg(PurOrdRctDetail/Master) + 维表(本库) → mdp_std_purchase_receipt。
  256. /// </summary>
  257. private async Task<int> TransformStandardAsync(string batchId, DateTime now, string? sourceSystem = null)
  258. {
  259. // 双源:sourceSystem 非空时仅统计/转换当前 source(切源 FULL Replace 用);为空时保持既有全源行为不变。
  260. var srcClause = sourceSystem == null ? "" : " AND p.source_system=@Src AND d.source_system=@Src";
  261. var countPars = new List<SugarParameter>();
  262. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  263. var rows = await _db.Ado.GetIntAsync(
  264. $"""
  265. SELECT COUNT(1)
  266. FROM mdp_stg_purchase_receipt p
  267. INNER JOIN mdp_stg_purchase_receipt d
  268. ON d.source_table='PurOrdRctMaster'
  269. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  270. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Receiver'))
  271. WHERE p.source_table='PurOrdRctDetail'
  272. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctType'))='rc'{srcClause}
  273. """, countPars);
  274. var insertSql =
  275. $"""
  276. INSERT INTO mdp_std_purchase_receipt
  277. (tenant_id, factory_id, source_system, domain, receiver, line, rct_date, supp, sort_name,
  278. item_num, item_name, item_spec, um, qty_ordered, qty_received, lot_serial, location,
  279. ord_nbr, ord_line, blanket_line, pur_ord, pur_line, sales_job, address1,
  280. req, req_line, dop_req, source_biz_key, sync_batch_id, sync_time)
  281. SELECT
  282. IFNULL(p.tenant_id, 0), 1, IFNULL(NULLIF(p.source_system,''), 'AIDOP'),
  283. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')),
  284. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')),
  285. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Line')) AS SIGNED),
  286. STR_TO_DATE(REPLACE(LEFT(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RctDate')),'null'),''), 19), 'T', ' '), '%Y-%m-%d %H:%i:%s'),
  287. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Supp')),
  288. TRIM(CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Supp')),''), ' ', IFNULL(s.SortName, ''))),
  289. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.ItemNum')), i.Descr, i.Descr1,
  290. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.UM')),
  291. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.QtyOrded')) AS DECIMAL(18,6)),
  292. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.QtyReceived')) AS DECIMAL(18,6)),
  293. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.LotSerial')),
  294. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Location')),
  295. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdNbr')),
  296. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdLine')) AS SIGNED),
  297. pd.BlanketLine, pd.PurOrd, pd.Line, pd.SalesJob, a.Address1,
  298. (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO' THEN '' ELSE sd.Req END),
  299. sd.ReqLine,
  300. (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO' THEN sd.Req ELSE dr.pr_billno END),
  301. IFNULL(NULLIF(p.source_biz_key,''), CONCAT(
  302. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')), '#',
  303. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')), '#',
  304. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Line')))),
  305. @BatchId, @Now
  306. FROM mdp_stg_purchase_receipt p
  307. INNER JOIN mdp_stg_purchase_receipt d
  308. ON d.source_table='PurOrdRctMaster'
  309. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  310. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Receiver'))
  311. LEFT JOIN ItemMaster i
  312. ON JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = i.Domain
  313. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.ItemNum')) = i.ItemNum
  314. LEFT JOIN SuppMaster s
  315. ON JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain')) = s.Domain
  316. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Supp')) = s.Supp
  317. LEFT JOIN ConsigneeAddressMaster a
  318. ON JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = a.Domain
  319. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Supp')) = a.Address
  320. AND a.Typed = 'Supp'
  321. LEFT JOIN PurOrdDetail pd
  322. ON JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = pd.Domain
  323. AND (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO'
  324. AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctNbr')),'')<>''
  325. THEN JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctNbr'))
  326. ELSE JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdNbr')) END)
  327. = (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO'
  328. AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctNbr')),'')<>''
  329. THEN pd.Contract ELSE pd.PurOrd END)
  330. AND (CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdType'))='DO'
  331. AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctNbr')),'')<>''
  332. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.BlanketLine')) AS SIGNED)
  333. ELSE CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdLine')) AS SIGNED) END) = pd.Line
  334. LEFT JOIN PurOrdDetail sd
  335. ON JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = sd.Domain
  336. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdNbr')) = sd.PurOrd
  337. AND CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdLine')) AS SIGNED) = sd.Line
  338. LEFT JOIN srm_pr_main dr
  339. ON CAST(dr.factory_id AS CHAR) = sd.Domain
  340. AND dr.SAP_pr_billno = sd.Req
  341. AND IFNULL(dr.SAP_pr_billno,'')<>''
  342. WHERE p.source_table='PurOrdRctDetail'
  343. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctType'))='rc'{srcClause}
  344. ON DUPLICATE KEY UPDATE
  345. factory_id=VALUES(factory_id), rct_date=VALUES(rct_date), supp=VALUES(supp), sort_name=VALUES(sort_name),
  346. item_num=VALUES(item_num), item_name=VALUES(item_name), item_spec=VALUES(item_spec), um=VALUES(um),
  347. qty_ordered=VALUES(qty_ordered), qty_received=VALUES(qty_received), lot_serial=VALUES(lot_serial), location=VALUES(location),
  348. ord_nbr=VALUES(ord_nbr), ord_line=VALUES(ord_line), blanket_line=VALUES(blanket_line),
  349. pur_ord=VALUES(pur_ord), pur_line=VALUES(pur_line), sales_job=VALUES(sales_job), address1=VALUES(address1),
  350. req=VALUES(req), req_line=VALUES(req_line), dop_req=VALUES(dop_req),
  351. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  352. """;
  353. var insPars = new List<SugarParameter>
  354. {
  355. new("@BatchId", batchId),
  356. new("@Now", now)
  357. };
  358. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  359. await _db.Ado.ExecuteCommandAsync(insertSql, insPars);
  360. return rows;
  361. }
  362. private async Task<long> InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  363. {
  364. await _db.Ado.ExecuteCommandAsync(
  365. """
  366. INSERT INTO mdp_transform_run_log
  367. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  368. VALUES (0, @JobCode, 'S5采购收货单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  369. """,
  370. new SugarParameter("@JobCode", JobCode),
  371. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  372. new SugarParameter("@BatchId", batchId),
  373. new SugarParameter("@StartTime", startedAt));
  374. return await _db.Ado.GetLongAsync(
  375. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  376. new List<SugarParameter> { new("@BatchId", batchId) });
  377. }
  378. private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, PurchaseReceiptMdpSyncResult result)
  379. {
  380. var finishedAt = DateTime.Now;
  381. await _db.Ado.ExecuteCommandAsync(
  382. """
  383. UPDATE mdp_transform_run_log
  384. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  385. stage_rows=0, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP
  386. WHERE id=@Id
  387. """,
  388. new SugarParameter("@EndTime", finishedAt),
  389. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  390. new SugarParameter("@StandardRows", result.StdRows),
  391. new SugarParameter("@Id", runLogId));
  392. }
  393. private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message)
  394. {
  395. var finishedAt = DateTime.Now;
  396. await _db.Ado.ExecuteCommandAsync(
  397. """
  398. UPDATE mdp_transform_run_log
  399. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  400. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  401. WHERE id=@Id
  402. """,
  403. new SugarParameter("@EndTime", finishedAt),
  404. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  405. new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message),
  406. new SugarParameter("@Id", runLogId));
  407. }
  408. private static string NormalizeTriggerType(string? triggerType)
  409. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  410. }
  411. /// <summary>采购收货单 MDP 同步转换结果。</summary>
  412. public sealed class PurchaseReceiptMdpSyncResult
  413. {
  414. public long RunLogId { get; set; }
  415. public string BatchId { get; set; } = string.Empty;
  416. public int StdRows { get; set; }
  417. }
  418. /// <summary>采购收货双模式入站结果(stg 抽数 + std 转换)。</summary>
  419. public sealed class PurchaseReceiptInboundResult
  420. {
  421. public string PullBatchId { get; set; } = string.Empty;
  422. public int RowsPulled { get; set; }
  423. public int RowsWrittenStg { get; set; }
  424. public string? NewCursor { get; set; }
  425. public string? PullMessage { get; set; }
  426. public string TransformBatchId { get; set; } = string.Empty;
  427. public int StdRows { get; set; }
  428. }