FqcMdpSyncService.cs 48 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814
  1. using System.Linq;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform;
  3. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  4. using Admin.NET.Plugin.AiDOP.Infrastructure;
  5. namespace Admin.NET.Plugin.AiDOP.FinishedWarehouse;
  6. /// <summary>
  7. /// S7 FQC 成品质量 数据中台只读同步转换服务(DOP 内部,独立于 KPI 管线)。
  8. ///
  9. /// B-①/C 迁通用管线(CTO 拍板:可接受批次快照延迟、全迁 std):
  10. /// 源(qms_fqcbj 任务 / qms_qcpp_inspbill 结果头 / qms_qcpp_inspbillst 明细) 经执行器灌 mdp_stg_fqc_pull(三表共表,source_table 区分)
  11. /// → transform 读 pull stg(MdpJsonSql 跨源类型兼容)→ mdp_std_fqc_task / mdp_std_fqc_result / mdp_std_fqc_result_detail。
  12. /// result 富化 SAP 工单号:本库 WorkOrdMaster 相关子查询 MIN(WorkOrd) WHERE Batch=sczld(维表本库,不入 stg)。
  13. /// result_detail 回填 std_result_id:glid → mdp_std_fqc_result.source_row_id。样本 j1..j50 逐列落 std(utf8mb3 控行宽)。
  14. ///
  15. /// 双源:本库实体 S7_FQC_TASK/RESULT/RESULT_DETAIL(AIDOPDEV_MYSQL, status=1);
  16. /// dopdemorq 第二源 *_SQLSERVER(DOPDEMORQ_SQLSERVER, status=0 就位不启用)。qms_* 无 UpdateTime → FULL。
  17. ///
  18. /// 约束:只读源/贴源,仅写 mdp_stg_fqc_pull / mdp_std_fqc_*;绝不 INSERT/UPDATE/DELETE 业务源表;不接 ApprovalFlow、不做录入/审核/处置。
  19. /// tenant_id 从 raw_data 保留(源含该列 → std 租户不回归;SQL Server 源无此列 → 回落入站上下文)。
  20. /// 源为 0 行时转换成功完成、处理数 0,不报错。切源 FULL Replace 三表各一次单事务、仅按 tenant、Pull 成功后才 destructive。
  21. /// </summary>
  22. public class FqcMdpSyncService : ITransient
  23. {
  24. private const string JobCode = "S7_FQC_MDP_SYNC";
  25. private const string InboundTaskEntityCode = "S7_FQC_TASK";
  26. private const string InboundResultEntityCode = "S7_FQC_RESULT";
  27. private const string InboundDetailEntityCode = "S7_FQC_RESULT_DETAIL";
  28. private const string SqlServerSourceCode = "DOPDEMORQ_SQLSERVER";
  29. private const string SqlServerTaskEntityCode = "S7_FQC_TASK_SQLSERVER";
  30. private const string SqlServerResultEntityCode = "S7_FQC_RESULT_SQLSERVER";
  31. private const string SqlServerDetailEntityCode = "S7_FQC_RESULT_DETAIL_SQLSERVER";
  32. /// <summary>165 SQL Server 源无 tenant_id 列;切源时须经账套映射解析。</summary>
  33. private const string SourceZtid = "pbxfxp";
  34. private const string StdTaskTable = "mdp_std_fqc_task";
  35. private const string StdResultTable = "mdp_std_fqc_result";
  36. private const string StdResultDetailTable = "mdp_std_fqc_result_detail";
  37. private readonly ISqlSugarClient _db;
  38. private readonly MdpSourcePullDispatcher _pullDispatcher;
  39. public FqcMdpSyncService(ISqlSugarClient db, MdpSourcePullDispatcher pullDispatcher)
  40. {
  41. _db = db;
  42. _pullDispatcher = pullDispatcher;
  43. }
  44. /// <summary>全量:本地 DB 执行器灌 pull stg(任务/结果/明细) → 标准层三表(读 stg)。</summary>
  45. public async Task<FqcMdpSyncResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
  46. {
  47. cancellationToken.ThrowIfCancellationRequested();
  48. await EnsureTablesAsync();
  49. await EnsurePullStgTableAsync();
  50. // 1E-0:并发收口——同库同时只允许一个 FQC 全量刷新(多实例/多次触发串行化,抢不到则跳过,交下一轮)
  51. var locked = await _db.Ado.GetIntAsync("SELECT GET_LOCK(@k, 15)", new List<SugarParameter> { new("@k", FqcRefreshLockKey) });
  52. if (locked != 1)
  53. return new FqcMdpSyncResult { BatchId = "SKIPPED_LOCK_BUSY", Skipped = true };
  54. var now = DateTime.Now;
  55. var batchId = $"S7_FQC_FULL_{now:yyyyMMddHHmmss}";
  56. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  57. var result = new FqcMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  58. try
  59. {
  60. // 本库 qms_* 含 tenant_id;ctx tenantId=0 仅兜底,Writer 优先取源行值。
  61. var pullCtx = new MdpPullContext
  62. {
  63. TenantId = 0,
  64. FullRefresh = true,
  65. TaskCode = "S7_FQC_INBOUND",
  66. BatchId = $"{batchId}_PULL"
  67. };
  68. var pull = await PopulateStgAsync(pullCtx, cancellationToken);
  69. result.StageRows = pull.RowsWritten;
  70. result.TaskRows = await TransformTaskStandardAsync(batchId, now);
  71. result.ResultRows = await TransformResultStandardAsync(batchId, now);
  72. result.DetailRows = await TransformResultDetailStandardAsync(batchId, now);
  73. // 1F:全量刷新后清理本库(AIDOP)贴源 orphan——源单据已删的 std 行(transform 只 upsert 不删,
  74. // pull stg 跨全量批不清致旧会话行残留 → std 累积孤儿 → 结果列表虚增行)。仅清 AIDOP 源,
  75. // SQLSERVER 双源(status=0 就位不启用)行按 source_system 排除,绝不误删。
  76. await PurgeStdOrphansAsync();
  77. await MarkRunSuccessAsync(runLogId, now, result);
  78. return result;
  79. }
  80. catch (Exception ex)
  81. {
  82. await MarkRunFailedAsync(runLogId, now, ex.Message);
  83. throw;
  84. }
  85. finally
  86. {
  87. await _db.Ado.ExecuteCommandAsync("SELECT RELEASE_LOCK(@k)", new List<SugarParameter> { new("@k", FqcRefreshLockKey) });
  88. }
  89. }
  90. /// <summary>FQC 全量刷新并发锁键(同库同时只跑一个)。</summary>
  91. private const string FqcRefreshLockKey = "aidop:s7:fqc:refresh";
  92. /// <summary>本库 MySQL 贴源 source_system(与 MdpDbPullExecutor 落 stg 的取值一致,决定 uk_source_key 归属)。</summary>
  93. private const string NativeSourceSystem = "AIDOPDEV_MYSQL";
  94. /// <summary>targeted sync 的单据作用域(租户 + 检验单 id + 报检单号)。</summary>
  95. private sealed class FqcBillScope
  96. {
  97. public long TenantId { get; init; }
  98. public long BillId { get; init; }
  99. /// <summary>qms_qcpp_inspbill.lydjbh = qms_fqcbj.FBILLNO;为空则跳过任务表同步。</summary>
  100. public string? SourceBillNo { get; init; }
  101. }
  102. /// <summary>把作用域参数补进 transform 的参数表(未启用作用域时不加,保持全量 SQL 原样)。</summary>
  103. private static void AddScopePars(List<SugarParameter> pars, FqcBillScope? scope)
  104. {
  105. if (scope == null) return;
  106. pars.Add(new SugarParameter("@ScopeTenant", scope.TenantId));
  107. pars.Add(new SugarParameter("@ScopeBillId", scope.BillId.ToString()));
  108. pars.Add(new SugarParameter("@ScopeBillNo", scope.SourceBillNo));
  109. }
  110. /// <summary>
  111. /// 单据级 targeted 同步(业务写入后立即可见,替代把 12s 全量 <see cref="RunFullAsync"/> 挂到高频提交路径)。
  112. ///
  113. /// 只处理 billId 这一张检验单:源 → stg(贴源 upsert)→ std 三表(复用全量同一套 transform,字段映射零分叉)。
  114. /// qms_qcpp_inspbill (id=billId) → mdp_std_fqc_result
  115. /// qms_qcpp_inspbillst(glid=billId) → mdp_std_fqc_result_detail
  116. /// qms_fqcbj (FBILLNO=inspbill.lydjbh) → mdp_std_fqc_task(submit-result/退回会改 FINSPECTSTATUS 等进度列)
  117. ///
  118. /// 边界:只读业务源、只写 mdp_stg_fqc_pull / mdp_std_fqc_*;全程 tenant_id 强约束(源查询与 transform 双重);
  119. /// 纯 upsert,不 DELETE、不跑 orphan purge、不抢全量 GET_LOCK、不产生 S7_FQC_FULL 批次。
  120. /// </summary>
  121. public async Task<FqcMdpSyncResult> SyncBillAsync(
  122. long billId, long tenantId, string triggerType = "BILL", CancellationToken cancellationToken = default)
  123. {
  124. cancellationToken.ThrowIfCancellationRequested();
  125. if (billId <= 0 || tenantId <= 0) return new FqcMdpSyncResult { BatchId = "SKIPPED_INVALID_SCOPE", Skipped = true };
  126. var now = DateTime.Now;
  127. var batchId = $"S7_FQC_BILL_{now:yyyyMMddHHmmssfff}_{billId}";
  128. var runLogId = await InsertRunLogAsync(batchId, now, triggerType);
  129. var result = new FqcMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  130. try
  131. {
  132. // 来源单号:决定是否需要同步任务表;同时用租户约束再证一次单据归属
  133. var sourceBillNo = await _db.Ado.SqlQuerySingleAsync<string>(
  134. "SELECT lydjbh FROM qms_qcpp_inspbill WHERE id=@id AND tenant_id=@tid LIMIT 1",
  135. new SugarParameter("@id", billId), new SugarParameter("@tid", tenantId));
  136. var scope = new FqcBillScope { TenantId = tenantId, BillId = billId, SourceBillNo = sourceBillNo };
  137. // ① 贴源:把该单据的源行刷进 stg(transform 一律读 stg,故必须先刷,否则同步的还是旧快照)
  138. result.StageRows = await UpsertStgFromSourceAsync(
  139. "qms_qcpp_inspbill", "s.id=@id AND s.tenant_id=@tid",
  140. new List<SugarParameter> { new("@id", billId), new("@tid", tenantId) }, batchId, now);
  141. result.StageRows += await UpsertStgFromSourceAsync(
  142. "qms_qcpp_inspbillst", "s.glid=@id AND s.tenant_id=@tid",
  143. new List<SugarParameter> { new("@id", billId), new("@tid", tenantId) }, batchId, now);
  144. if (!string.IsNullOrWhiteSpace(sourceBillNo))
  145. result.StageRows += await UpsertStgFromSourceAsync(
  146. "qms_fqcbj", "s.FBILLNO=@bjbh AND s.tenant_id=@tid",
  147. new List<SugarParameter> { new("@bjbh", sourceBillNo), new("@tid", tenantId) }, batchId, now);
  148. // ② 标准层:复用全量 transform,仅加单据作用域。result 必须先于 detail(detail 回填 std_result_id)
  149. result.ResultRows = await TransformResultStandardAsync(batchId, now, null, scope);
  150. result.DetailRows = await TransformResultDetailStandardAsync(batchId, now, null, scope);
  151. if (!string.IsNullOrWhiteSpace(sourceBillNo))
  152. result.TaskRows = await TransformTaskStandardAsync(batchId, now, null, scope);
  153. await MarkRunSuccessAsync(runLogId, now, result);
  154. return result;
  155. }
  156. catch (Exception ex)
  157. {
  158. await MarkRunFailedAsync(runLogId, now, ex.Message);
  159. throw;
  160. }
  161. }
  162. /// <summary>
  163. /// 单表贴源 upsert:源行 → mdp_stg_fqc_pull,raw_data 按源表<b>全列</b>动态生成 JSON_OBJECT,
  164. /// 与 MdpDbPullExecutor 的信封形状一致(source_row_id=source_biz_key=主键,JSON 键=源列名),
  165. /// 因此绝不会因为 transform 未来新增字段而漂移。命中既有 uk_source_key 则原地更新,不产生重复 stg 行。
  166. /// </summary>
  167. private async Task<int> UpsertStgFromSourceAsync(
  168. string sourceTable, string whereSql, List<SugarParameter> pars, string batchId, DateTime now)
  169. {
  170. var cols = await _db.Ado.SqlQueryAsync<string>(
  171. "SELECT COLUMN_NAME FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t ORDER BY ORDINAL_POSITION",
  172. new List<SugarParameter> { new("@t", sourceTable) });
  173. // 列名来自 information_schema 且强制标识符白名单,非用户输入,无注入面
  174. cols = cols.Where(c => !string.IsNullOrWhiteSpace(c) && System.Text.RegularExpressions.Regex.IsMatch(c, "^[A-Za-z0-9_]+$")).ToList();
  175. if (cols.Count == 0) return 0;
  176. var jsonObj = "JSON_OBJECT(" + string.Join(", ", cols.Select(c => $"'{c}', s.`{c}`")) + ")";
  177. var allPars = new List<SugarParameter>(pars)
  178. {
  179. new("@BatchId", batchId),
  180. new("@Now", now),
  181. new("@SrcSys", NativeSourceSystem),
  182. new("@SrcTab", sourceTable),
  183. };
  184. return await _db.Ado.ExecuteCommandAsync(
  185. $"""
  186. INSERT INTO mdp_stg_fqc_pull
  187. (tenant_id, source_system, source_table, source_row_id, source_biz_key, raw_data, sync_batch_id, sync_time, process_status)
  188. SELECT s.tenant_id, @SrcSys, @SrcTab, CAST(s.id AS CHAR), CAST(s.id AS CHAR), {jsonObj}, @BatchId, @Now, 'PENDING'
  189. FROM `{sourceTable}` s
  190. WHERE {whereSql}
  191. ON DUPLICATE KEY UPDATE
  192. tenant_id=VALUES(tenant_id), source_row_id=VALUES(source_row_id), raw_data=VALUES(raw_data),
  193. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  194. """,
  195. allPars);
  196. }
  197. /// <summary>
  198. /// 全量刷新后清理本库(AIDOP)贴源 orphan:源单据(qms_qcpp_inspbill / qms_fqcbj / qms_qcpp_inspbillst)
  199. /// 已删除但 std 层仍残留的行。仅清 source_system='AIDOP'(含历史 NULL/空)的行,SQLSERVER 双源行不受影响;
  200. /// detail 无 source_system 列,按父 result(glid=result.source_row_id)是否存在判定,与源系统无关且安全。
  201. /// </summary>
  202. private async Task PurgeStdOrphansAsync()
  203. {
  204. // 仅限本库 MySQL 源(source_system=AIDOPDEV_MYSQL, mdp_source.source_code);DEMO 合成源、
  205. // DOPDEMORQ_SQLSERVER 双源均按 source_system 排除,绝不误删。
  206. const string nativeScope = "s.source_system = 'AIDOPDEV_MYSQL'";
  207. // 结果头:源 qms_qcpp_inspbill(source_row_id=inspbill.id)
  208. await _db.Ado.ExecuteCommandAsync(
  209. $"""
  210. DELETE s FROM {StdResultTable} s
  211. LEFT JOIN qms_qcpp_inspbill q ON q.id = CAST(s.source_row_id AS UNSIGNED) AND q.tenant_id = s.tenant_id
  212. WHERE {nativeScope} AND q.id IS NULL
  213. """);
  214. // 任务:源 qms_fqcbj(source_row_id=fqcbj.id)
  215. await _db.Ado.ExecuteCommandAsync(
  216. $"""
  217. DELETE s FROM {StdTaskTable} s
  218. LEFT JOIN qms_fqcbj q ON q.id = CAST(s.source_row_id AS UNSIGNED) AND q.tenant_id = s.tenant_id
  219. WHERE {nativeScope} AND q.id IS NULL
  220. """);
  221. // 明细:本库源 qms_qcpp_inspbillst(source_row_id=inspbillst.id, 数值型);非本库(DEMO 的 source_row_id
  222. // 形如 DEMO-FQC-SRC-*)由数值正则排除,绝不误删。父结果被上面清理后本库孤儿明细亦随源缺失被清。
  223. await _db.Ado.ExecuteCommandAsync(
  224. $"""
  225. DELETE d FROM {StdResultDetailTable} d
  226. LEFT JOIN qms_qcpp_inspbillst q ON q.id = CAST(d.source_row_id AS UNSIGNED) AND q.tenant_id = d.tenant_id
  227. WHERE d.source_row_id REGEXP '^[0-9]+$' AND q.id IS NULL
  228. """);
  229. }
  230. /// <summary>双模式入站:执行器抽三表 → pull stg,再跑标准层三表(读 stg)。</summary>
  231. public async Task<FqcInboundResult> RunInboundAsync(
  232. long tenantId = 0,
  233. bool fullRefresh = false,
  234. CancellationToken cancellationToken = default)
  235. {
  236. cancellationToken.ThrowIfCancellationRequested();
  237. await EnsureTablesAsync();
  238. await EnsurePullStgTableAsync();
  239. var now = DateTime.Now;
  240. // 入参 tenantId:本库源含 tenant_id 时可传 0(Writer 优先取源行)。
  241. var pullCtx = new MdpPullContext
  242. {
  243. TenantId = tenantId,
  244. FullRefresh = fullRefresh,
  245. TaskCode = "S7_FQC_INBOUND",
  246. BatchId = $"S7_FQC_IN_{now:yyyyMMddHHmmss}"
  247. };
  248. var pull = await PopulateStgAsync(pullCtx, cancellationToken);
  249. var batchId = $"S7_FQC_STD_{now:yyyyMMddHHmmss}";
  250. var runLogId = await InsertRunLogAsync(batchId, now, "INBOUND");
  251. var result = new FqcMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  252. try
  253. {
  254. result.TaskRows = await TransformTaskStandardAsync(batchId, now);
  255. result.ResultRows = await TransformResultStandardAsync(batchId, now);
  256. result.DetailRows = await TransformResultDetailStandardAsync(batchId, now);
  257. await MarkRunSuccessAsync(runLogId, now, result);
  258. }
  259. catch (Exception ex)
  260. {
  261. await MarkRunFailedAsync(runLogId, now, ex.Message);
  262. throw;
  263. }
  264. return new FqcInboundResult
  265. {
  266. PullBatchId = pullCtx.BatchId,
  267. RowsPulled = pull.RowsPulled,
  268. RowsWrittenStg = pull.RowsWritten,
  269. NewCursor = pull.NewCursor,
  270. TransformBatchId = batchId,
  271. TaskRows = result.TaskRows,
  272. ResultRows = result.ResultRows,
  273. DetailRows = result.DetailRows
  274. };
  275. }
  276. /// <summary>
  277. /// 双源切换 FULL Replace(三表):从指定 DB 源(默认 DOPDEMORQ_SQLSERVER)PullAll 全量灌 pull stg,
  278. /// 成功后 task/result/result_detail 各一次单事务 FULL 重建(按 tenant 精确隔离)。
  279. /// 顺序:result 先(写 std 结果,供明细回填 std_result_id),再 result_detail;task 独立。
  280. /// SQLSERVER 实体默认 status=0(就位不启用);dopdemorq 三表当前空,未实际执行 destructive 切换。
  281. /// </summary>
  282. public async Task<FqcInboundResult> RunSourceSwitchFullAsync(
  283. string sourceCode = SqlServerSourceCode,
  284. string taskEntityCode = SqlServerTaskEntityCode,
  285. string resultEntityCode = SqlServerResultEntityCode,
  286. string detailEntityCode = SqlServerDetailEntityCode,
  287. long tenantId = 0,
  288. CancellationToken cancellationToken = default)
  289. {
  290. cancellationToken.ThrowIfCancellationRequested();
  291. await EnsureTablesAsync();
  292. await EnsurePullStgTableAsync();
  293. // 165 SQL Server 源无 tenant_id;显式 tenantId≤0 时经账套映射解析。
  294. tenantId = AidopSourceTenantMap.ResolveTenantId(SourceZtid, tenantId);
  295. var now = DateTime.Now;
  296. var pullCtx = new MdpPullContext
  297. {
  298. TenantId = tenantId,
  299. FullRefresh = true,
  300. TaskCode = "S7_FQC_INBOUND",
  301. BatchId = $"S7_FQC_SW_{now:yyyyMMddHHmmss}"
  302. };
  303. var task = await _pullDispatcher.PullAllByEntityCodeAsync(taskEntityCode, pullCtx, cancellationToken);
  304. var res = await _pullDispatcher.PullAllByEntityCodeAsync(resultEntityCode, pullCtx, cancellationToken);
  305. var det = await _pullDispatcher.PullAllByEntityCodeAsync(detailEntityCode, pullCtx, cancellationToken);
  306. var batchId = $"S7_FQC_SWSTD_{now:yyyyMMddHHmmss}";
  307. var runLogId = await InsertRunLogAsync(batchId, now, "SOURCE_SWITCH");
  308. var result = new FqcMdpSyncResult { BatchId = batchId, RunLogId = runLogId };
  309. try
  310. {
  311. result.TaskRows = await MdpStdFullReplace.ReplaceAsync(
  312. _db, StdTaskTable, tenantId, extraWhere: null,
  313. insertScopedAsync: () => TransformTaskStandardAsync(batchId, now, sourceCode), cancellationToken);
  314. result.ResultRows = await MdpStdFullReplace.ReplaceAsync(
  315. _db, StdResultTable, tenantId, extraWhere: null,
  316. insertScopedAsync: () => TransformResultStandardAsync(batchId, now, sourceCode), cancellationToken);
  317. result.DetailRows = await MdpStdFullReplace.ReplaceAsync(
  318. _db, StdResultDetailTable, tenantId, extraWhere: null,
  319. insertScopedAsync: () => TransformResultDetailStandardAsync(batchId, now, sourceCode), cancellationToken);
  320. await MarkRunSuccessAsync(runLogId, now, result);
  321. }
  322. catch (Exception ex)
  323. {
  324. await MarkRunFailedAsync(runLogId, now, ex.Message);
  325. throw;
  326. }
  327. return new FqcInboundResult
  328. {
  329. PullBatchId = pullCtx.BatchId,
  330. RowsPulled = task.RowsPulled + res.RowsPulled + det.RowsPulled,
  331. RowsWrittenStg = task.RowsWritten + res.RowsWritten + det.RowsWritten,
  332. NewCursor = det.NewCursor ?? res.NewCursor ?? task.NewCursor,
  333. TransformBatchId = batchId,
  334. TaskRows = result.TaskRows,
  335. ResultRows = result.ResultRows,
  336. DetailRows = result.DetailRows
  337. };
  338. }
  339. private async Task<(int RowsPulled, int RowsWritten, string? NewCursor)> PopulateStgAsync(
  340. MdpPullContext pullCtx, CancellationToken cancellationToken)
  341. {
  342. // 1E-0:改分页抽尽,消除单页 5000 行截断(源表 > batch_size 时单页只灌前 5000 → std 缺尾)
  343. var task = await _pullDispatcher.PullAllByEntityCodeAsync(InboundTaskEntityCode, pullCtx, cancellationToken);
  344. var res = await _pullDispatcher.PullAllByEntityCodeAsync(InboundResultEntityCode, pullCtx, cancellationToken);
  345. var det = await _pullDispatcher.PullAllByEntityCodeAsync(InboundDetailEntityCode, pullCtx, cancellationToken);
  346. return (
  347. task.RowsPulled + res.RowsPulled + det.RowsPulled,
  348. task.RowsWritten + res.RowsWritten + det.RowsWritten,
  349. det.NewCursor ?? res.NewCursor ?? task.NewCursor);
  350. }
  351. private async Task EnsurePullStgTableAsync()
  352. {
  353. await _db.Ado.ExecuteCommandAsync(
  354. """
  355. CREATE TABLE IF NOT EXISTS mdp_stg_fqc_pull (
  356. id BIGINT PRIMARY KEY AUTO_INCREMENT,
  357. tenant_id BIGINT NOT NULL,
  358. source_system VARCHAR(50) NULL,
  359. source_table VARCHAR(200),
  360. source_row_id VARCHAR(200),
  361. source_biz_key VARCHAR(300) NULL,
  362. raw_data JSON,
  363. sync_batch_id VARCHAR(100),
  364. sync_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP,
  365. process_status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
  366. process_message VARCHAR(500) NULL,
  367. create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
  368. update_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  369. UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key),
  370. KEY idx_batch (sync_batch_id),
  371. KEY idx_src (source_table, source_row_id),
  372. KEY idx_tenant (tenant_id)
  373. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='S7 FQC 执行器贴源层(任务/结果/明细共表,source_table 区分)'
  374. """);
  375. }
  376. /// <summary>
  377. /// 防御式建表(幂等):归口源空表(qms_fqcbj/qms_qcpp_inspbillst,方案 A,供执行器贴源) + 任务 std + 结果 std + 明细 std(j1..j50 utf8mb3 控行宽)。
  378. /// qms_qcpp_inspbill(结果头) 由既有 FQC 批次建,沿用原 Service 假设已存在。
  379. /// </summary>
  380. private async Task EnsureTablesAsync()
  381. {
  382. // 归口源空表(方案 A,DDL 例外、幂等、零 DML):同构自旧系统 dopdemo,待上游抽数;供执行器 pull 贴源。
  383. await _db.Ado.ExecuteCommandAsync(
  384. """
  385. CREATE TABLE IF NOT EXISTS qms_fqcbj (
  386. id BIGINT NOT NULL PRIMARY KEY COMMENT '主表id',
  387. tenant_id BIGINT NOT NULL COMMENT 'tenant_id WP0-D4',
  388. FBILLNO VARCHAR(80) NULL COMMENT '单据编号',
  389. FBILLSTATUS VARCHAR(50) NULL COMMENT '单据状态',
  390. FCREATORID BIGINT NULL, FMODIFIERID BIGINT NULL, FAUDITORID BIGINT NULL,
  391. FAUDITDATE DATETIME NULL, FMODIFYTIME DATETIME NULL, FCREATETIME DATETIME NULL,
  392. FORGID BIGINT NULL COMMENT '申请组织', FQUALITYORG BIGINT NULL,
  393. FAPPLYUSER VARCHAR(50) NULL COMMENT '申请人', FAPPLYTIME DATETIME NULL COMMENT '申请时间',
  394. FBILLTYPE VARCHAR(50) NULL COMMENT '单据类型', FBIZTYPE BIGINT NULL COMMENT '业务类型',
  395. FINSPECORGID BIGINT NULL, FCOMMENT TEXT NULL COMMENT '备注', FAUTHORIZEOBJID BIGINT NULL, FINTERFACEID VARCHAR(50) NULL,
  396. cpmc VARCHAR(255) NULL COMMENT '产品名称', cpxh VARCHAR(255) NULL COMMENT '产品型号',
  397. scph VARCHAR(255) NULL COMMENT '生产批号', sl DECIMAL(10,0) NULL COMMENT '数量',
  398. jfbh TEXT NULL COMMENT '击发编号', wlbm VARCHAR(255) NULL COMMENT '物料编码',
  399. jyfzr VARCHAR(255) NULL COMMENT '检验负责人', yxj VARCHAR(255) NULL COMMENT '检验优先级',
  400. jykssj DATETIME NULL COMMENT '检验开始时间', jywcsj DATETIME NULL COMMENT '检验完成时间',
  401. mjpc VARCHAR(255) NULL, tsyq VARCHAR(255) NULL,
  402. FINSPECTSTATUS VARCHAR(255) NULL COMMENT '检验进度', sczld VARCHAR(255) NULL COMMENT '生产指令单',
  403. sczldsl DECIMAL(10,0) NULL COMMENT '指令单数量', lydjbh VARCHAR(100) NULL COMMENT '来源单据编号',
  404. ISscrap INT NULL, CirculationCard TEXT NULL
  405. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='FQC检验任务(报检) 归口自旧系统dopdemo,空表待上游抽数'
  406. """);
  407. var jGkDdl = string.Join(", ", Enumerable.Range(1, 50).Select(i => $"j{i} VARCHAR(255) NULL"));
  408. await _db.Ado.ExecuteCommandAsync(
  409. $"""
  410. CREATE TABLE IF NOT EXISTS qms_qcpp_inspbillst (
  411. id BIGINT NOT NULL PRIMARY KEY COMMENT '子表id',
  412. tenant_id BIGINT NOT NULL COMMENT 'tenant_id WP0-D4',
  413. glid BIGINT NULL COMMENT '关联id(=qms_qcpp_inspbill.id)',
  414. jyxm VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '检验项目',
  415. sx VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '上限',
  416. xx VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '下限',
  417. bz VARCHAR(255) NULL COMMENT '备注', czf VARCHAR(255) NULL,
  418. {jGkDdl},
  419. pd VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '判定',
  420. ybl VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '样本量',
  421. txfl BIGINT NULL COMMENT '指标分类', bztk TEXT NULL, jsyq TEXT NULL COMMENT '技术要求',
  422. jyyqsbgj TEXT NULL COMMENT '检验仪器设备工具', jyff TEXT NULL, jytk TEXT NULL,
  423. jsbz TEXT NULL COMMENT '接收标准', cyfa TEXT NULL COMMENT '抽样方案', jl TEXT NULL,
  424. jlqjbh VARCHAR(255) NULL COMMENT '计量器具编号', jlbh VARCHAR(255) NULL COMMENT '计量器具编号',
  425. jf1 VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '击发1',
  426. jf2 VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '击发2',
  427. jf3 VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '击发3',
  428. jf4 VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '击发4',
  429. jf5 VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '击发5',
  430. NonNumericalTypeOK BIGINT NULL COMMENT '非数值OK', NonNumericalTypeNG BIGINT NULL COMMENT '非数值NG',
  431. ThresholdSwitch VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '自动判定按钮',
  432. ThresholdSwitchInt BIGINT NULL COMMENT '自动判定按钮',
  433. SizeSpecifications VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL COMMENT '规范标准',
  434. ybl1 VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL,
  435. xh INT NULL COMMENT '序号',
  436. KEY idx_qms_qcpp_inspbillst_glid (glid)
  437. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='FQC检验单明细(品检项目) 归口自旧系统dopdemo,空表待上游抽数'
  438. """);
  439. await _db.Ado.ExecuteCommandAsync(
  440. """
  441. CREATE TABLE IF NOT EXISTS mdp_std_fqc_task (
  442. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  443. tenant_id BIGINT NOT NULL DEFAULT 0,
  444. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  445. bill_no VARCHAR(80) NULL,
  446. production_order_no VARCHAR(255) NULL,
  447. apply_org_id BIGINT NULL,
  448. apply_time DATETIME NULL,
  449. material_code VARCHAR(255) NULL,
  450. product_name VARCHAR(255) NULL,
  451. product_model VARCHAR(255) NULL,
  452. production_batch_no VARCHAR(255) NULL,
  453. order_qty DECIMAL(18,6) NULL,
  454. qty DECIMAL(18,6) NULL,
  455. applicant VARCHAR(50) NULL,
  456. remark TEXT NULL,
  457. priority VARCHAR(255) NULL,
  458. inspect_start_time DATETIME NULL,
  459. inspect_finish_time DATETIME NULL,
  460. inspector VARCHAR(255) NULL,
  461. inspect_progress VARCHAR(255) NULL,
  462. source_row_id VARCHAR(100) NOT NULL,
  463. source_biz_key VARCHAR(200) NULL,
  464. sync_batch_id VARCHAR(100) NOT NULL,
  465. sync_time DATETIME NOT NULL,
  466. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  467. UNIQUE KEY uk_mdp_std_fqc_task (tenant_id, source_row_id),
  468. KEY idx_mdp_std_fqc_task_bill (tenant_id, bill_no),
  469. KEY idx_mdp_std_fqc_task_apply (tenant_id, apply_time)
  470. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S7 FQC 检验任务标准层'
  471. """);
  472. await _db.Ado.ExecuteCommandAsync(
  473. """
  474. CREATE TABLE IF NOT EXISTS mdp_std_fqc_result (
  475. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  476. tenant_id BIGINT NOT NULL DEFAULT 0,
  477. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  478. bill_no VARCHAR(80) NULL,
  479. inspect_time DATETIME NULL,
  480. production_batch_no VARCHAR(255) NULL,
  481. material_code VARCHAR(255) NULL,
  482. material_name VARCHAR(255) NULL,
  483. spec VARCHAR(255) NULL,
  484. inspector VARCHAR(255) NULL,
  485. inspect_qty DECIMAL(18,6) NULL,
  486. source_bill_no VARCHAR(100) NULL,
  487. production_order_no VARCHAR(255) NULL,
  488. order_qty DECIMAL(18,6) NULL,
  489. sap_work_order_no VARCHAR(64) NULL,
  490. judgment VARCHAR(32) NULL,
  491. qualified_qty DECIMAL(18,6) NULL,
  492. unqualified_qty DECIMAL(18,6) NULL,
  493. fire_no TEXT NULL,
  494. inspect_spec_no VARCHAR(255) NULL,
  495. inspect_spec_version VARCHAR(255) NULL,
  496. source_row_id VARCHAR(100) NOT NULL,
  497. source_biz_key VARCHAR(200) NULL,
  498. sync_batch_id VARCHAR(100) NOT NULL,
  499. sync_time DATETIME NOT NULL,
  500. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  501. UNIQUE KEY uk_mdp_std_fqc_result (tenant_id, source_row_id),
  502. KEY idx_mdp_std_fqc_result_bill (tenant_id, bill_no),
  503. KEY idx_mdp_std_fqc_result_inspect (tenant_id, inspect_time)
  504. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S7 FQC 检验单(结果)标准层'
  505. """);
  506. var jDdl = string.Join(",\n ", Enumerable.Range(1, 50).Select(i =>
  507. $"j{i} VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL"));
  508. await _db.Ado.ExecuteCommandAsync(
  509. $"""
  510. CREATE TABLE IF NOT EXISTS mdp_std_fqc_result_detail (
  511. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  512. tenant_id BIGINT NOT NULL DEFAULT 0,
  513. std_result_id BIGINT NULL,
  514. glid VARCHAR(100) NULL,
  515. seq_no INT NULL,
  516. inspection_item TEXT NULL,
  517. accept_standard TEXT NULL,
  518. tech_requirement TEXT NULL,
  519. sampling_plan TEXT NULL,
  520. sample_qty VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL,
  521. upper_limit VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL,
  522. lower_limit VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL,
  523. size_specification VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL,
  524. judgment VARCHAR(255) CHARACTER SET utf8mb3 COLLATE utf8mb3_general_ci NULL,
  525. non_numeric_ok BIGINT NULL,
  526. non_numeric_ng BIGINT NULL,
  527. remark TEXT NULL,
  528. {jDdl},
  529. source_row_id VARCHAR(100) NOT NULL,
  530. source_biz_key VARCHAR(200) NULL,
  531. sync_batch_id VARCHAR(100) NOT NULL,
  532. sync_time DATETIME NOT NULL,
  533. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  534. UNIQUE KEY uk_mdp_std_fqc_result_dtl (tenant_id, source_row_id),
  535. KEY idx_mdp_std_fqc_result_dtl_head (tenant_id, glid),
  536. KEY idx_mdp_std_fqc_result_dtl_stdres (std_result_id)
  537. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S7 FQC 检验单明细标准层(j1-j50样本)'
  538. """);
  539. }
  540. /// <summary>标准化任务:pull stg(qms_fqcbj) → mdp_std_fqc_task。返回处理行数。</summary>
  541. private async Task<int> TransformTaskStandardAsync(string batchId, DateTime now, string? sourceSystem = null, FqcBillScope? scope = null)
  542. {
  543. var srcClause = sourceSystem == null ? "" : " AND t.source_system=@Src";
  544. var tTenant = MdpJsonSql.TenantFromStgJsonCol("t", MdpJsonSql.Int("t", "tenant_id"));
  545. var where = $"t.source_table='qms_fqcbj'{srcClause} AND {MdpJsonSql.TenantGuard(tTenant)}";
  546. // 单据级作用域(targeted sync):任务经报检单号关联(qms_fqcbj.FBILLNO = inspbill.lydjbh),并强制租户相等
  547. if (scope != null)
  548. where += $" AND ({tTenant})=@ScopeTenant AND {MdpJsonSql.Str("t", "FBILLNO")}=@ScopeBillNo";
  549. var countPars = new List<SugarParameter>();
  550. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  551. AddScopePars(countPars, scope);
  552. var rows = await _db.Ado.GetIntAsync(
  553. $"SELECT COUNT(1) FROM mdp_stg_fqc_pull t WHERE {where}", countPars);
  554. var insertSql =
  555. $"""
  556. INSERT INTO mdp_std_fqc_task
  557. (tenant_id, source_system, bill_no, production_order_no, apply_org_id, apply_time, material_code,
  558. product_name, product_model, production_batch_no, order_qty, qty, applicant, remark, priority,
  559. inspect_start_time, inspect_finish_time, inspector, inspect_progress,
  560. source_row_id, source_biz_key, sync_batch_id, sync_time)
  561. SELECT
  562. {tTenant}, IFNULL(NULLIF(t.source_system,''), 'AIDOP'),
  563. {MdpJsonSql.Str("t", "FBILLNO")}, {MdpJsonSql.Str("t", "sczld")}, {MdpJsonSql.Int("t", "FORGID")}, {MdpJsonSql.DateTimeSec("t", "FAPPLYTIME")}, {MdpJsonSql.Str("t", "wlbm")},
  564. {MdpJsonSql.Str("t", "cpmc")}, {MdpJsonSql.Str("t", "cpxh")}, {MdpJsonSql.Str("t", "scph")}, {MdpJsonSql.Dec("t", "sczldsl", 18, 6)}, {MdpJsonSql.Dec("t", "sl", 18, 6)}, {MdpJsonSql.Str("t", "FAPPLYUSER")}, {MdpJsonSql.Str("t", "FCOMMENT")}, {MdpJsonSql.Str("t", "yxj")},
  565. {MdpJsonSql.DateTimeSec("t", "jykssj")}, {MdpJsonSql.DateTimeSec("t", "jywcsj")}, {MdpJsonSql.Str("t", "jyfzr")}, {MdpJsonSql.Str("t", "FINSPECTSTATUS")},
  566. IFNULL(t.source_row_id, {MdpJsonSql.Str("t", "id")}), IFNULL(NULLIF(t.source_biz_key,''), {MdpJsonSql.Str("t", "FBILLNO")}),
  567. @BatchId, @Now
  568. FROM mdp_stg_fqc_pull t
  569. WHERE {where}
  570. ON DUPLICATE KEY UPDATE
  571. bill_no=VALUES(bill_no), production_order_no=VALUES(production_order_no), apply_org_id=VALUES(apply_org_id),
  572. apply_time=VALUES(apply_time), material_code=VALUES(material_code), product_name=VALUES(product_name),
  573. product_model=VALUES(product_model), production_batch_no=VALUES(production_batch_no), order_qty=VALUES(order_qty),
  574. qty=VALUES(qty), applicant=VALUES(applicant), remark=VALUES(remark), priority=VALUES(priority),
  575. inspect_start_time=VALUES(inspect_start_time), inspect_finish_time=VALUES(inspect_finish_time),
  576. inspector=VALUES(inspector), inspect_progress=VALUES(inspect_progress),
  577. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  578. """;
  579. var insPars = new List<SugarParameter> { new("@BatchId", batchId), new("@Now", now) };
  580. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  581. AddScopePars(insPars, scope);
  582. await _db.Ado.ExecuteCommandAsync(insertSql, insPars);
  583. return rows;
  584. }
  585. /// <summary>标准化结果:pull stg(qms_qcpp_inspbill) + WorkOrdMaster(本库富化 SAP 工单号) → mdp_std_fqc_result。</summary>
  586. private async Task<int> TransformResultStandardAsync(string batchId, DateTime now, string? sourceSystem = null, FqcBillScope? scope = null)
  587. {
  588. var srcClause = sourceSystem == null ? "" : " AND a.source_system=@Src";
  589. var aTenant = MdpJsonSql.TenantFromStgJsonCol("a", MdpJsonSql.Int("a", "tenant_id"));
  590. var where = $"a.source_table='qms_qcpp_inspbill'{srcClause} AND {MdpJsonSql.TenantGuard(aTenant)}";
  591. // 单据级作用域(targeted sync):按 stg 的 source_row_id(=inspbill.id)收敛,并强制租户相等,杜绝跨租户命中
  592. if (scope != null)
  593. where += $" AND ({aTenant})=@ScopeTenant AND IFNULL(a.source_row_id, {MdpJsonSql.Str("a", "id")})=@ScopeBillId";
  594. var sapSub = $"(SELECT MIN(wo.WorkOrd) FROM WorkOrdMaster wo WHERE wo.Batch = {MdpJsonSql.Str("a", "sczld")})";
  595. // 判定码表归一(FIX-A2):源 pd 存两套编码——业务写入 0/1(FqcInspBillFlowService 校验口径 0=合格/1=不合格),
  596. // 历史归口数据直存中文。std 面向展示统一为中文;非 0/1 的历史值 ELSE 原样透传,绝不破坏既有数据。
  597. // 仅映射整单 pd;明细行 pd(qms_qcpp_inspbillst.pd)源本身即中文,不映射。源表 pd 保持 0/1 不动。
  598. var judgmentExpr = $"CASE {MdpJsonSql.Str("a", "pd")} WHEN '0' THEN '合格' WHEN '1' THEN '不合格' ELSE {MdpJsonSql.Str("a", "pd")} END";
  599. var countPars = new List<SugarParameter>();
  600. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  601. AddScopePars(countPars, scope);
  602. var rows = await _db.Ado.GetIntAsync(
  603. $"SELECT COUNT(1) FROM mdp_stg_fqc_pull a WHERE {where}", countPars);
  604. var insertSql =
  605. $"""
  606. INSERT INTO mdp_std_fqc_result
  607. (tenant_id, source_system, bill_no, inspect_time, production_batch_no, material_code, material_name, spec,
  608. inspector, inspect_qty, source_bill_no, production_order_no, order_qty, sap_work_order_no, judgment,
  609. qualified_qty, unqualified_qty, fire_no, inspect_spec_no, inspect_spec_version,
  610. source_row_id, source_biz_key, sync_batch_id, sync_time)
  611. SELECT
  612. {aTenant}, IFNULL(NULLIF(a.source_system,''), 'AIDOP'),
  613. {MdpJsonSql.Str("a", "FBILLNO")}, {MdpJsonSql.DateTimeSec("a", "FINSPEENDDATE")}, {MdpJsonSql.Str("a", "scph")}, {MdpJsonSql.Str("a", "FMATERIALCFG")}, {MdpJsonSql.Str("a", "wlmc")}, {MdpJsonSql.Str("a", "ggxh")},
  614. {MdpJsonSql.Str("a", "jyr")}, {MdpJsonSql.Dec("a", "jysj", 18, 6)}, {MdpJsonSql.Str("a", "lydjbh")}, {MdpJsonSql.Str("a", "sczld")}, {MdpJsonSql.Dec("a", "sczldsl", 18, 6)}, {sapSub}, {judgmentExpr},
  615. {MdpJsonSql.Dec("a", "hgsl", 18, 6)}, {MdpJsonSql.Dec("a", "bhgsl", 18, 6)}, {MdpJsonSql.Str("a", "jfbh")}, {MdpJsonSql.Str("a", "jgbh")}, {MdpJsonSql.Str("a", "jgbb")},
  616. IFNULL(a.source_row_id, {MdpJsonSql.Str("a", "id")}), IFNULL(NULLIF(a.source_biz_key,''), {MdpJsonSql.Str("a", "FBILLNO")}),
  617. @BatchId, @Now
  618. FROM mdp_stg_fqc_pull a
  619. WHERE {where}
  620. ON DUPLICATE KEY UPDATE
  621. bill_no=VALUES(bill_no), inspect_time=VALUES(inspect_time), production_batch_no=VALUES(production_batch_no),
  622. material_code=VALUES(material_code), material_name=VALUES(material_name), spec=VALUES(spec),
  623. inspector=VALUES(inspector), inspect_qty=VALUES(inspect_qty), source_bill_no=VALUES(source_bill_no),
  624. production_order_no=VALUES(production_order_no), order_qty=VALUES(order_qty), sap_work_order_no=VALUES(sap_work_order_no),
  625. judgment=VALUES(judgment), qualified_qty=VALUES(qualified_qty), unqualified_qty=VALUES(unqualified_qty),
  626. fire_no=VALUES(fire_no), inspect_spec_no=VALUES(inspect_spec_no), inspect_spec_version=VALUES(inspect_spec_version),
  627. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  628. """;
  629. var insPars = new List<SugarParameter> { new("@BatchId", batchId), new("@Now", now) };
  630. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  631. AddScopePars(insPars, scope);
  632. await _db.Ado.ExecuteCommandAsync(insertSql, insPars);
  633. return rows;
  634. }
  635. /// <summary>标准化明细:pull stg(qms_qcpp_inspbillst) → mdp_std_fqc_result_detail(回填 std_result_id,j1..j50 逐列)。</summary>
  636. private async Task<int> TransformResultDetailStandardAsync(string batchId, DateTime now, string? sourceSystem = null, FqcBillScope? scope = null)
  637. {
  638. var srcClause = sourceSystem == null ? "" : " AND d.source_system=@Src";
  639. var dTenant = MdpJsonSql.TenantFromStgJsonCol("d", MdpJsonSql.Int("d", "tenant_id"));
  640. // 单据级作用域(targeted sync):明细按 glid(=inspbill.id)收敛,并强制租户相等
  641. var scopeClause = scope == null
  642. ? ""
  643. : $"\n AND ({dTenant})=@ScopeTenant AND {MdpJsonSql.Str("d", "glid")}=@ScopeBillId";
  644. var fromWhere =
  645. $"""
  646. FROM mdp_stg_fqc_pull d
  647. LEFT JOIN mdp_std_fqc_result r
  648. ON r.tenant_id = {dTenant} AND r.source_row_id = {MdpJsonSql.Str("d", "glid")}
  649. WHERE d.source_table='qms_qcpp_inspbillst'{srcClause}
  650. AND {MdpJsonSql.TenantGuard(dTenant)}{scopeClause}
  651. """;
  652. var countPars = new List<SugarParameter>();
  653. if (sourceSystem != null) countPars.Add(new SugarParameter("@Src", sourceSystem));
  654. AddScopePars(countPars, scope);
  655. var rows = await _db.Ado.GetIntAsync($"SELECT COUNT(1) {fromWhere}", countPars);
  656. var jCols = string.Join(", ", Enumerable.Range(1, 50).Select(i => $"j{i}"));
  657. var jSelect = string.Join(", ", Enumerable.Range(1, 50).Select(i => MdpJsonSql.Str("d", $"j{i}")));
  658. var jUpdate = string.Join(", ", Enumerable.Range(1, 50).Select(i => $"j{i}=VALUES(j{i})"));
  659. var insertSql =
  660. $"""
  661. INSERT INTO mdp_std_fqc_result_detail
  662. (tenant_id, std_result_id, glid, seq_no, inspection_item, accept_standard, tech_requirement, sampling_plan,
  663. sample_qty, upper_limit, lower_limit, size_specification, judgment, non_numeric_ok, non_numeric_ng, remark,
  664. {jCols},
  665. source_row_id, source_biz_key, sync_batch_id, sync_time)
  666. SELECT
  667. {dTenant}, r.id, {MdpJsonSql.Str("d", "glid")}, {MdpJsonSql.Int("d", "xh")},
  668. {MdpJsonSql.Str("d", "jyxm")}, {MdpJsonSql.Str("d", "jsbz")}, {MdpJsonSql.Str("d", "jsyq")}, {MdpJsonSql.Str("d", "cyfa")},
  669. {MdpJsonSql.Str("d", "ybl")}, {MdpJsonSql.Str("d", "sx")}, {MdpJsonSql.Str("d", "xx")}, {MdpJsonSql.Str("d", "SizeSpecifications")}, {MdpJsonSql.Str("d", "pd")}, {MdpJsonSql.Int("d", "NonNumericalTypeOK")}, {MdpJsonSql.Int("d", "NonNumericalTypeNG")}, {MdpJsonSql.Str("d", "bz")},
  670. {jSelect},
  671. IFNULL(d.source_row_id, {MdpJsonSql.Str("d", "id")}), IFNULL(NULLIF(d.source_biz_key,''), {MdpJsonSql.Str("d", "id")}),
  672. @BatchId, @Now
  673. {fromWhere}
  674. ON DUPLICATE KEY UPDATE
  675. std_result_id=VALUES(std_result_id), glid=VALUES(glid), seq_no=VALUES(seq_no),
  676. inspection_item=VALUES(inspection_item), accept_standard=VALUES(accept_standard), tech_requirement=VALUES(tech_requirement),
  677. sampling_plan=VALUES(sampling_plan), sample_qty=VALUES(sample_qty), upper_limit=VALUES(upper_limit),
  678. lower_limit=VALUES(lower_limit), size_specification=VALUES(size_specification), judgment=VALUES(judgment),
  679. non_numeric_ok=VALUES(non_numeric_ok), non_numeric_ng=VALUES(non_numeric_ng), remark=VALUES(remark),
  680. {jUpdate},
  681. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  682. """;
  683. var insPars = new List<SugarParameter> { new("@BatchId", batchId), new("@Now", now) };
  684. if (sourceSystem != null) insPars.Add(new SugarParameter("@Src", sourceSystem));
  685. AddScopePars(insPars, scope);
  686. await _db.Ado.ExecuteCommandAsync(insertSql, insPars);
  687. return rows;
  688. }
  689. private async Task<long> InsertRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  690. {
  691. await _db.Ado.ExecuteCommandAsync(
  692. """
  693. INSERT INTO mdp_transform_run_log
  694. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  695. VALUES (0, @JobCode, 'S7 FQC MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  696. """,
  697. new SugarParameter("@JobCode", JobCode),
  698. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  699. new SugarParameter("@BatchId", batchId),
  700. new SugarParameter("@StartTime", startedAt));
  701. return await _db.Ado.GetLongAsync(
  702. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  703. new List<SugarParameter> { new("@BatchId", batchId) });
  704. }
  705. private async Task MarkRunSuccessAsync(long runLogId, DateTime startedAt, FqcMdpSyncResult result)
  706. {
  707. var finishedAt = DateTime.Now;
  708. await _db.Ado.ExecuteCommandAsync(
  709. """
  710. UPDATE mdp_transform_run_log
  711. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  712. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=0, update_time=CURRENT_TIMESTAMP
  713. WHERE id=@Id
  714. """,
  715. new SugarParameter("@EndTime", finishedAt),
  716. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  717. new SugarParameter("@StageRows", result.StageRows),
  718. new SugarParameter("@StandardRows", result.TaskRows + result.ResultRows + result.DetailRows),
  719. new SugarParameter("@Id", runLogId));
  720. }
  721. private async Task MarkRunFailedAsync(long runLogId, DateTime startedAt, string message)
  722. {
  723. var finishedAt = DateTime.Now;
  724. await _db.Ado.ExecuteCommandAsync(
  725. """
  726. UPDATE mdp_transform_run_log
  727. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  728. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  729. WHERE id=@Id
  730. """,
  731. new SugarParameter("@EndTime", finishedAt),
  732. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  733. new SugarParameter("@ErrorMessage", message.Length > 2000 ? message[..2000] : message),
  734. new SugarParameter("@Id", runLogId));
  735. }
  736. private static string NormalizeTriggerType(string? triggerType)
  737. => string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  738. }
  739. /// <summary>S7 FQC MDP 同步转换结果。</summary>
  740. public sealed class FqcMdpSyncResult
  741. {
  742. public long RunLogId { get; set; }
  743. public string BatchId { get; set; } = string.Empty;
  744. public int StageRows { get; set; }
  745. public int TaskRows { get; set; }
  746. public int ResultRows { get; set; }
  747. public int DetailRows { get; set; }
  748. /// <summary>1E-0:并发锁未抢到、本次跳过(交由持锁者或下一轮刷新)。</summary>
  749. public bool Skipped { get; set; }
  750. }
  751. /// <summary>S7 FQC 双模式入站结果。</summary>
  752. public sealed class FqcInboundResult
  753. {
  754. public string PullBatchId { get; set; } = string.Empty;
  755. public int RowsPulled { get; set; }
  756. public int RowsWrittenStg { get; set; }
  757. public string? NewCursor { get; set; }
  758. public string TransformBatchId { get; set; } = string.Empty;
  759. public int TaskRows { get; set; }
  760. public int ResultRows { get; set; }
  761. public int DetailRows { get; set; }
  762. }