ProductionReceiptMdpSyncService.cs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  2. namespace Admin.NET.Plugin.AiDOP.MaterialWarehouse;
  3. /// <summary>
  4. /// S5 生产入库单 数据中台只读同步转换服务(DOP 内部,扁平单表,独立于 S5MdpSyncTransformService 的 KPI 管线)。
  5. ///
  6. /// 源:aidopdev NbrMaster(n) + NbrDetail(d),业务类型 Type='WOI'(生产入库);
  7. /// 扁平 LIST join:n.Domain=d.Domain AND n.Nbr=d.Nbr(一行=一条 NbrDetail 明细)。
  8. /// 维表 DepartmentMaster(p) / ItemMaster(i) / LocationMaster(lt,lf) 全 LEFT JOIN。
  9. /// LocationMaster 无 Domain 列(DOP 重建:domain_code/location/descr),按 tenant_id+location join。
  10. /// 链路(Phase 1):执行器抽 NbrMaster/NbrDetail → mdp_stg_production_receipt → mdp_std_production_receipt。
  11. /// 双模式入站:mdp_entity=S7_NBR_MASTER/S7_NBR_DETAIL → Writer → stg;标准层读 stg;维表仍直连本库。
  12. ///
  13. /// 约束:
  14. /// - 只读源/贴源,仅写 mdp_stg_production_receipt / mdp_std_production_receipt;绝不写 NbrMaster/NbrDetail。
  15. /// - 过滤 Type='WOI' AND IsActive=1。
  16. /// - WOI 为 0 行时转换成功完成、处理数为 0,不报错。
  17. /// </summary>
  18. public class ProductionReceiptMdpSyncService : ITransient
  19. {
  20. private const string JobCode = "S5_PRODUCTION_RECEIPT_MDP_SYNC";
  21. private const string InboundEntityCode = "S7_NBR_MASTER";
  22. private const string InboundDetailEntityCode = "S7_NBR_DETAIL";
  23. private readonly ISqlSugarClient _db;
  24. private readonly MdpSourcePullDispatcher _pullDispatcher;
  25. public ProductionReceiptMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher)
  26. {
  27. _db = db;
  28. _pullDispatcher = pullDispatcher;
  29. }
  30. /// <summary>全量:本地 DB 执行器灌 stg → 标准层(读 stg)。</summary>
  31. public async Task<ProductionReceiptMdpSyncResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
  32. {
  33. cancellationToken.ThrowIfCancellationRequested();
  34. await EnsureTablesAsync();
  35. await EnsureStgTableAsync();
  36. var now = DateTime.Now;
  37. var batchId = $"S5_PROD_RCPT_FULL_{now:yyyyMMddHHmmss}";
  38. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  39. var result = new ProductionReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  40. try
  41. {
  42. var pullCtx = new MdpPullContext
  43. {
  44. TenantId = 0,
  45. FullRefresh = true,
  46. TaskCode = "S7_PRODUCTION_RECEIPT_INBOUND",
  47. BatchId = $"{batchId}_PULL"
  48. };
  49. await PopulateStgAsync(pullCtx, cancellationToken);
  50. result.Rows = await TransformStandardAsync(batchId, now);
  51. await MarkRunSuccessAsync(runLogId, now, result);
  52. return result;
  53. }
  54. catch (Exception ex)
  55. {
  56. await MarkRunFailedAsync(runLogId, now, ex.Message);
  57. throw;
  58. }
  59. }
  60. /// <summary>
  61. /// S7 成品入库双模式入站:执行器抽主/明细 → stg,再跑 WOI 标准层(读 stg)。
  62. /// </summary>
  63. public async Task<ProductionReceiptInboundResult> RunInboundAsync(
  64. long tenantId = 0,
  65. bool fullRefresh = false,
  66. CancellationToken cancellationToken = default)
  67. {
  68. cancellationToken.ThrowIfCancellationRequested();
  69. await EnsureTablesAsync();
  70. await EnsureStgTableAsync();
  71. var now = DateTime.Now;
  72. var pullCtx = new MdpPullContext
  73. {
  74. TenantId = tenantId,
  75. FullRefresh = fullRefresh,
  76. TaskCode = "S7_PRODUCTION_RECEIPT_INBOUND",
  77. BatchId = $"S7_NBR_IN_{now:yyyyMMddHHmmss}"
  78. };
  79. var pull = await PopulateStgAsync(pullCtx, cancellationToken);
  80. var batchId = $"S7_NBR_STD_{now:yyyyMMddHHmmss}";
  81. var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
  82. var result = new ProductionReceiptMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  83. try
  84. {
  85. result.Rows = await TransformStandardAsync(batchId, now);
  86. await MarkRunSuccessAsync(runLogId, now, result);
  87. }
  88. catch (Exception ex)
  89. {
  90. await MarkRunFailedAsync(runLogId, now, ex.Message);
  91. throw;
  92. }
  93. return new ProductionReceiptInboundResult
  94. {
  95. PullBatchId = pullCtx.BatchId,
  96. RowsPulled = pull.RowsPulled,
  97. RowsWrittenStg = pull.RowsWritten,
  98. NewCursor = pull.NewCursor,
  99. TransformBatchId = result.BatchId,
  100. StdRows = result.Rows
  101. };
  102. }
  103. private async Task<(int RowsPulled, int RowsWritten, string? NewCursor)> PopulateStgAsync(
  104. MdpPullContext pullCtx, CancellationToken cancellationToken)
  105. {
  106. var master = await _pullDispatcher.PullByEntityCodeAsync(InboundEntityCode, pullCtx, cancellationToken);
  107. var detail = await _pullDispatcher.PullByEntityCodeAsync(InboundDetailEntityCode, pullCtx, cancellationToken);
  108. return (
  109. master.RowsPulled + detail.RowsPulled,
  110. master.RowsWritten + detail.RowsWritten,
  111. detail.NewCursor ?? master.NewCursor);
  112. }
  113. private async Task EnsureStgTableAsync()
  114. {
  115. await _db.Ado.ExecuteCommandAsync(
  116. """
  117. CREATE TABLE IF NOT EXISTS mdp_stg_production_receipt (
  118. id BIGINT PRIMARY KEY AUTO_INCREMENT,
  119. tenant_id BIGINT NOT NULL,
  120. source_system VARCHAR(50) NULL,
  121. source_table VARCHAR(200),
  122. source_row_id VARCHAR(200),
  123. source_biz_key VARCHAR(300) NULL,
  124. raw_data JSON,
  125. sync_batch_id VARCHAR(100),
  126. sync_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP,
  127. process_status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
  128. process_message VARCHAR(500) NULL,
  129. create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
  130. update_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  131. UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key),
  132. KEY idx_batch (sync_batch_id),
  133. KEY idx_src (source_table, source_row_id),
  134. KEY idx_tenant (tenant_id)
  135. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='S7成品入库执行器贴源层'
  136. """);
  137. }
  138. /// <summary>防御式建表(与 UpdateScripts/1.0.214.sql 同构,幂等)。</summary>
  139. private async Task EnsureTablesAsync()
  140. {
  141. await _db.Ado.ExecuteCommandAsync(
  142. """
  143. CREATE TABLE IF NOT EXISTS mdp_std_production_receipt (
  144. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  145. tenant_id BIGINT NOT NULL DEFAULT 0,
  146. factory_id BIGINT NULL DEFAULT 1,
  147. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  148. domain VARCHAR(80) NOT NULL,
  149. master_rec_id INT NULL,
  150. detail_rec_id INT NOT NULL,
  151. nbr VARCHAR(24) NULL,
  152. line SMALLINT NOT NULL DEFAULT 0,
  153. receipt_date DATETIME NULL,
  154. status VARCHAR(8) NULL,
  155. status_desc VARCHAR(20) NULL,
  156. remark VARCHAR(200) NULL,
  157. prod_line VARCHAR(8) NULL,
  158. work_ord VARCHAR(64) NULL,
  159. erp_work_ord VARCHAR(60) NULL,
  160. department VARCHAR(8) NULL,
  161. department_desc VARCHAR(255) NULL,
  162. applicant_name VARCHAR(12) NULL,
  163. item_num VARCHAR(24) NULL,
  164. item_name VARCHAR(200) NULL,
  165. item_spec VARCHAR(200) NULL,
  166. um VARCHAR(8) NULL,
  167. location_to VARCHAR(8) NULL,
  168. location_to_desc VARCHAR(255) NULL,
  169. lot_serial VARCHAR(120) NULL,
  170. qty_rec DECIMAL(18,5) NULL DEFAULT 0,
  171. qty_to DECIMAL(18,5) NULL DEFAULT 0,
  172. location_from VARCHAR(8) NULL,
  173. location_from_desc VARCHAR(255) NULL,
  174. ord_nbr VARCHAR(48) NULL,
  175. source_biz_key VARCHAR(200) NULL,
  176. sync_batch_id VARCHAR(100) NOT NULL,
  177. sync_time DATETIME NOT NULL,
  178. create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
  179. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  180. UNIQUE KEY uk_mdp_std_prod_receipt (tenant_id, domain, detail_rec_id),
  181. KEY idx_mdp_std_prod_rcpt_nbr (tenant_id, nbr),
  182. KEY idx_mdp_std_prod_rcpt_date (tenant_id, receipt_date),
  183. KEY idx_mdp_std_prod_rcpt_workord (tenant_id, work_ord),
  184. KEY idx_mdp_std_prod_rcpt_erp (tenant_id, erp_work_ord),
  185. KEY idx_mdp_std_prod_rcpt_lot (tenant_id, lot_serial),
  186. KEY idx_mdp_std_prod_rcpt_item (tenant_id, item_num),
  187. KEY idx_mdp_std_prod_rcpt_locto (tenant_id, location_to)
  188. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S5生产入库单标准层(扁平,一行=一条明细)'
  189. """);
  190. }
  191. /// <summary>
  192. /// 标准化:stg(NbrMaster/NbrDetail) + 维表(本库) → mdp_std_production_receipt。
  193. /// </summary>
  194. private async Task<int> TransformStandardAsync(string batchId, DateTime now)
  195. {
  196. var rows = await _db.Ado.GetIntAsync(
  197. """
  198. SELECT COUNT(1)
  199. FROM mdp_stg_production_receipt n
  200. INNER JOIN mdp_stg_production_receipt d
  201. ON d.source_table='NbrDetail'
  202. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  203. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr'))
  204. WHERE n.source_table='NbrMaster'
  205. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI'
  206. AND CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive')) AS SIGNED)=1
  207. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL
  208. """);
  209. await _db.Ado.ExecuteCommandAsync(
  210. """
  211. INSERT INTO mdp_std_production_receipt
  212. (tenant_id, factory_id, source_system, domain, master_rec_id, detail_rec_id, nbr, line,
  213. receipt_date, status, status_desc, remark, prod_line, work_ord, erp_work_ord,
  214. department, department_desc, applicant_name, item_num, item_name, item_spec, um,
  215. location_to, location_to_desc, lot_serial, qty_rec, qty_to, location_from, location_from_desc,
  216. ord_nbr, source_biz_key, sync_batch_id, sync_time)
  217. SELECT
  218. IFNULL(n.tenant_id, 0), 1, IFNULL(NULLIF(n.source_system,''), 'AIDOP'),
  219. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')),
  220. CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.RecID')) AS SIGNED),
  221. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS SIGNED),
  222. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')),
  223. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS SIGNED),
  224. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Date')),'null'),''), '%Y-%m-%d %H:%i:%s'),
  225. UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
  226. UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
  227. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Remark')),
  228. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ProdLine')),
  229. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.WorkOrd')),
  230. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Address')),
  231. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),
  232. TRIM(CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),''), ' ', IFNULL(p.Descr, ''))),
  233. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Name')),
  234. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')), i.Descr, i.Descr1, i.UM,
  235. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo')), lt.descr,
  236. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LotSerial')) AS CHAR),
  237. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyRec')) AS DECIMAL(18,5)),
  238. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyTo')) AS DECIMAL(18,5)),
  239. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom')), lf.descr,
  240. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')),
  241. CONCAT(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')), '#',
  242. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID'))),
  243. @BatchId, @Now
  244. FROM mdp_stg_production_receipt n
  245. INNER JOIN mdp_stg_production_receipt d
  246. ON d.source_table='NbrDetail'
  247. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  248. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr'))
  249. LEFT JOIN DepartmentMaster p
  250. ON p.Domain = JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain'))
  251. AND p.Department = JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department'))
  252. LEFT JOIN ItemMaster i
  253. ON i.Domain = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  254. AND i.ItemNum = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum'))
  255. LEFT JOIN LocationMaster lt
  256. ON lt.tenant_id = IFNULL(d.tenant_id, 0)
  257. AND lt.location = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo'))
  258. LEFT JOIN LocationMaster lf
  259. ON lf.tenant_id = IFNULL(d.tenant_id, 0)
  260. AND lf.location = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom'))
  261. WHERE n.source_table='NbrMaster'
  262. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI'
  263. AND CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive')) AS SIGNED)=1
  264. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL
  265. ON DUPLICATE KEY UPDATE
  266. factory_id=VALUES(factory_id), master_rec_id=VALUES(master_rec_id), nbr=VALUES(nbr), line=VALUES(line),
  267. receipt_date=VALUES(receipt_date), status=VALUES(status), status_desc=VALUES(status_desc), remark=VALUES(remark),
  268. prod_line=VALUES(prod_line), work_ord=VALUES(work_ord), erp_work_ord=VALUES(erp_work_ord),
  269. department=VALUES(department), department_desc=VALUES(department_desc), applicant_name=VALUES(applicant_name),
  270. item_num=VALUES(item_num), item_name=VALUES(item_name), item_spec=VALUES(item_spec), um=VALUES(um),
  271. location_to=VALUES(location_to), location_to_desc=VALUES(location_to_desc), lot_serial=VALUES(lot_serial),
  272. qty_rec=VALUES(qty_rec), qty_to=VALUES(qty_to), location_from=VALUES(location_from),
  273. location_from_desc=VALUES(location_from_desc), ord_nbr=VALUES(ord_nbr),
  274. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  275. """,
  276. new SugarParameter("@BatchId", batchId),
  277. new SugarParameter("@Now", now));
  278. return rows;
  279. }
  280. private async Task<long> InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  281. {
  282. await _db.Ado.ExecuteCommandAsync(
  283. """
  284. INSERT INTO mdp_transform_run_log
  285. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  286. VALUES (0, @JobCode, 'S5生产入库单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  287. """,
  288. new SugarParameter("@JobCode", JobCode),
  289. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  290. new SugarParameter("@BatchId", batchId),
  291. new SugarParameter("@StartTime", startedAt));
  292. return await _db.Ado.GetLongAsync(
  293. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  294. new List<SugarParameter> { new("@BatchId", batchId) });
  295. }
  296. private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, ProductionReceiptMdpSyncResult result)
  297. {
  298. var finishedAt = DateTime.Now;
  299. await _db.Ado.ExecuteCommandAsync(
  300. """
  301. UPDATE mdp_transform_run_log
  302. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  303. stage_rows=0, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP
  304. WHERE id=@Id
  305. """,
  306. new SugarParameter("@EndTime", finishedAt),
  307. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  308. new SugarParameter("@StandardRows", result.Rows),
  309. new SugarParameter("@Id", runLogId));
  310. }
  311. private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message)
  312. {
  313. var finishedAt = DateTime.Now;
  314. await _db.Ado.ExecuteCommandAsync(
  315. """
  316. UPDATE mdp_transform_run_log
  317. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  318. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  319. WHERE id=@Id
  320. """,
  321. new SugarParameter("@EndTime", finishedAt),
  322. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  323. new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message),
  324. new SugarParameter("@Id", runLogId));
  325. }
  326. private static string NormalizeTriggerType(string? triggerType)
  327. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  328. }
  329. /// <summary>生产入库单 MDP 同步转换结果。</summary>
  330. public sealed class ProductionReceiptMdpSyncResult
  331. {
  332. public long RunLogId { get; set; }
  333. public string BatchId { get; set; } = string.Empty;
  334. public int Rows { get; set; }
  335. }
  336. /// <summary>S7 成品入库双模式入站结果。</summary>
  337. public sealed class ProductionReceiptInboundResult
  338. {
  339. public string PullBatchId { get; set; } = string.Empty;
  340. public int RowsPulled { get; set; }
  341. public int RowsWrittenStg { get; set; }
  342. public string? NewCursor { get; set; }
  343. public string TransformBatchId { get; set; } = string.Empty;
  344. public int StdRows { get; set; }
  345. }