PurchaseReceiptMdpSyncService.cs 23 KB

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