IpqcInspectionMdpSyncService.cs 46 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813
  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.Manufacturing;
  5. /// <summary>
  6. /// S6 过程检验单(IPQC)数据中台只读同步转换服务(DOP 内部,独立于 S6MdpSyncTransformService 的 KPI 管线)。
  7. ///
  8. /// 源:aidopdev.qms_gcjyd(过程检验单头,31 列)+ aidopdev.qms_gcjydzb(明细,100 列),归口自旧系统 dopdemo。
  9. /// 链路:执行器 → mdp_stg_ipqc_pull(头 source_table='qms_gcjyd' + 明细 source_table='qms_gcjydzb',共表)→ std(头/明细各读 pull stg)。
  10. /// 双模式入站:mdp_entity=S6_IPQC_INSPECTION_HEAD / S6_IPQC_INSPECTION_DETAIL → Writer → mdp_stg_ipqc_pull;头/明细标准层均读 pull stg(MdpJsonSql 跨源类型兼容)。
  11. /// SyncHeadStagingAsync / SyncDetailStagingAsync 保留为 vestigial 可观测贴源写入(mdp_stg_ipqc_inspection(_detail)),不再被 transform 消费。
  12. ///
  13. /// 归口说明(CTO 拍板方案 A):
  14. /// qms_gcjyd/qms_gcjydzb 原仅存于旧系统 dopdemo(0 行)、aidopdev 无此表,且 dopdemo 非运行时连接。
  15. /// 本服务 EnsureTablesAsync 在 aidopdev 建同构空表(IF NOT EXISTS)作为归口源;
  16. /// 后续由上游一次性抽数填充,届时无需改代码即随刷新流通。
  17. ///
  18. /// 约束:
  19. /// - 只读 qms_gcjyd/qms_gcjydzb,仅写 mdp_stg_*/mdp_std_*/run_log;绝不 INSERT/UPDATE/DELETE 业务源单据表。
  20. /// - 源表无 tenant/org/factory 列 → 标准层按常量 tenant_id=0, factory_id=1 落库。
  21. /// - 明细头关联:qms_gcjydzb.glid -> qms_gcjyd.id;标准层经 std_head_id 回填 mdp_std_ipqc_inspection.id。
  22. /// - 明细只转换已确认子集列(检验项目/依据/工序/方法/技术标准/上下限/测量仪器/检验人/检验时间/样本量/检验数量/不合格数量/判定/备注等);
  23. /// j1~j70 动态检验值列语义不明、后置不转换、不落 typed 字段。
  24. /// - result_judgement(jgpd/pd)为字典码,原样透出,语义解码后置。
  25. /// - 源为 0 行时转换成功完成、处理数为 0,不报错。
  26. /// </summary>
  27. public class IpqcInspectionMdpSyncService : ITransient
  28. {
  29. private const string JobCode = "S6_IPQC_INSPECTION_MDP_SYNC";
  30. private const string InboundEntityCode = "S6_IPQC_INSPECTION_HEAD";
  31. // Phase 2:Detail(qms_gcjydzb) 迁通用管线,经执行器灌 mdp_stg_ipqc_pull(source_table='qms_gcjydzb',与 Head 共表),transform 读 pull stg。
  32. private const string InboundDetailEntityCode = "S6_IPQC_INSPECTION_DETAIL";
  33. // 双源(dopdemorq SQL Server):Head(qms_gcjyd)+ Detail(qms_gcjydzb),均无 UpdateTime → FULL。实体默认 status=0(就位不启用)。
  34. private const string SqlServerSourceCode = "DOPDEMORQ_SQLSERVER";
  35. private const string SqlServerHeadEntityCode = "S6_IPQC_INSPECTION_HEAD_SQLSERVER";
  36. private const string SqlServerDetailEntityCode = "S6_IPQC_INSPECTION_DETAIL_SQLSERVER";
  37. /// <summary>qms_gcjyd/qms_gcjydzb 无 tenant_id 列;165 侧亦同,须经账套映射解析。</summary>
  38. private const string SourceZtid = "pbxfxp";
  39. private readonly ISqlSugarClient _db;
  40. private readonly TransformRunLogFinalizer _runLogFinalizer;
  41. private readonly MdpSourcePullDispatcher _pullDispatcher;
  42. private readonly List<string> _tenantSkipWarnings = new();
  43. public IpqcInspectionMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher, TransformRunLogFinalizer runLogFinalizer)
  44. {
  45. _db = db;
  46. _runLogFinalizer = runLogFinalizer;
  47. _pullDispatcher = pullDispatcher;
  48. }
  49. /// <summary>
  50. /// 全量:执行器灌头 pull-stg + 明细贴源 → 标准化(头/明细读 stg)。
  51. /// </summary>
  52. public async Task<IpqcInspectionMdpSyncResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
  53. {
  54. cancellationToken.ThrowIfCancellationRequested();
  55. await EnsureTablesAsync();
  56. await EnsurePullStgTableAsync();
  57. var now = DateTime.Now;
  58. var batchId = $"S6_IPQC_INSP_FULL_{now:yyyyMMddHHmmss}";
  59. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  60. var result = new IpqcInspectionMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  61. _tenantSkipWarnings.Clear();
  62. try
  63. {
  64. // 源表无 tenant_id;ctx 须传入经账套映射解析的真实租户。
  65. var pullCtx = new MdpPullContext
  66. {
  67. TenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, 0),
  68. FullRefresh = true,
  69. TaskCode = "S6_IPQC_INSPECTION_INBOUND",
  70. BatchId = $"{batchId}_PULL"
  71. };
  72. var pullHead = await _pullDispatcher.PullByEntityCodeAsync(InboundEntityCode, pullCtx, cancellationToken);
  73. var pullDetail = await _pullDispatcher.PullByEntityCodeAsync(InboundDetailEntityCode, pullCtx, cancellationToken);
  74. result.HeadStageRows = pullHead.RowsWritten;
  75. result.DetailStageRows = pullDetail.RowsWritten;
  76. // 保留富贴源头/明细表写入(vestigial 可观测);头/明细 std 均以 pull stg 为准
  77. await SyncHeadStagingAsync(batchId, now);
  78. await SyncDetailStagingAsync(batchId, now);
  79. result.HeadStandardRows = await TransformHeadStandardAsync(batchId, now);
  80. result.DetailStandardRows = await TransformDetailStandardAsync(batchId, now);
  81. result.TenantSkipWarnings = _tenantSkipWarnings.ToList();
  82. await MarkRunSuccessAsync(runLogId, now, result);
  83. return result;
  84. }
  85. catch (Exception ex)
  86. {
  87. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  88. if (!_runLogFinalizer.IsHostStopping)
  89. await MarkRunFailedAsync(runLogId, now, ex.Message);
  90. throw;
  91. }
  92. finally
  93. {
  94. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  95. }
  96. }
  97. /// <summary>本库 MySQL 贴源 source_system(与 MdpDbPullExecutor 落 stg 的取值一致,决定 uk_source_key 归属)。</summary>
  98. private const string NativeSourceSystem = "AIDOPDEV_MYSQL";
  99. /// <summary>targeted sync 的单据作用域(租户 + 检验单 id)。</summary>
  100. private sealed class IpqcBillScope
  101. {
  102. public long TenantId { get; init; }
  103. public long InspectionId { get; init; }
  104. }
  105. /// <summary>
  106. /// 单据级 targeted 同步(业务写入后立即可见),替代把 Full/Inbound 全表扫描挂到高频提交路径。
  107. ///
  108. /// 只处理 inspectionId 这一张过程检验单的<b>头表</b>:
  109. /// qms_gcjyd (id=inspectionId AND tenant_id=tenantId) → mdp_stg_ipqc_pull → mdp_std_ipqc_inspection
  110. ///
  111. /// <b>刻意不同步明细</b>(qms_gcjydzb → mdp_std_ipqc_inspection_detail):该表全仓无任何应用
  112. /// INSERT/UPDATE/DELETE,源不变则标准层不可能因业务写入而滞后,纳入即为无依据的额外开销。
  113. ///
  114. /// 边界:只读业务源、只写 mdp_stg_ipqc_pull / mdp_std_ipqc_inspection / run_log;
  115. /// 租户双重收口(源查询 + transform 均按 tenantId);纯 upsert,无 DELETE / 无 orphan purge /
  116. /// 无 MdpStdFullReplace / 不复用 RunFullAsync 的 pbxfxp 账套映射租户。
  117. /// </summary>
  118. public async Task<IpqcInspectionMdpSyncResult> SyncInspectionAsync(
  119. long inspectionId, long tenantId, string triggerType = "BILL", CancellationToken cancellationToken = default)
  120. {
  121. cancellationToken.ThrowIfCancellationRequested();
  122. if (inspectionId <= 0 || tenantId <= 0)
  123. return new IpqcInspectionMdpSyncResult { BatchId = "SKIPPED_INVALID_SCOPE" };
  124. await EnsureTablesAsync();
  125. await EnsurePullStgTableAsync();
  126. var now = DateTime.Now;
  127. var batchId = $"S6_IPQC_BILL_{now:yyyyMMddHHmmssfff}_{inspectionId}";
  128. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  129. var result = new IpqcInspectionMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  130. _tenantSkipWarnings.Clear();
  131. try
  132. {
  133. var scope = new IpqcBillScope { TenantId = tenantId, InspectionId = inspectionId };
  134. // ① 贴源:把该单据的源行刷进 pull stg(transform 一律读 stg,不先刷就还是旧快照)
  135. result.HeadStageRows = await UpsertHeadStgFromSourceAsync(scope, batchId, now);
  136. // ② 标准层:复用全量同一套 transform,仅加单据作用域,字段映射零分叉
  137. result.HeadStandardRows = await TransformHeadStandardAsync(batchId, now, null, scope);
  138. result.TenantSkipWarnings = _tenantSkipWarnings.ToList();
  139. await MarkRunSuccessAsync(runLogId, now, result);
  140. return result;
  141. }
  142. catch (Exception ex)
  143. {
  144. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  145. if (!_runLogFinalizer.IsHostStopping)
  146. await MarkRunFailedAsync(runLogId, now, ex.Message);
  147. throw;
  148. }
  149. finally
  150. {
  151. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  152. }
  153. }
  154. /// <summary>
  155. /// 单据级头表贴源 upsert:qms_gcjyd 源行 → mdp_stg_ipqc_pull。
  156. /// raw_data 按源表<b>全列</b>动态生成 JSON_OBJECT,与 MdpDbPullExecutor/MdpStagingWriter 的信封形状一致
  157. /// (JSON 键 = 源列名;source_row_id = id;source_biz_key = djbh,空则回落 id,与 MdpStagingWriter.BuildBizKey
  158. /// 的 "biz_key_expr=djbh,取不到则回落 sourceRowId" 语义一致),因此命中既有 uk_source_key 行原地更新,
  159. /// 不产生重复 stg 行,也不会因 transform 将来新增字段而漂移。
  160. /// </summary>
  161. private async Task<int> UpsertHeadStgFromSourceAsync(IpqcBillScope scope, string batchId, DateTime now)
  162. {
  163. var cols = await _db.Ado.SqlQueryAsync<string>(
  164. "SELECT COLUMN_NAME FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME='qms_gcjyd' ORDER BY ORDINAL_POSITION");
  165. // 列名取自 information_schema 且强制标识符白名单,非用户输入,无注入面
  166. cols = cols.Where(c => !string.IsNullOrWhiteSpace(c) && System.Text.RegularExpressions.Regex.IsMatch(c, "^[A-Za-z0-9_]+$")).ToList();
  167. if (cols.Count == 0) return 0;
  168. var jsonObj = "JSON_OBJECT(" + string.Join(", ", cols.Select(c => $"'{c}', s.`{c}`")) + ")";
  169. return await MdpSchemaAligner.ExecuteAsync(_db,
  170. $"""
  171. INSERT INTO mdp_stg_ipqc_pull
  172. (tenant_id, source_system, source_table, source_row_id, source_biz_key, raw_data, sync_batch_id, sync_time, process_status, create_time)
  173. SELECT s.tenant_id, @SrcSys, 'qms_gcjyd', CAST(s.id AS CHAR),
  174. IFNULL(NULLIF(s.djbh,''), CAST(s.id AS CHAR)), {jsonObj},
  175. @BatchId, @Now, 'PENDING', @Now
  176. FROM qms_gcjyd s
  177. WHERE s.id=@Id AND s.tenant_id=@TenantId
  178. ON DUPLICATE KEY UPDATE
  179. tenant_id=VALUES(tenant_id), source_row_id=VALUES(source_row_id), raw_data=VALUES(raw_data),
  180. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  181. process_status='PENDING', update_time=CURRENT_TIMESTAMP
  182. """,
  183. new List<SugarParameter>
  184. {
  185. new("@SrcSys", NativeSourceSystem),
  186. new("@BatchId", batchId),
  187. new("@Now", now),
  188. new("@Id", scope.InspectionId),
  189. new("@TenantId", scope.TenantId),
  190. });
  191. }
  192. /// <summary>
  193. /// 双模式入站:执行器抽数落 pull-stg,再跑标准层(头读 pull stg)。
  194. /// </summary>
  195. public async Task<IpqcInspectionInboundResult> RunInboundAsync(
  196. long tenantId = 0,
  197. bool fullRefresh = false,
  198. CancellationToken cancellationToken = default)
  199. {
  200. cancellationToken.ThrowIfCancellationRequested();
  201. await EnsureTablesAsync();
  202. await EnsurePullStgTableAsync();
  203. var now = DateTime.Now;
  204. // 源表无 tenant_id;显式 tenantId≤0 时经账套映射解析真实租户。
  205. var resolvedTenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId);
  206. var pullCtx = new MdpPullContext
  207. {
  208. TenantId = resolvedTenantId,
  209. FullRefresh = fullRefresh,
  210. TaskCode = "S6_IPQC_INSPECTION_INBOUND",
  211. BatchId = $"S6_IPQC_IN_{now:yyyyMMddHHmmss}"
  212. };
  213. var pullHead = await _pullDispatcher.PullByEntityCodeAsync(InboundEntityCode, pullCtx, cancellationToken);
  214. var pullDetail = await _pullDispatcher.PullByEntityCodeAsync(InboundDetailEntityCode, pullCtx, cancellationToken);
  215. var batchId = $"S6_IPQC_STD_{now:yyyyMMddHHmmss}";
  216. var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
  217. var result = new IpqcInspectionMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  218. _tenantSkipWarnings.Clear();
  219. try
  220. {
  221. result.HeadStageRows = pullHead.RowsWritten;
  222. result.DetailStageRows = pullDetail.RowsWritten;
  223. await SyncHeadStagingAsync(batchId, now);
  224. await SyncDetailStagingAsync(batchId, now);
  225. result.HeadStandardRows = await TransformHeadStandardAsync(batchId, now);
  226. result.DetailStandardRows = await TransformDetailStandardAsync(batchId, now);
  227. result.TenantSkipWarnings = _tenantSkipWarnings.ToList();
  228. await MarkRunSuccessAsync(runLogId, now, result);
  229. }
  230. catch (Exception ex)
  231. {
  232. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  233. if (!_runLogFinalizer.IsHostStopping)
  234. await MarkRunFailedAsync(runLogId, now, ex.Message);
  235. throw;
  236. }
  237. finally
  238. {
  239. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  240. }
  241. return new IpqcInspectionInboundResult
  242. {
  243. PullBatchId = pullCtx.BatchId,
  244. RowsPulled = pullHead.RowsPulled + pullDetail.RowsPulled,
  245. RowsWrittenStg = pullHead.RowsWritten + pullDetail.RowsWritten,
  246. NewCursor = pullDetail.NewCursor ?? pullHead.NewCursor,
  247. TransformBatchId = result.BatchId,
  248. HeadStageRows = result.HeadStageRows,
  249. HeadStandardRows = result.HeadStandardRows,
  250. DetailStageRows = result.DetailStageRows,
  251. DetailStandardRows = result.DetailStandardRows
  252. };
  253. }
  254. /// <summary>
  255. /// 双源切换 FULL Replace(Phase 1,仅 Head):从指定 DB 源(默认 DOPDEMORQ_SQLSERVER)PullAll 全量灌 pull-stg,
  256. /// 成功后在单事务内以「仅当前源」结果 FULL 重建 mdp_std_ipqc_inspection(按 tenant 精确隔离),消除旧源独有业务键残留。
  257. /// qms_gcjyd 无 UpdateTime → FULL;用 PullAllByEntityCodeAsync + FullRefresh=true。
  258. /// <b>只处理 Head</b>:不调用 SyncHeadStagingAsync(vestigial 观测写入)/ SyncDetailStagingAsync / TransformDetailStandardAsync;
  259. /// IPQC Detail(qms_gcjydzb) 属 bespoke,留 Phase 2。
  260. /// SQLSERVER 实体默认 status=0(就位不启用);dopdemorq 该表当前空,Phase 1 不实际执行 destructive 切换。
  261. /// </summary>
  262. public async Task<IpqcInspectionInboundResult> RunHeadSourceSwitchFullAsync(
  263. string sourceCode = SqlServerSourceCode,
  264. string headEntityCode = SqlServerHeadEntityCode,
  265. long tenantId = 0,
  266. CancellationToken cancellationToken = default)
  267. {
  268. cancellationToken.ThrowIfCancellationRequested();
  269. await EnsureTablesAsync();
  270. await EnsurePullStgTableAsync();
  271. // 165 SQL Server 源无 tenant_id;显式 tenantId≤0 时经账套映射解析。
  272. tenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId);
  273. var now = DateTime.Now;
  274. // 1) PullAll 抽尽 + fullRefresh=true。Pull 失败会抛异常,std 未动。
  275. var pullCtx = new MdpPullContext
  276. {
  277. TenantId = tenantId,
  278. FullRefresh = true,
  279. TaskCode = "S6_IPQC_INSPECTION_INBOUND",
  280. BatchId = $"S6_IPQC_SW_{now:yyyyMMddHHmmss}"
  281. };
  282. var pull = await _pullDispatcher.PullAllByEntityCodeAsync(headEntityCode, pullCtx, cancellationToken);
  283. // 2) FULL Replace(仅头 std):事务内 DELETE 当前 tenant,再 INSERT 仅当前 source_system 的头结果。
  284. var batchId = $"S6_IPQC_SWSTD_{now:yyyyMMddHHmmss}";
  285. var runLogId = await InsertRunLogAsync(batchId, now, "SOURCE_SWITCH");
  286. var result = new IpqcInspectionMdpSyncResult { BatchId = batchId, RunLogId = runLogId, HeadStageRows = pull.RowsWritten };
  287. try
  288. {
  289. result.HeadStandardRows = await MdpStdFullReplace.ReplaceAsync(
  290. _db, "mdp_std_ipqc_inspection", tenantId, extraWhere: null,
  291. insertScopedAsync: () => TransformHeadStandardAsync(batchId, now, sourceCode),
  292. cancellationToken);
  293. await MarkRunSuccessAsync(runLogId, now, result);
  294. }
  295. catch (Exception ex)
  296. {
  297. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  298. if (!_runLogFinalizer.IsHostStopping)
  299. await MarkRunFailedAsync(runLogId, now, ex.Message);
  300. throw;
  301. }
  302. finally
  303. {
  304. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  305. }
  306. return new IpqcInspectionInboundResult
  307. {
  308. PullBatchId = pullCtx.BatchId,
  309. RowsPulled = pull.RowsPulled,
  310. RowsWrittenStg = pull.RowsWritten,
  311. NewCursor = pull.NewCursor,
  312. TransformBatchId = batchId,
  313. HeadStageRows = result.HeadStageRows,
  314. HeadStandardRows = result.HeadStandardRows,
  315. DetailStageRows = 0,
  316. DetailStandardRows = 0
  317. };
  318. }
  319. /// <summary>
  320. /// 双源切换 FULL Replace(Detail):从指定 DB 源(默认 DOPDEMORQ_SQLSERVER)PullAll 全量灌 pull-stg(source_table='qms_gcjydzb'),
  321. /// 成功后在单事务内以「仅当前源」结果 FULL 重建 mdp_std_ipqc_inspection_detail(按 tenant 精确隔离),消除旧源独有业务键残留。
  322. /// qms_gcjydzb 无 UpdateTime → FULL;用 PullAllByEntityCodeAsync + FullRefresh=true。std_head_id 回填依赖 Head std 已存在(LEFT JOIN,缺则 NULL)。
  323. /// <b>只处理 Detail</b>:不碰 Head std(切 Head 用 RunHeadSourceSwitchFullAsync)。
  324. /// SQLSERVER 实体默认 status=0(就位不启用);dopdemorq 该表当前空,Phase 2 不实际执行 destructive 切换。
  325. /// </summary>
  326. public async Task<IpqcInspectionInboundResult> RunDetailSourceSwitchFullAsync(
  327. string sourceCode = SqlServerSourceCode,
  328. string detailEntityCode = SqlServerDetailEntityCode,
  329. long tenantId = 0,
  330. CancellationToken cancellationToken = default)
  331. {
  332. cancellationToken.ThrowIfCancellationRequested();
  333. await EnsureTablesAsync();
  334. await EnsurePullStgTableAsync();
  335. // 165 SQL Server 源无 tenant_id;显式 tenantId≤0 时经账套映射解析。
  336. tenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId);
  337. var now = DateTime.Now;
  338. // 1) PullAll 抽尽 + fullRefresh=true。Pull 失败会抛异常,std 未动。
  339. var pullCtx = new MdpPullContext
  340. {
  341. TenantId = tenantId,
  342. FullRefresh = true,
  343. TaskCode = "S6_IPQC_INSPECTION_INBOUND",
  344. BatchId = $"S6_IPQC_DTL_SW_{now:yyyyMMddHHmmss}"
  345. };
  346. var pull = await _pullDispatcher.PullAllByEntityCodeAsync(detailEntityCode, pullCtx, cancellationToken);
  347. // 2) FULL Replace(仅明细 std):事务内 DELETE 当前 tenant,再 INSERT 仅当前 source_system 的明细结果。
  348. var batchId = $"S6_IPQC_DTL_SWSTD_{now:yyyyMMddHHmmss}";
  349. var runLogId = await InsertRunLogAsync(batchId, now, "SOURCE_SWITCH");
  350. var result = new IpqcInspectionMdpSyncResult { BatchId = batchId, RunLogId = runLogId, DetailStageRows = pull.RowsWritten };
  351. try
  352. {
  353. result.DetailStandardRows = await MdpStdFullReplace.ReplaceAsync(
  354. _db, "mdp_std_ipqc_inspection_detail", tenantId, extraWhere: null,
  355. insertScopedAsync: () => TransformDetailStandardAsync(batchId, now, sourceCode),
  356. cancellationToken);
  357. await MarkRunSuccessAsync(runLogId, now, result);
  358. }
  359. catch (Exception ex)
  360. {
  361. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  362. if (!_runLogFinalizer.IsHostStopping)
  363. await MarkRunFailedAsync(runLogId, now, ex.Message);
  364. throw;
  365. }
  366. finally
  367. {
  368. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  369. }
  370. return new IpqcInspectionInboundResult
  371. {
  372. PullBatchId = pullCtx.BatchId,
  373. RowsPulled = pull.RowsPulled,
  374. RowsWrittenStg = pull.RowsWritten,
  375. NewCursor = pull.NewCursor,
  376. TransformBatchId = batchId,
  377. HeadStageRows = 0,
  378. HeadStandardRows = 0,
  379. DetailStageRows = result.DetailStageRows,
  380. DetailStandardRows = result.DetailStandardRows
  381. };
  382. }
  383. private async Task EnsurePullStgTableAsync()
  384. {
  385. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_ipqc_pull"));
  386. }
  387. /// <summary>防御式建表(幂等):归口源空表 qms_gcjyd/qms_gcjydzb + 贴源层(头/明细) + 标准层(头/明细)。</summary>
  388. private async Task EnsureTablesAsync()
  389. {
  390. // 归口源表(方案 A):同构自旧系统 dopdemo.qms_gcjyd 头表,空表,待上游抽数填充。
  391. // id 不设 AUTO_INCREMENT:抽数将携带旧系统主键。
  392. await MdpSchemaAligner.ExecuteAsync(_db,
  393. """
  394. CREATE TABLE IF NOT EXISTS qms_gcjyd (
  395. id BIGINT NOT NULL PRIMARY KEY COMMENT '主表id',
  396. tenant_id BIGINT NOT NULL COMMENT 'tenant_id',
  397. djbh VARCHAR(255) NULL COMMENT '单据编号',
  398. cplx VARCHAR(255) NULL COMMENT '产品型号',
  399. scph VARCHAR(255) NULL COMMENT '生产批号',
  400. jyrq VARCHAR(255) NULL COMMENT '检验日期',
  401. lydjbh VARCHAR(255) NULL COMMENT '来源单据编号',
  402. jgpd BIGINT NULL COMMENT '结果判定',
  403. jyhgsl DECIMAL(10,0) NULL COMMENT '检验合格数量',
  404. jybhgsl DECIMAL(10,0) NULL COMMENT '检验不合格数量',
  405. fj TEXT NULL COMMENT '附件',
  406. bz TEXT NULL COMMENT '备注',
  407. jyr VARCHAR(255) NULL COMMENT '检验人',
  408. gxbm VARCHAR(255) NULL COMMENT '工序编码',
  409. gxmc VARCHAR(255) NULL COMMENT '工序名称',
  410. sczyry VARCHAR(255) NULL COMMENT '生产作业人员',
  411. ybl DECIMAL(10,0) NULL COMMENT '样本量',
  412. bdbh VARCHAR(255) NULL COMMENT '表单编号',
  413. bbh VARCHAR(255) NULL COMMENT '版本号',
  414. sxrq DATETIME NULL COMMENT '生效日期',
  415. wlbm VARCHAR(255) NULL COMMENT '物料编码',
  416. wlmc VARCHAR(255) NULL COMMENT '物料名称',
  417. lysl DECIMAL(10,0) NULL COMMENT '留样数量',
  418. phsl DECIMAL(10,0) NULL COMMENT '破坏数量',
  419. jygfid VARCHAR(255) NULL COMMENT '检验规范',
  420. jygfbbm VARCHAR(255) NULL COMMENT '没有用',
  421. jgbb VARCHAR(255) NULL COMMENT '检规版本',
  422. jgbh VARCHAR(255) NULL COMMENT '检规编号',
  423. title VARCHAR(255) NULL,
  424. jyr1 VARCHAR(255) NULL,
  425. ChangeStatus INT NULL,
  426. status VARCHAR(255) NULL COMMENT '检验状态'
  427. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='过程检验单(头) 归口自旧系统dopdemo,空表待上游抽数'
  428. """);
  429. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_ipqc_inspection"));
  430. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_ipqc_inspection"));
  431. // 归口源明细表(方案 A):同构自旧系统 dopdemo.qms_gcjydzb(100 列),空表,待上游抽数填充。
  432. // 关联键 glid -> qms_gcjyd.id。id 不设 AUTO_INCREMENT:抽数将携带旧系统主键。
  433. await MdpSchemaAligner.ExecuteAsync(_db,
  434. """
  435. CREATE TABLE IF NOT EXISTS qms_gcjydzb (
  436. id BIGINT NOT NULL PRIMARY KEY COMMENT '子表id',
  437. tenant_id BIGINT NOT NULL COMMENT 'tenant_id',
  438. glid BIGINT NULL COMMENT '关联id',
  439. dwsj TEXT NULL, jyxm TEXT NULL, jybz TEXT NULL, jyyj TEXT NULL,
  440. j1 TEXT NULL, j2 TEXT NULL, j3 TEXT NULL, j4 TEXT NULL, j5 TEXT NULL,
  441. ljbh TEXT NULL, sczyry TEXT NULL, fj TEXT NULL, bz TEXT NULL,
  442. j6 TEXT NULL, j7 TEXT NULL, j8 TEXT NULL, j9 TEXT NULL, j10 TEXT NULL,
  443. j11 TEXT NULL, j12 TEXT NULL, j13 TEXT NULL, j14 TEXT NULL, j15 TEXT NULL,
  444. j16 TEXT NULL, j17 TEXT NULL, j18 TEXT NULL, j19 TEXT NULL, j20 TEXT NULL,
  445. j21 TEXT NULL, j22 TEXT NULL, j23 TEXT NULL, j24 TEXT NULL, j25 TEXT NULL,
  446. j26 TEXT NULL, j27 TEXT NULL, j28 TEXT NULL, j29 TEXT NULL, j30 TEXT NULL,
  447. j31 TEXT NULL, j32 TEXT NULL, j33 TEXT NULL, j34 TEXT NULL, j35 TEXT NULL,
  448. j36 TEXT NULL, j37 TEXT NULL, j38 TEXT NULL, j39 TEXT NULL, j40 TEXT NULL,
  449. pd BIGINT NULL COMMENT '结果判定',
  450. jysl DECIMAL(10,0) NULL COMMENT '检验数量',
  451. bhgsl DECIMAL(10,0) NULL COMMENT '不合格数量',
  452. j41 TEXT NULL, j42 TEXT NULL, j43 TEXT NULL, j44 TEXT NULL, j45 TEXT NULL,
  453. j46 TEXT NULL, j47 TEXT NULL, j48 TEXT NULL, j49 TEXT NULL, j50 TEXT NULL,
  454. j51 TEXT NULL, j52 TEXT NULL, j53 TEXT NULL, j54 TEXT NULL, j55 TEXT NULL,
  455. j56 TEXT NULL, j57 TEXT NULL, j58 TEXT NULL, j59 TEXT NULL, j60 TEXT NULL,
  456. j61 TEXT NULL, j62 TEXT NULL, j63 TEXT NULL, j64 TEXT NULL, j65 TEXT NULL,
  457. j66 TEXT NULL, j67 TEXT NULL, j68 TEXT NULL, j69 TEXT NULL, j70 TEXT NULL,
  458. nature TEXT NULL, jsbz TEXT NULL, jypc TEXT NULL, ggbh TEXT NULL,
  459. ProcessCode TEXT NULL, ProcessName TEXT NULL, jyff TEXT NULL,
  460. ybl INT NULL COMMENT '样本量',
  461. NonNumericalTypeOK INT NULL COMMENT '非数值型OK',
  462. NonNumericalTypeNG INT NULL COMMENT '非数值型NG',
  463. Inspector TEXT NULL,
  464. InspectionTime DATETIME NULL COMMENT '检验时间',
  465. jgid BIGINT NULL,
  466. AbnormalNumber TEXT NULL, MeasurementInstrument TEXT NULL, sx TEXT NULL, xx TEXT NULL,
  467. KEY idx_qms_gcjydzb_glid (glid)
  468. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='过程检验单(明细) 归口自旧系统dopdemo,空表待上游抽数'
  469. """);
  470. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_stg_ipqc_inspection_detail"));
  471. await MdpSchemaAligner.ExecuteAsync(_db, MdpSchemaDefinition.Ddl("mdp_std_ipqc_inspection_detail"));
  472. }
  473. /// <summary>贴源头:qms_gcjyd -> mdp_stg_ipqc_inspection(raw_data JSON 快照)。跳过无有效 tenant_id 的行。返回写入行数。</summary>
  474. private async Task<int> SyncHeadStagingAsync(string batchId, DateTime now)
  475. {
  476. var total = await _db.Ado.GetIntAsync("SELECT COUNT(1) FROM qms_gcjyd");
  477. var skipped = await _db.Ado.GetIntAsync("SELECT COUNT(1) FROM qms_gcjyd m WHERE NULLIF(m.tenant_id, 0) IS NULL");
  478. if (skipped > 0)
  479. _tenantSkipWarnings.Add($"qms_gcjyd 贴源跳过 {skipped} 行:无法解析 tenant_id");
  480. await MdpSchemaAligner.ExecuteAsync(_db,
  481. """
  482. INSERT INTO mdp_stg_ipqc_inspection
  483. (tenant_id, factory_id, source_system, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time, process_status, raw_data)
  484. SELECT
  485. NULLIF(m.tenant_id, 0), 1, 'AIDOP_NATIVE', 'qms_gcjyd',
  486. CAST(m.id AS CHAR), CAST(m.djbh AS CHAR), @BatchId, @Now, 'PENDING',
  487. JSON_OBJECT(
  488. 'id', m.id, 'djbh', m.djbh, 'cplx', m.cplx, 'scph', m.scph, 'lydjbh', m.lydjbh,
  489. 'jgpd', m.jgpd, 'fj', CAST(m.fj AS CHAR), 'bz', CAST(m.bz AS CHAR), 'jyr', m.jyr,
  490. 'gxbm', m.gxbm, 'gxmc', m.gxmc, 'sczyry', m.sczyry, 'ybl', m.ybl,
  491. 'bdbh', m.bdbh, 'bbh', m.bbh, 'sxrq', m.sxrq, 'wlbm', m.wlbm, 'wlmc', m.wlmc,
  492. 'jgbb', m.jgbb, 'jgbh', m.jgbh, 'status', m.status
  493. )
  494. FROM qms_gcjyd m
  495. WHERE NULLIF(m.tenant_id, 0) IS NOT NULL
  496. ON DUPLICATE KEY UPDATE
  497. source_biz_key=VALUES(source_biz_key), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  498. process_status=VALUES(process_status), raw_data=VALUES(raw_data), update_time=CURRENT_TIMESTAMP
  499. """,
  500. new SugarParameter("@BatchId", batchId),
  501. new SugarParameter("@Now", now));
  502. return total - skipped;
  503. }
  504. /// <summary>标准化头:mdp_stg_ipqc_pull(qms_gcjyd) → mdp_std_ipqc_inspection。</summary>
  505. private async Task<int> TransformHeadStandardAsync(string batchId, DateTime now, string? sourceSystem = null, IpqcBillScope? scope = null)
  506. {
  507. // 双源:sourceSystem 非空时仅统计/转换当前 source(切源 FULL Replace 用);为空时保持既有全源行为不变。
  508. var countPars = new List<SugarParameter>();
  509. var countSrc = "";
  510. var insertSrc = "";
  511. if (sourceSystem != null)
  512. {
  513. countSrc = " AND source_system=@Src";
  514. insertSrc = " AND m.source_system=@Src";
  515. countPars.Add(new SugarParameter("@Src", sourceSystem));
  516. }
  517. var mTenant = MdpJsonSql.TenantFromStg("m");
  518. // 单据作用域(targeted sync):按 source_row_id 收敛到单张检验单,并强制租户相等。
  519. // 租户显式来自业务单,不依赖 RunFullAsync 的 pbxfxp 账套映射上下文。scope 为 null 时全量行为逐字不变。
  520. if (scope != null)
  521. {
  522. var scopeSql = $" AND m.source_row_id=@ScopeRowId AND ({mTenant})=@ScopeTenant";
  523. countSrc += scopeSql;
  524. insertSrc += scopeSql;
  525. countPars.Add(new SugarParameter("@ScopeRowId", scope.InspectionId.ToString()));
  526. countPars.Add(new SugarParameter("@ScopeTenant", scope.TenantId));
  527. }
  528. var rows = await _db.Ado.GetIntAsync(
  529. $"SELECT COUNT(1) FROM mdp_stg_ipqc_pull m WHERE m.source_table='qms_gcjyd'{countSrc} AND {MdpJsonSql.TenantGuard(mTenant)}", countPars);
  530. var skipped = await _db.Ado.GetIntAsync(
  531. $"SELECT COUNT(1) FROM mdp_stg_ipqc_pull m WHERE m.source_table='qms_gcjyd'{countSrc} AND NOT {MdpJsonSql.TenantGuard(mTenant)}", countPars);
  532. if (skipped > 0)
  533. _tenantSkipWarnings.Add($"mdp_std_ipqc_inspection 跳过 {skipped} 行:无法解析 tenant_id");
  534. var insertSql =
  535. $"""
  536. INSERT INTO mdp_std_ipqc_inspection
  537. (tenant_id, factory_id, source_system, bill_no, product_model, production_batch_no, production_work_order,
  538. result_judgement, attachment, remark, inspector, process_code, process_name, production_person,
  539. sample_qty, form_no, version_no, effective_date, material_code, material_name,
  540. inspec_standard_version, inspec_standard_code, inspection_status,
  541. source_row_id, source_biz_key, sync_batch_id, sync_time)
  542. SELECT
  543. {mTenant}, 1, {MdpSourceIdentity.Resolve("m.source_table", "m.source_system")},
  544. IFNULL({MdpJsonSql.Str("m", "djbh")}, {MdpJsonSql.Str("m", "id")}),
  545. {MdpJsonSql.Str("m", "cplx")},
  546. {MdpJsonSql.Str("m", "scph")},
  547. {MdpJsonSql.Str("m", "lydjbh")},
  548. {MdpJsonSql.Str("m", "jgpd")},
  549. {MdpJsonSql.Str("m", "fj")},
  550. {MdpJsonSql.Str("m", "bz")},
  551. {MdpJsonSql.Str("m", "jyr")},
  552. {MdpJsonSql.Str("m", "gxbm")},
  553. {MdpJsonSql.Str("m", "gxmc")},
  554. {MdpJsonSql.Str("m", "sczyry")},
  555. {MdpJsonSql.Dec("m", "ybl", 18, 6)},
  556. {MdpJsonSql.Str("m", "bdbh")},
  557. {MdpJsonSql.Str("m", "bbh")},
  558. {MdpJsonSql.DateTimeSec("m", "sxrq")},
  559. {MdpJsonSql.Str("m", "wlbm")},
  560. {MdpJsonSql.Str("m", "wlmc")},
  561. {MdpJsonSql.Str("m", "jgbb")},
  562. {MdpJsonSql.Str("m", "jgbh")},
  563. {MdpJsonSql.Str("m", "status")},
  564. IFNULL(m.source_row_id, JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.id'))),
  565. IFNULL(NULLIF(m.source_biz_key,''), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.djbh'))),
  566. @BatchId, @Now
  567. FROM mdp_stg_ipqc_pull m
  568. WHERE m.source_table='qms_gcjyd'{insertSrc} AND {MdpJsonSql.TenantGuard(mTenant)}
  569. ON DUPLICATE KEY UPDATE
  570. bill_no=VALUES(bill_no), product_model=VALUES(product_model), production_batch_no=VALUES(production_batch_no),
  571. production_work_order=VALUES(production_work_order), result_judgement=VALUES(result_judgement),
  572. attachment=VALUES(attachment), remark=VALUES(remark), inspector=VALUES(inspector),
  573. process_code=VALUES(process_code), process_name=VALUES(process_name), production_person=VALUES(production_person),
  574. sample_qty=VALUES(sample_qty), form_no=VALUES(form_no), version_no=VALUES(version_no),
  575. effective_date=VALUES(effective_date), material_code=VALUES(material_code), material_name=VALUES(material_name),
  576. inspec_standard_version=VALUES(inspec_standard_version), inspec_standard_code=VALUES(inspec_standard_code),
  577. inspection_status=VALUES(inspection_status),
  578. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  579. """;
  580. var insPars = new List<SugarParameter>
  581. {
  582. new("@BatchId", batchId),
  583. new("@Now", now)
  584. };
  585. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  586. if (scope != null)
  587. {
  588. insPars.Add(new SugarParameter("@ScopeRowId", scope.InspectionId.ToString()));
  589. insPars.Add(new SugarParameter("@ScopeTenant", scope.TenantId));
  590. }
  591. await MdpSchemaAligner.ExecuteAsync(_db, insertSql, insPars);
  592. return rows;
  593. }
  594. /// <summary>贴源明细:qms_gcjydzb -> mdp_stg_ipqc_inspection_detail(raw_data JSON 快照,确定子集列)。跳过无有效 tenant_id 的行。</summary>
  595. private async Task<int> SyncDetailStagingAsync(string batchId, DateTime now)
  596. {
  597. var total = await _db.Ado.GetIntAsync("SELECT COUNT(1) FROM qms_gcjydzb");
  598. var skipped = await _db.Ado.GetIntAsync("SELECT COUNT(1) FROM qms_gcjydzb d WHERE NULLIF(d.tenant_id, 0) IS NULL");
  599. if (skipped > 0)
  600. _tenantSkipWarnings.Add($"qms_gcjydzb 贴源跳过 {skipped} 行:无法解析 tenant_id");
  601. await MdpSchemaAligner.ExecuteAsync(_db,
  602. """
  603. INSERT INTO mdp_stg_ipqc_inspection_detail
  604. (tenant_id, factory_id, source_system, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time, process_status, raw_data)
  605. SELECT
  606. NULLIF(d.tenant_id, 0), 1, 'AIDOP_NATIVE', 'qms_gcjydzb',
  607. CAST(d.id AS CHAR), CAST(d.glid AS CHAR), @BatchId, @Now, 'PENDING',
  608. JSON_OBJECT(
  609. 'id', d.id, 'glid', d.glid, 'jyxm', CAST(d.jyxm AS CHAR), 'jyyj', CAST(d.jyyj AS CHAR),
  610. 'jsbz', CAST(d.jsbz AS CHAR), 'jyff', CAST(d.jyff AS CHAR), 'sx', CAST(d.sx AS CHAR), 'xx', CAST(d.xx AS CHAR),
  611. 'ProcessCode', CAST(d.ProcessCode AS CHAR), 'ProcessName', CAST(d.ProcessName AS CHAR),
  612. 'MeasurementInstrument', CAST(d.MeasurementInstrument AS CHAR), 'Inspector', CAST(d.Inspector AS CHAR),
  613. 'InspectionTime', d.InspectionTime, 'ybl', d.ybl, 'jysl', d.jysl, 'bhgsl', d.bhgsl,
  614. 'pd', d.pd, 'NonNumericalTypeOK', d.NonNumericalTypeOK, 'NonNumericalTypeNG', d.NonNumericalTypeNG,
  615. 'bz', CAST(d.bz AS CHAR), 'fj', CAST(d.fj AS CHAR)
  616. )
  617. FROM qms_gcjydzb d
  618. WHERE NULLIF(d.tenant_id, 0) IS NOT NULL
  619. ON DUPLICATE KEY UPDATE
  620. source_biz_key=VALUES(source_biz_key), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  621. process_status=VALUES(process_status), raw_data=VALUES(raw_data), update_time=CURRENT_TIMESTAMP
  622. """,
  623. new SugarParameter("@BatchId", batchId),
  624. new SugarParameter("@Now", now));
  625. return total - skipped;
  626. }
  627. /// <summary>
  628. /// 标准化明细:mdp_stg_ipqc_pull(qms_gcjydzb) → mdp_std_ipqc_inspection_detail(回填 std_head_id)。
  629. /// 全 typed 列经 MdpJsonSql 跨源类型兼容(ybl/jysl/bhgsl/NonNumerical* Dec/Int 加 NULL 守卫、InspectionTime 兼容 ISO-T)。
  630. /// sourceSystem 非空时仅统计/转换当前 source(切源 FULL Replace 用);为空时保持既有全源行为。
  631. /// </summary>
  632. private async Task<int> TransformDetailStandardAsync(string batchId, DateTime now, string? sourceSystem = null)
  633. {
  634. var srcClause = sourceSystem == null ? "" : " AND d.source_system=@Src";
  635. var dTenant = MdpJsonSql.TenantFromStg("d");
  636. var countPars = new List<SugarParameter>();
  637. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  638. var rows = await _db.Ado.GetIntAsync(
  639. $"SELECT COUNT(1) FROM mdp_stg_ipqc_pull d WHERE d.source_table='qms_gcjydzb'{srcClause} AND {MdpJsonSql.TenantGuard(dTenant)}", countPars);
  640. var skipped = await _db.Ado.GetIntAsync(
  641. $"SELECT COUNT(1) FROM mdp_stg_ipqc_pull d WHERE d.source_table='qms_gcjydzb'{srcClause} AND NOT {MdpJsonSql.TenantGuard(dTenant)}", countPars);
  642. if (skipped > 0)
  643. _tenantSkipWarnings.Add($"mdp_std_ipqc_inspection_detail 跳过 {skipped} 行:无法解析 tenant_id");
  644. var insertSql =
  645. $"""
  646. INSERT INTO mdp_std_ipqc_inspection_detail
  647. (tenant_id, std_head_id, source_head_id, inspection_item, inspection_basis, technical_standard, inspection_method,
  648. upper_limit, lower_limit, process_code, process_name, measurement_instrument, inspector, inspection_time,
  649. sample_qty, inspection_qty, unqualified_qty, result_judgement, non_numeric_ok, non_numeric_ng, remark, attachment,
  650. source_row_id, source_biz_key, sync_batch_id, sync_time)
  651. SELECT
  652. {dTenant}, h.id,
  653. {MdpJsonSql.Str("d", "glid")},
  654. {MdpJsonSql.Str("d", "jyxm")},
  655. {MdpJsonSql.Str("d", "jyyj")},
  656. {MdpJsonSql.Str("d", "jsbz")},
  657. {MdpJsonSql.Str("d", "jyff")},
  658. {MdpJsonSql.Str("d", "sx")},
  659. {MdpJsonSql.Str("d", "xx")},
  660. {MdpJsonSql.Str("d", "ProcessCode")},
  661. {MdpJsonSql.Str("d", "ProcessName")},
  662. {MdpJsonSql.Str("d", "MeasurementInstrument")},
  663. {MdpJsonSql.Str("d", "Inspector")},
  664. {MdpJsonSql.DateTimeSec("d", "InspectionTime")},
  665. {MdpJsonSql.Int("d", "ybl")},
  666. {MdpJsonSql.Dec("d", "jysl", 18, 6)},
  667. {MdpJsonSql.Dec("d", "bhgsl", 18, 6)},
  668. {MdpJsonSql.Str("d", "pd")},
  669. {MdpJsonSql.Int("d", "NonNumericalTypeOK")},
  670. {MdpJsonSql.Int("d", "NonNumericalTypeNG")},
  671. {MdpJsonSql.Str("d", "bz")},
  672. {MdpJsonSql.Str("d", "fj")},
  673. IFNULL(d.source_row_id, {MdpJsonSql.Str("d", "id")}),
  674. IFNULL(NULLIF(d.source_biz_key,''), {MdpJsonSql.Str("d", "glid")}),
  675. @BatchId, @Now
  676. FROM mdp_stg_ipqc_pull d
  677. LEFT JOIN mdp_std_ipqc_inspection h
  678. ON h.tenant_id = {dTenant}
  679. AND h.source_row_id = {MdpJsonSql.Str("d", "glid")}
  680. WHERE d.source_table='qms_gcjydzb'{srcClause} AND {MdpJsonSql.TenantGuard(dTenant)}
  681. ON DUPLICATE KEY UPDATE
  682. std_head_id=VALUES(std_head_id), source_head_id=VALUES(source_head_id), inspection_item=VALUES(inspection_item),
  683. inspection_basis=VALUES(inspection_basis), technical_standard=VALUES(technical_standard), inspection_method=VALUES(inspection_method),
  684. upper_limit=VALUES(upper_limit), lower_limit=VALUES(lower_limit), process_code=VALUES(process_code), process_name=VALUES(process_name),
  685. measurement_instrument=VALUES(measurement_instrument), inspector=VALUES(inspector), inspection_time=VALUES(inspection_time),
  686. sample_qty=VALUES(sample_qty), inspection_qty=VALUES(inspection_qty), unqualified_qty=VALUES(unqualified_qty),
  687. result_judgement=VALUES(result_judgement), non_numeric_ok=VALUES(non_numeric_ok), non_numeric_ng=VALUES(non_numeric_ng),
  688. remark=VALUES(remark), attachment=VALUES(attachment),
  689. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  690. """;
  691. var insPars = new List<SugarParameter> { new("@BatchId", batchId), new("@Now", now) };
  692. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  693. await MdpSchemaAligner.ExecuteAsync(_db, insertSql, insPars);
  694. return rows;
  695. }
  696. private async Task<long> InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  697. {
  698. await MdpSchemaAligner.ExecuteAsync(_db,
  699. """
  700. INSERT INTO mdp_transform_run_log
  701. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  702. VALUES (0, @JobCode, 'S6过程检验单MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  703. """,
  704. new SugarParameter("@JobCode", JobCode),
  705. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  706. new SugarParameter("@BatchId", batchId),
  707. new SugarParameter("@StartTime", startedAt));
  708. return await _db.Ado.GetLongAsync(
  709. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  710. new List<SugarParameter> { new("@BatchId", batchId) });
  711. }
  712. private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, IpqcInspectionMdpSyncResult result)
  713. {
  714. var finishedAt = DateTime.Now;
  715. await MdpSchemaAligner.ExecuteAsync(_db,
  716. """
  717. UPDATE mdp_transform_run_log
  718. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  719. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP
  720. WHERE id=@Id
  721. """,
  722. new SugarParameter("@EndTime", finishedAt),
  723. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  724. new SugarParameter("@StageRows", result.HeadStageRows + result.DetailStageRows),
  725. new SugarParameter("@StandardRows", result.HeadStandardRows + result.DetailStandardRows),
  726. new SugarParameter("@Id", runLogId));
  727. }
  728. private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message)
  729. {
  730. var finishedAt = DateTime.Now;
  731. await MdpSchemaAligner.ExecuteAsync(_db,
  732. """
  733. UPDATE mdp_transform_run_log
  734. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  735. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  736. WHERE id=@Id
  737. """,
  738. new SugarParameter("@EndTime", finishedAt),
  739. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  740. new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message),
  741. new SugarParameter("@Id", runLogId));
  742. }
  743. private static string NormalizeTriggerType(string? triggerType)
  744. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  745. }
  746. /// <summary>过程检验单 MDP 同步转换结果。</summary>
  747. public sealed class IpqcInspectionMdpSyncResult
  748. {
  749. public long RunLogId { get; set; }
  750. public string BatchId { get; set; } = string.Empty;
  751. public int HeadStageRows { get; set; }
  752. public int HeadStandardRows { get; set; }
  753. public int DetailStageRows { get; set; }
  754. public int DetailStandardRows { get; set; }
  755. public List<string> TenantSkipWarnings { get; set; } = new();
  756. }
  757. /// <summary>S6 IPQC 双模式入站结果。</summary>
  758. public sealed class IpqcInspectionInboundResult
  759. {
  760. public string PullBatchId { get; set; } = string.Empty;
  761. public int RowsPulled { get; set; }
  762. public int RowsWrittenStg { get; set; }
  763. public string? NewCursor { get; set; }
  764. public string TransformBatchId { get; set; } = string.Empty;
  765. public int HeadStageRows { get; set; }
  766. public int HeadStandardRows { get; set; }
  767. public int DetailStageRows { get; set; }
  768. public int DetailStandardRows { get; set; }
  769. }