ProductionIssueMdpSyncService.cs 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397
  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 内部,头级单表)。
  6. ///
  7. /// B-① 双源通用管线(2026-07-29 由 bespoke 直读 NbrMaster 迁入):
  8. /// 源(NbrMaster) --MdpDbPullExecutor--> mdp_stg_production_issue(raw_data JSON) --transform--> mdp_std_production_issue。
  9. /// 入站实体:S5_PRODUCTION_ISSUE_MASTER(源=AIDOPDEV_MYSQL) / S5_PRODUCTION_ISSUE_MASTER_SQLSERVER(源=DOPDEMORQ_SQLSERVER)。
  10. /// 维表 DepartmentMaster 仍本库 LEFT JOIN(与 mode-A 采购/生产收货同构;department_desc 富化,部门 code 由事实表保留)。
  11. ///
  12. /// 约束:
  13. /// - 只读源/贴源,仅写 mdp_stg_production_issue / mdp_std_production_issue;绝不写 NbrMaster;不碰 WorkOrderPickBillService 写路径。
  14. /// - 头过滤 Type='SM' AND IsActive(领料无 IsReturn 语义);SM=领料单业务类型(Type SM/WOI/WOD/CA 互斥)。
  15. /// - 跨源类型兼容经 MdpJsonSql(datetime ISO-T / bit true-false / JSON null 归一),不改 MDP 核心。
  16. /// - 切源单 active source + FULL Replace(决策1A);不做 merge、不做回写/outbox(决策2A)。
  17. /// - SM 为 0 行时转换成功完成、处理数为 0,不报错。
  18. /// </summary>
  19. public class ProductionIssueMdpSyncService : ITransient
  20. {
  21. private const string JobCode = "S5_PRODUCTION_ISSUE_MDP_SYNC";
  22. private const string InboundEntityCode = "S5_PRODUCTION_ISSUE_MASTER";
  23. // 双源(dopdemorq SQL Server):第二源与其入站实体。实体默认 status=0(就位不启用)。
  24. private const string SqlServerSourceCode = "DOPDEMORQ_SQLSERVER";
  25. private const string SqlServerEntityCode = "S5_PRODUCTION_ISSUE_MASTER_SQLSERVER";
  26. private readonly ISqlSugarClient _db;
  27. private readonly MdpSourcePullDispatcher _pullDispatcher;
  28. public ProductionIssueMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher)
  29. {
  30. _db = db;
  31. _pullDispatcher = pullDispatcher;
  32. }
  33. /// <summary>全量:本库 NbrMaster 执行器灌 stg → 标准层(读 stg)。</summary>
  34. public async Task<ProductionIssueMdpSyncResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
  35. {
  36. cancellationToken.ThrowIfCancellationRequested();
  37. await EnsureTablesAsync();
  38. await EnsureStgTableAsync();
  39. var now = DateTime.Now;
  40. var batchId = $"S5_PROD_ISSUE_FULL_{now:yyyyMMddHHmmss}";
  41. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  42. var result = new ProductionIssueMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  43. try
  44. {
  45. var pullCtx = new MdpPullContext
  46. {
  47. TenantId = 0,
  48. FullRefresh = true,
  49. TaskCode = "S5_PRODUCTION_ISSUE_INBOUND",
  50. BatchId = $"{batchId}_PULL"
  51. };
  52. await PopulateStgAsync(InboundEntityCode, pullCtx, cancellationToken);
  53. result.HeadRows = await TransformHeadStandardAsync(batchId, now);
  54. await MarkRunSuccessAsync(runLogId, now, result);
  55. return result;
  56. }
  57. catch (Exception ex)
  58. {
  59. await MarkRunFailedAsync(runLogId, now, ex.Message);
  60. throw;
  61. }
  62. }
  63. /// <summary>双模式入站:执行器抽头落 stg,再跑标准层(读 stg)。</summary>
  64. public async Task<ProductionIssueInboundResult> RunInboundAsync(
  65. long tenantId = 0, bool fullRefresh = false, CancellationToken cancellationToken = default)
  66. {
  67. cancellationToken.ThrowIfCancellationRequested();
  68. await EnsureTablesAsync();
  69. await EnsureStgTableAsync();
  70. var now = DateTime.Now;
  71. var pullCtx = new MdpPullContext
  72. {
  73. TenantId = tenantId,
  74. FullRefresh = fullRefresh,
  75. TaskCode = "S5_PRODUCTION_ISSUE_INBOUND",
  76. BatchId = $"S5_PROD_ISSUE_IN_{now:yyyyMMddHHmmss}"
  77. };
  78. var pull = await PopulateStgAsync(InboundEntityCode, pullCtx, cancellationToken);
  79. var batchId = $"S5_PROD_ISSUE_STD_{now:yyyyMMddHHmmss}";
  80. var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
  81. var result = new ProductionIssueMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  82. try
  83. {
  84. result.HeadRows = await TransformHeadStandardAsync(batchId, now);
  85. await MarkRunSuccessAsync(runLogId, now, result);
  86. }
  87. catch (Exception ex)
  88. {
  89. await MarkRunFailedAsync(runLogId, now, ex.Message);
  90. throw;
  91. }
  92. return new ProductionIssueInboundResult
  93. {
  94. PullBatchId = pullCtx.BatchId,
  95. RowsPulled = pull.RowsPulled,
  96. RowsWrittenStg = pull.RowsWritten,
  97. NewCursor = pull.NewCursor,
  98. TransformBatchId = result.BatchId,
  99. StdRows = result.HeadRows
  100. };
  101. }
  102. /// <summary>
  103. /// 双源切换 FULL Replace(Phase 1):从指定 DB 源(默认 DOPDEMORQ_SQLSERVER)PullAll 全量灌 stg,
  104. /// 成功后在单事务内以「仅当前源」结果 FULL 重建 mdp_std_production_issue(按 tenant 精确隔离),消除旧源独有业务键残留。
  105. /// PullAllByEntityCodeAsync 抽尽 + FullRefresh=true;transform 仅当前 source_system;Pull 成功后才进入 destructive replace。
  106. /// SQLSERVER 实体默认 status=0(就位不启用);本批 FULL-only,不启自动增量。
  107. /// </summary>
  108. public async Task<ProductionIssueInboundResult> RunSourceSwitchFullAsync(
  109. string sourceCode = SqlServerSourceCode,
  110. string entityCode = SqlServerEntityCode,
  111. long tenantId = 0,
  112. CancellationToken cancellationToken = default)
  113. {
  114. cancellationToken.ThrowIfCancellationRequested();
  115. await EnsureTablesAsync();
  116. await EnsureStgTableAsync();
  117. var now = DateTime.Now;
  118. var pullCtx = new MdpPullContext
  119. {
  120. TenantId = tenantId,
  121. FullRefresh = true,
  122. TaskCode = "S5_PRODUCTION_ISSUE_INBOUND",
  123. BatchId = $"S5_PROD_ISSUE_SW_{now:yyyyMMddHHmmss}"
  124. };
  125. var pull = await _pullDispatcher.PullAllByEntityCodeAsync(entityCode, pullCtx, cancellationToken);
  126. var batchId = $"S5_PROD_ISSUE_SWSTD_{now:yyyyMMddHHmmss}";
  127. var runLogId = await InsertRunLogAsync(batchId, now, "SOURCE_SWITCH");
  128. var result = new ProductionIssueMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  129. try
  130. {
  131. result.HeadRows = await MdpStdFullReplace.ReplaceAsync(
  132. _db, "mdp_std_production_issue", tenantId, extraWhere: null,
  133. insertScopedAsync: () => TransformHeadStandardAsync(batchId, now, sourceCode),
  134. cancellationToken);
  135. await MarkRunSuccessAsync(runLogId, now, result);
  136. }
  137. catch (Exception ex)
  138. {
  139. await MarkRunFailedAsync(runLogId, now, ex.Message);
  140. throw;
  141. }
  142. return new ProductionIssueInboundResult
  143. {
  144. PullBatchId = pullCtx.BatchId,
  145. RowsPulled = pull.RowsPulled,
  146. RowsWrittenStg = pull.RowsWritten,
  147. NewCursor = pull.NewCursor,
  148. TransformBatchId = batchId,
  149. StdRows = result.HeadRows
  150. };
  151. }
  152. private async Task<(int RowsPulled, int RowsWritten, string? NewCursor)> PopulateStgAsync(
  153. string entityCode, MdpPullContext pullCtx, CancellationToken cancellationToken)
  154. {
  155. var r = await _pullDispatcher.PullByEntityCodeAsync(entityCode, pullCtx, cancellationToken);
  156. return (r.RowsPulled, r.RowsWritten, r.NewCursor);
  157. }
  158. private async Task EnsureStgTableAsync()
  159. {
  160. await _db.Ado.ExecuteCommandAsync(
  161. """
  162. CREATE TABLE IF NOT EXISTS mdp_stg_production_issue (
  163. id BIGINT PRIMARY KEY AUTO_INCREMENT,
  164. tenant_id BIGINT NOT NULL,
  165. source_system VARCHAR(50) NULL,
  166. source_table VARCHAR(200),
  167. source_row_id VARCHAR(200),
  168. source_biz_key VARCHAR(300) NULL,
  169. raw_data JSON,
  170. sync_batch_id VARCHAR(100),
  171. sync_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP,
  172. process_status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
  173. process_message VARCHAR(500) NULL,
  174. create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
  175. update_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  176. UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key),
  177. KEY idx_batch (sync_batch_id),
  178. KEY idx_src (source_table, source_row_id),
  179. KEY idx_tenant (tenant_id)
  180. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='S5生产领料执行器贴源层'
  181. """);
  182. }
  183. /// <summary>防御式建表(与 UpdateScripts DDL 同构,幂等)。</summary>
  184. private async Task EnsureTablesAsync()
  185. {
  186. await _db.Ado.ExecuteCommandAsync(
  187. """
  188. CREATE TABLE IF NOT EXISTS mdp_std_production_issue (
  189. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  190. tenant_id BIGINT NOT NULL DEFAULT 0,
  191. factory_id BIGINT NULL DEFAULT 1,
  192. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  193. domain VARCHAR(80) NOT NULL,
  194. rec_id INT NOT NULL,
  195. nbr VARCHAR(24) NULL,
  196. issue_date DATETIME NULL,
  197. status VARCHAR(8) NULL,
  198. status_desc VARCHAR(20) NULL,
  199. work_ord VARCHAR(64) NULL,
  200. department VARCHAR(8) NULL,
  201. department_desc VARCHAR(255) NULL,
  202. qty_ord DECIMAL(18,5) NULL DEFAULT 0,
  203. prod_line VARCHAR(8) NULL,
  204. applicant_name VARCHAR(12) NULL,
  205. issue_user VARCHAR(255) NULL,
  206. user1 TEXT NULL,
  207. remark VARCHAR(200) NULL,
  208. create_user VARCHAR(24) NULL,
  209. source_create_time DATETIME NULL,
  210. pretreatment_state VARCHAR(50) NULL,
  211. trans_type VARCHAR(24) NULL,
  212. trans_type_text VARCHAR(50) NULL,
  213. eff_date DATETIME NULL,
  214. address VARCHAR(120) NULL,
  215. source_biz_key VARCHAR(200) NULL,
  216. sync_batch_id VARCHAR(100) NOT NULL,
  217. sync_time DATETIME NOT NULL,
  218. create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
  219. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  220. UNIQUE KEY uk_mdp_std_prod_issue (tenant_id, domain, rec_id),
  221. KEY idx_mdp_std_prod_issue_nbr (tenant_id, nbr),
  222. KEY idx_mdp_std_prod_issue_date (tenant_id, issue_date),
  223. KEY idx_mdp_std_prod_issue_status (tenant_id, status),
  224. KEY idx_mdp_std_prod_issue_workord (tenant_id, work_ord),
  225. KEY idx_mdp_std_prod_issue_prodline (tenant_id, prod_line),
  226. KEY idx_mdp_std_prod_issue_dept (tenant_id, department)
  227. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S5生产领料单头标准层'
  228. """);
  229. }
  230. /// <summary>
  231. /// 标准化头:stg(NbrMaster,SM) + DepartmentMaster(本库) → mdp_std_production_issue。
  232. /// sourceSystem 非空时仅转当前源(切源 FULL Replace 用);为空时全源(本库全量)。跨源类型经 MdpJsonSql 兼容。
  233. /// </summary>
  234. private async Task<int> TransformHeadStandardAsync(string batchId, DateTime now, string? sourceSystem = null)
  235. {
  236. var srcClause = sourceSystem == null ? "" : " AND m.source_system=@Src";
  237. // 表达式片段(跨源类型兼容)
  238. var domainE = MdpJsonSql.Str("m", "Domain");
  239. var deptE = MdpJsonSql.Str("m", "Department");
  240. var statusUpper = $"UPPER(IFNULL({MdpJsonSql.Str("m", "Status")},''))";
  241. var pretreatE = MdpJsonSql.Str("m", "PretreatmentState");
  242. var transTypeE = MdpJsonSql.Str("m", "TransType");
  243. var countPars = new List<SugarParameter>();
  244. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  245. var rows = await _db.Ado.GetIntAsync(
  246. $"SELECT COUNT(1) FROM mdp_stg_production_issue m WHERE m.source_table='NbrMaster' " +
  247. $"AND {MdpJsonSql.Str("m", "Type")}='SM' AND {MdpJsonSql.BoolTrue("m", "IsActive")}{srcClause}",
  248. countPars);
  249. var insertSql =
  250. $"""
  251. INSERT INTO mdp_std_production_issue
  252. (tenant_id, factory_id, source_system, domain, rec_id, nbr, issue_date, status, status_desc,
  253. work_ord, department, department_desc, qty_ord, prod_line, applicant_name, issue_user,
  254. user1, remark, create_user, source_create_time, pretreatment_state, trans_type, trans_type_text,
  255. eff_date, address, source_biz_key, sync_batch_id, sync_time)
  256. SELECT
  257. IFNULL({MdpJsonSql.Int("m", "tenant_id")}, IFNULL(m.tenant_id, 0)), 1, IFNULL(NULLIF(m.source_system,''), 'AIDOP'),
  258. {domainE}, {MdpJsonSql.Int("m", "RecID")}, {MdpJsonSql.Str("m", "Nbr")},
  259. {MdpJsonSql.DateTimeSec("m", "Date")},
  260. {statusUpper},
  261. CASE
  262. WHEN IFNULL({pretreatE},'')<>'' THEN {pretreatE}
  263. WHEN {statusUpper}='C' THEN '已下架'
  264. WHEN {statusUpper}='A' THEN '备料中'
  265. ELSE ''
  266. END,
  267. {MdpJsonSql.Str("m", "WorkOrd")}, {deptE},
  268. TRIM(CONCAT(IFNULL({deptE},''), ' ', IFNULL(d.Descr, ''))),
  269. {MdpJsonSql.Dec("m", "QtyOrd", 18, 5)}, {MdpJsonSql.Str("m", "ProdLine")}, {MdpJsonSql.Str("m", "Name")},
  270. {MdpJsonSql.Str("m", "User1")},
  271. {MdpJsonSql.Str("m", "User1")}, {MdpJsonSql.Str("m", "Remark")}, {MdpJsonSql.Str("m", "CreateUser")},
  272. {MdpJsonSql.DateTimeSec("m", "CreateTime")},
  273. IFNULL({pretreatE}, ''),
  274. CASE WHEN {transTypeE}='Z61' THEN '补料' ELSE '正常' END,
  275. CASE WHEN {transTypeE}='PrevProcess' THEN '需要前处理' ELSE '' END,
  276. {MdpJsonSql.DateTimeSec("m", "EffDate")}, {MdpJsonSql.Str("m", "Address")},
  277. IFNULL(NULLIF(m.source_biz_key,''), CONCAT(IFNULL({domainE},''), ':', IFNULL({MdpJsonSql.Str("m", "RecID")},''))),
  278. @BatchId, @Now
  279. FROM mdp_stg_production_issue m
  280. LEFT JOIN DepartmentMaster d ON d.Domain = {domainE} AND d.Department = {deptE}
  281. WHERE m.source_table='NbrMaster'
  282. AND {MdpJsonSql.Str("m", "Type")}='SM'
  283. AND {MdpJsonSql.BoolTrue("m", "IsActive")}{srcClause}
  284. ON DUPLICATE KEY UPDATE
  285. factory_id=VALUES(factory_id), source_system=VALUES(source_system), nbr=VALUES(nbr), issue_date=VALUES(issue_date),
  286. status=VALUES(status), status_desc=VALUES(status_desc), work_ord=VALUES(work_ord),
  287. department=VALUES(department), department_desc=VALUES(department_desc), qty_ord=VALUES(qty_ord),
  288. prod_line=VALUES(prod_line), applicant_name=VALUES(applicant_name), issue_user=VALUES(issue_user),
  289. user1=VALUES(user1), remark=VALUES(remark), create_user=VALUES(create_user),
  290. source_create_time=VALUES(source_create_time), pretreatment_state=VALUES(pretreatment_state),
  291. trans_type=VALUES(trans_type), trans_type_text=VALUES(trans_type_text),
  292. eff_date=VALUES(eff_date), address=VALUES(address),
  293. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  294. """;
  295. var insPars = new List<SugarParameter>
  296. {
  297. new("@BatchId", batchId),
  298. new("@Now", now)
  299. };
  300. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  301. await _db.Ado.ExecuteCommandAsync(insertSql, insPars);
  302. return rows;
  303. }
  304. private async Task<long> InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  305. {
  306. await _db.Ado.ExecuteCommandAsync(
  307. """
  308. INSERT INTO mdp_transform_run_log
  309. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  310. VALUES (0, @JobCode, 'S5生产领料单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  311. """,
  312. new SugarParameter("@JobCode", JobCode),
  313. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  314. new SugarParameter("@BatchId", batchId),
  315. new SugarParameter("@StartTime", startedAt));
  316. return await _db.Ado.GetLongAsync(
  317. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  318. new List<SugarParameter> { new("@BatchId", batchId) });
  319. }
  320. private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, ProductionIssueMdpSyncResult result)
  321. {
  322. var finishedAt = DateTime.Now;
  323. await _db.Ado.ExecuteCommandAsync(
  324. """
  325. UPDATE mdp_transform_run_log
  326. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  327. stage_rows=0, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP
  328. WHERE id=@Id
  329. """,
  330. new SugarParameter("@EndTime", finishedAt),
  331. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  332. new SugarParameter("@StandardRows", result.HeadRows),
  333. new SugarParameter("@Id", runLogId));
  334. }
  335. private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message)
  336. {
  337. var finishedAt = DateTime.Now;
  338. await _db.Ado.ExecuteCommandAsync(
  339. """
  340. UPDATE mdp_transform_run_log
  341. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  342. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  343. WHERE id=@Id
  344. """,
  345. new SugarParameter("@EndTime", finishedAt),
  346. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  347. new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message),
  348. new SugarParameter("@Id", runLogId));
  349. }
  350. private static string NormalizeTriggerType(string? triggerType)
  351. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  352. }
  353. /// <summary>生产领料单 MDP 同步转换结果。</summary>
  354. public sealed class ProductionIssueMdpSyncResult
  355. {
  356. public long RunLogId { get; set; }
  357. public string BatchId { get; set; } = string.Empty;
  358. public int HeadRows { get; set; }
  359. }
  360. /// <summary>生产领料双源入站结果(stg 抽数 + std 转换)。</summary>
  361. public sealed class ProductionIssueInboundResult
  362. {
  363. public string PullBatchId { get; set; } = string.Empty;
  364. public int RowsPulled { get; set; }
  365. public int RowsWrittenStg { get; set; }
  366. public string? NewCursor { get; set; }
  367. public string TransformBatchId { get; set; } = string.Empty;
  368. public int StdRows { get; set; }
  369. }