S2MdpSyncTransformService.cs 59 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  3. using Admin.NET.Plugin.AiDOP.SmartOps;
  4. namespace Admin.NET.Plugin.AiDOP.Production;
  5. /// <summary>
  6. /// S2 生产排程 MDP 同步、标准化、DWD 与 KPI 计算服务。
  7. /// Phase2:源→stg 经统一执行器;std/DWD/KPI 仍读 stg。
  8. /// </summary>
  9. public class S2MdpSyncTransformService : ITransient
  10. {
  11. private const string JobCode = "S2_MDP_SYNC_TRANSFORM";
  12. private readonly ISqlSugarClient _db;
  13. private readonly TransformRunLogFinalizer _runLogFinalizer;
  14. private readonly SmartOpsKpiAtomicBuildService _atomicBuild;
  15. private readonly MdpModuleStagingPuller _stagingPuller;
  16. private readonly IKpiTargetResolver _kpiTargetResolver;
  17. private MdpRebuildScope _runScope = null!;
  18. public S2MdpSyncTransformService(
  19. ISqlSugarClient db,
  20. SmartOpsKpiAtomicBuildService atomicBuild,
  21. MdpModuleStagingPuller stagingPuller,
  22. IKpiTargetResolver kpiTargetResolver,
  23. TransformRunLogFinalizer runLogFinalizer)
  24. {
  25. _db = db;
  26. _runLogFinalizer = runLogFinalizer;
  27. _atomicBuild = atomicBuild;
  28. _stagingPuller = stagingPuller;
  29. _kpiTargetResolver = kpiTargetResolver;
  30. }
  31. public Task<S2MdpSyncTransformResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO") =>
  32. throw new InvalidOperationException("S2 全量重算必须指定租户/工厂作用域,请走 ModuleRebuildService.Enqueue");
  33. public async Task<S2MdpSyncTransformResult> RunFullAsync(
  34. MdpRebuildScope scope,
  35. CancellationToken cancellationToken = default,
  36. string triggerType = "AUTO",
  37. long? rebuildJobId = null,
  38. Func<ModuleProgressUpdate, Task>? reportProgress = null)
  39. {
  40. _ = rebuildJobId;
  41. _runScope = MdpRebuildScope.Create("S2", scope.TenantId, scope.FactoryId);
  42. cancellationToken.ThrowIfCancellationRequested();
  43. var now = DateTime.Now;
  44. var batchId = $"S2_MDP_FULL_{_runScope.TenantId}_{_runScope.FactoryId}_{now:yyyyMMddHHmmss}";
  45. var runLogId = await InsertTransformRunLogAsync(batchId, now, triggerType);
  46. var result = new S2MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  47. try
  48. {
  49. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Preparing, 0, 5, "准备运行环境"));
  50. await EnsureS2RuntimeObjectsAsync();
  51. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Staging, 1, 10, "正在拉取源数据"));
  52. result.StageRows = await SyncStagingAsync(batchId, now, cancellationToken);
  53. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Staging, 1, 35, "拉取源数据完成", result.StageRows, ModuleRebuildStages.Staging));
  54. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Standard, 2, 35, "正在标准化数据"));
  55. result.StandardRows = await TransformStandardAsync(batchId, now, cancellationToken);
  56. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Standard, 2, 50, "标准化数据完成", result.StandardRows, ModuleRebuildStages.Standard));
  57. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Dwd, 3, 50, "正在生成 DWD 明细"));
  58. result.DwdRows = await BuildDwdAsync(batchId, now, cancellationToken);
  59. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Dwd, 3, 65, "生成 DWD 明细完成", result.DwdRows, ModuleRebuildStages.Dwd));
  60. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Kpi, 4, 65, "正在重算 KPI"));
  61. result.KpiRows = await BuildS2KpiValuesAsync(batchId, now, cancellationToken);
  62. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Kpi, 4, 80, "重算 KPI 完成", result.KpiRows, ModuleRebuildStages.Kpi));
  63. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Atomic, 5, 82, "正在重建原子数据"));
  64. result.AtomicRows = await _atomicBuild.BuildWorkScheduleDomainForAllDatesAsync(
  65. _runScope.TenantId, _runScope.FactoryId, batchId, cancellationToken);
  66. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Atomic, 5, 98, "重建原子数据完成", result.AtomicRows, ModuleRebuildStages.Atomic));
  67. await ReportAsync(reportProgress, new ModuleProgressUpdate(ModuleRebuildStages.Finalizing, 5, 99, "正在收尾"));
  68. await MarkTransformRunSuccessAsync(runLogId, now, result);
  69. return result;
  70. }
  71. catch (Exception ex)
  72. {
  73. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  74. if (!_runLogFinalizer.IsHostStopping)
  75. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  76. throw;
  77. }
  78. finally
  79. {
  80. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  81. }
  82. }
  83. private static async Task ReportAsync(Func<ModuleProgressUpdate, Task>? report, ModuleProgressUpdate update)
  84. {
  85. if (report == null) return;
  86. await report(update);
  87. }
  88. private async Task EnsureS2RuntimeObjectsAsync()
  89. {
  90. foreach (var sql in S2MdpDdl.SqlBlocks)
  91. {
  92. await _db.Ado.ExecuteCommandAsync(sql);
  93. }
  94. }
  95. private async Task<int> SyncStagingAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  96. {
  97. _ = now;
  98. EnsureRunScope();
  99. return await _stagingPuller.PullEntitiesAsync(
  100. S2MdpEntityConfig.All.Select(x => x.EntityCode),
  101. batchId, _runScope.TenantId, fullRefresh: true, taskCode: "S2_MDP_INBOUND", cancellationToken,
  102. factoryId: _runScope.FactoryId, requireMatchingSourceTenant: true);
  103. }
  104. public async Task<S2MdpSyncTransformResult> RunInboundAsync(
  105. IEnumerable<string>? entityCodes = null,
  106. long tenantId = 0,
  107. bool fullRefresh = false,
  108. CancellationToken cancellationToken = default)
  109. {
  110. cancellationToken.ThrowIfCancellationRequested();
  111. if (tenantId <= 0)
  112. throw new InvalidOperationException("S2 inbound 必须指定有效 tenantId");
  113. _runScope = MdpRebuildScope.Create("S2", tenantId, 1);
  114. var now = DateTime.Now;
  115. var batchId = $"S2_MDP_IN_{_runScope.TenantId}_{now:yyyyMMddHHmmss}";
  116. var runLogId = await InsertTransformRunLogAsync(batchId, now, "INBOUND");
  117. var result = new S2MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  118. try
  119. {
  120. await EnsureS2RuntimeObjectsAsync();
  121. var codes = entityCodes?.ToList() ?? S2MdpEntityConfig.All.Select(x => x.EntityCode).ToList();
  122. result.StageRows = await _stagingPuller.PullEntitiesAsync(
  123. codes, batchId, tenantId, fullRefresh, "S2_MDP_INBOUND", cancellationToken);
  124. result.StandardRows = await TransformStandardAsync(batchId, now, cancellationToken);
  125. result.DwdRows = await BuildDwdAsync(batchId, now, cancellationToken);
  126. result.KpiRows = await BuildS2KpiValuesAsync(batchId, now, cancellationToken);
  127. result.AtomicRows = await _atomicBuild.BuildWorkScheduleDomainForAllDatesAsync(
  128. _runScope.TenantId, _runScope.FactoryId, batchId, cancellationToken);
  129. await MarkTransformRunSuccessAsync(runLogId, now, result);
  130. return result;
  131. }
  132. catch (Exception ex)
  133. {
  134. // 宿主关停不是转换失败:交给 finally 收口为 ABORTED,不污染 FAILED 语义。
  135. if (!_runLogFinalizer.IsHostStopping)
  136. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  137. throw;
  138. }
  139. finally
  140. {
  141. await _runLogFinalizer.FinalizeIfHostStoppingAsync(runLogId, now);
  142. }
  143. }
  144. [Obsolete("Phase3: use MdpModuleStagingPuller / executors")]
  145. private async Task<int> SyncOneEntityAsync(S2MdpEntityConfig entity, string batchId, DateTime now)
  146. {
  147. var entityRow = await _db.Ado.SqlQuerySingleAsync<S2MdpEntityRow>(
  148. "SELECT id AS Id, entity_name AS EntityName FROM mdp_entity WHERE tenant_id=0 AND entity_code=@EntityCode LIMIT 1",
  149. new SugarParameter("@EntityCode", entity.EntityCode));
  150. if (entityRow == null) throw Oops.Oh($"未找到 MDP 实体配置:{entity.EntityCode}");
  151. var columns = await _db.Ado.SqlQueryAsync<S2ColumnRow>(
  152. """
  153. SELECT COLUMN_NAME AS ColumnName
  154. FROM information_schema.COLUMNS
  155. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@TableName
  156. ORDER BY ORDINAL_POSITION
  157. """,
  158. new SugarParameter("@TableName", entity.SourceTable));
  159. if (columns.Count == 0) throw Oops.Oh($"未找到源表:{entity.SourceTable}");
  160. var names = columns.Select(u => u.ColumnName).ToList();
  161. var tenantExpr = names.Any(u => string.Equals(u, "tenant_id", StringComparison.OrdinalIgnoreCase))
  162. ? $"IFNULL(s.`{FindColumn(names, "tenant_id")}`,0)"
  163. : "0";
  164. var factoryExpr = names.Any(u => string.Equals(u, "factory_id", StringComparison.OrdinalIgnoreCase))
  165. ? $"s.`{FindColumn(names, "factory_id")}`"
  166. : "NULL";
  167. var domainExpr = names.Any(u => string.Equals(u, "Domain", StringComparison.OrdinalIgnoreCase))
  168. ? $"s.`{FindColumn(names, "Domain")}`"
  169. : "NULL";
  170. var sourceRowExpr = names.Any(u => string.Equals(u, entity.SourceRowIdExpression, StringComparison.OrdinalIgnoreCase))
  171. ? $"s.`{FindColumn(names, entity.SourceRowIdExpression)}`"
  172. : entity.SourceRowIdExpression;
  173. var rawDataExpr = BuildJsonObjectExpression(names);
  174. var rowsRead = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM `{entity.SourceTable}`");
  175. var logId = await InsertSyncLogAsync(entityRow.Id, entityRow.EntityName, batchId, rowsRead);
  176. var started = DateTime.Now;
  177. try
  178. {
  179. var affected = await _db.Ado.ExecuteCommandAsync(
  180. $"""
  181. INSERT INTO `{entity.TargetTable}`
  182. (tenant_id, factory_id, source_system, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time, process_status, raw_data)
  183. SELECT
  184. {tenantExpr},
  185. COALESCE({factoryExpr}, {domainExpr}),
  186. 'AIDOP',
  187. @SourceTable,
  188. CAST({sourceRowExpr} AS CHAR),
  189. CAST(COALESCE({entity.SourceBizKeyExpression}, CAST({sourceRowExpr} AS CHAR)) AS CHAR),
  190. @BatchId,
  191. @Now,
  192. 'PENDING',
  193. {rawDataExpr}
  194. FROM `{entity.SourceTable}` s
  195. ON DUPLICATE KEY UPDATE
  196. tenant_id=VALUES(tenant_id),
  197. factory_id=VALUES(factory_id),
  198. sync_batch_id=VALUES(sync_batch_id),
  199. sync_time=VALUES(sync_time),
  200. process_status=VALUES(process_status),
  201. raw_data=VALUES(raw_data),
  202. update_time=CURRENT_TIMESTAMP
  203. """,
  204. new SugarParameter("@SourceTable", entity.SourceTable),
  205. new SugarParameter("@BatchId", batchId),
  206. new SugarParameter("@Now", now));
  207. await MarkSyncLogSuccessAsync(logId, started, affected);
  208. return rowsRead;
  209. }
  210. catch (Exception ex)
  211. {
  212. await MarkSyncLogFailedAsync(logId, started, ex.Message);
  213. throw;
  214. }
  215. }
  216. private async Task<int> TransformStandardAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  217. {
  218. var total = 0;
  219. foreach (var command in BuildStandardCommands(batchId, now))
  220. {
  221. cancellationToken.ThrowIfCancellationRequested();
  222. total += await _db.Ado.ExecuteCommandAsync(command.Sql, command.Parameters);
  223. }
  224. return total;
  225. }
  226. private async Task<int> BuildDwdAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  227. {
  228. var total = 0;
  229. foreach (var command in BuildDwdCommands(batchId, now))
  230. {
  231. cancellationToken.ThrowIfCancellationRequested();
  232. total += await _db.Ado.ExecuteCommandAsync(command.Sql, command.Parameters);
  233. }
  234. return total;
  235. }
  236. private async Task<int> BuildS2KpiValuesAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  237. {
  238. var affected = 0;
  239. for (var dayOffset = 13; dayOffset >= 0; dayOffset--)
  240. {
  241. var statDate = now.Date.AddDays(-dayOffset);
  242. var rows = await CalculateS2KpiValuesAsync(batchId, statDate);
  243. foreach (var row in rows)
  244. {
  245. cancellationToken.ThrowIfCancellationRequested();
  246. affected += await UpsertS2KpiValueAsync(row, statDate, now);
  247. }
  248. }
  249. return affected;
  250. }
  251. private IEnumerable<S2MdpSqlCommand> BuildStandardCommands(string batchId, DateTime now)
  252. {
  253. yield return Cmd(
  254. """
  255. INSERT INTO mdp_std_work_order_schedule
  256. (tenant_id, factory_id, source_system, work_order, sales_order_no, sales_order_entry_id, item_code, item_name, site_code, status,
  257. priority, urgent_flag, qty_ordered, qty_completed, order_date, due_date, release_date, prod_line,
  258. specification, lot_serial, drawing_no, project, work_order_type, labor_variance,
  259. source_biz_key, sync_batch_id, sync_time)
  260. SELECT s.tenant_id,
  261. CASE WHEN s.factory_id REGEXP '^[0-9]+$' THEN CAST(s.factory_id AS UNSIGNED) ELSE 1 END,
  262. 'AIDOP',
  263. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrd')),
  264. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.SalesJob')),
  265. -- 订单行 ID 桥:WorkOrdMaster.BusinessID,>0 时恒等于 crm_seorderentry.Id
  266. -- (全库实测 28/28 命中同租户订单行、0 孤儿)。源值 0 表示「无关联订单行」,
  267. -- 必须归一为 NULL —— 0 不是合法订单行 Id,落 0 会在桥上造出假关联。
  268. -- sales_order_no(=SalesJob) 只是文本,禁止用它做归因。
  269. NULLIF(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.BusinessID')) REGEXP '^[0-9]+$'
  270. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.BusinessID')) AS SIGNED) END, 0),
  271. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemNum')),
  272. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemName')),
  273. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Site')),
  274. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Status')),
  275. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Priority')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Priority')) AS DECIMAL(18,6)) END,
  276. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Urgent')) IN ('1','true','True') THEN 1 ELSE 0 END,
  277. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyOrded')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyOrded')) AS DECIMAL(18,6)) ELSE 0 END,
  278. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyCompleted')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyCompleted')) AS DECIMAL(18,6)) ELSE 0 END,
  279. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdDate')), 'null'), ''),
  280. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.DueDate')), 'null'), ''),
  281. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ReleaseDate')), 'null'), ''),
  282. JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ProdLine')),
  283. -- Specification(型号)= ItemMaster.Descr1 中台 enrich((Domain,ItemNum) 去重取一,实测同键 Descr1 无冲突;
  284. -- 严格按 Domain+ItemNum 关联,不做 ItemNum-only 去重,避免跨工厂串型号):
  285. im.Descr1,
  286. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LotSerial')), 'null'), ''),
  287. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Drawing')), 'null'), ''),
  288. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Project')), 'null'), ''),
  289. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WoTyped')), 'null'), ''),
  290. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LbrVar')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LbrVar')) AS DECIMAL(18,6)) END,
  291. s.source_biz_key, @BatchId, @Now
  292. FROM mdp_stg_schedule s
  293. LEFT JOIN (
  294. SELECT `Domain`, ItemNum, MIN(Descr1) AS Descr1
  295. FROM ItemMaster
  296. GROUP BY `Domain`, ItemNum
  297. ) im ON im.`Domain` = JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain'))
  298. AND im.ItemNum = JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ItemNum'))
  299. WHERE s.source_table='WorkOrdMaster' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.WorkOrd')), '') <> ''
  300. ON DUPLICATE KEY UPDATE
  301. sales_order_no=VALUES(sales_order_no), sales_order_entry_id=VALUES(sales_order_entry_id),
  302. item_code=VALUES(item_code), item_name=VALUES(item_name),
  303. site_code=VALUES(site_code), status=VALUES(status), priority=VALUES(priority), urgent_flag=VALUES(urgent_flag),
  304. qty_ordered=VALUES(qty_ordered), qty_completed=VALUES(qty_completed), order_date=VALUES(order_date),
  305. due_date=VALUES(due_date), release_date=VALUES(release_date), prod_line=VALUES(prod_line),
  306. specification=VALUES(specification), lot_serial=VALUES(lot_serial), drawing_no=VALUES(drawing_no),
  307. project=VALUES(project), work_order_type=VALUES(work_order_type), labor_variance=VALUES(labor_variance),
  308. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  309. """, batchId, now);
  310. yield return Cmd(
  311. """
  312. INSERT INTO mdp_std_operation_schedule
  313. (tenant_id, factory_id, source_system, work_order, op_no, work_center, line_code, item_code,
  314. plan_date, prod_date, start_time, end_time, ord_qty, comp_qty, run_crew, employee, status,
  315. source_biz_key, sync_batch_id, sync_time)
  316. SELECT tenant_id,
  317. CASE WHEN factory_id REGEXP '^[0-9]+$' THEN CAST(factory_id AS UNSIGNED) ELSE 1 END,
  318. 'AIDOP',
  319. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrds')),
  320. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Op')),
  321. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkCtr')),
  322. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Line')),
  323. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ItemNum')),
  324. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.PlanDate')), 'null'), ''),
  325. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ProdDate')), 'null'), ''),
  326. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''),
  327. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.EndTime')), 'null'), ''),
  328. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.OrdQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.OrdQty')) AS DECIMAL(18,6)) ELSE 0 END,
  329. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompQty')) AS DECIMAL(18,6)) ELSE 0 END,
  330. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RunCrew')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.RunCrew')) AS DECIMAL(18,6)) END,
  331. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Employee')),
  332. -- 日计划下达状态:'Y'=已下达(可执行),空=未下达。原样投影 PeriodSequenceDet.Status。
  333. -- 它只是【状态佐证】—— Stage-2 End 的时间 Authority 一律取
  334. -- S2_DAILY_PLAN_RELEASE SUCCESS 的 end_time,绝不能用 update_time 当下达时间
  335. -- (update_time 会被后续任意一次 UPDATE 覆盖)。
  336. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Status')),
  337. source_biz_key, @BatchId, @Now
  338. FROM mdp_stg_schedule
  339. WHERE source_table='PeriodSequenceDet' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrds')), '') <> ''
  340. ON DUPLICATE KEY UPDATE
  341. work_center=VALUES(work_center), line_code=VALUES(line_code), item_code=VALUES(item_code),
  342. plan_date=VALUES(plan_date), prod_date=VALUES(prod_date), start_time=VALUES(start_time), end_time=VALUES(end_time),
  343. ord_qty=VALUES(ord_qty), comp_qty=VALUES(comp_qty), run_crew=VALUES(run_crew), employee=VALUES(employee),
  344. status=VALUES(status),
  345. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  346. """, batchId, now);
  347. // ⚠️ 同表的另一个来源:ScheduleResultOpMaster —— 它没有「日计划下达」这个概念。
  348. // 因此既不投影 status,ODKU 里也【刻意不列】status,
  349. // 否则会把 PeriodSequenceDet 分支刚投影进来的下达状态覆盖成 NULL。
  350. // (两来源 source_biz_key 格式不同,uk=(tenant_id,source_biz_key) 实测不碰撞;此为纵深防御。)
  351. yield return Cmd(
  352. """
  353. INSERT INTO mdp_std_operation_schedule
  354. (tenant_id, factory_id, source_system, work_order, op_no, work_center, line_code, item_code,
  355. plan_date, prod_date, start_time, end_time, ord_qty, comp_qty, run_crew, employee,
  356. source_biz_key, sync_batch_id, sync_time)
  357. SELECT tenant_id,
  358. CASE WHEN factory_id REGEXP '^[0-9]+$' THEN CAST(factory_id AS UNSIGNED) ELSE 1 END,
  359. 'AIDOP',
  360. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrd')),
  361. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Op')),
  362. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkCtr')),
  363. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Line')),
  364. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ItemNum')),
  365. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkDate')), 'null'), ''),
  366. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkDate')), 'null'), ''),
  367. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkStartTime')), 'null'), ''),
  368. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkEndTime')), 'null'), ''),
  369. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkQty')) AS DECIMAL(18,6)) ELSE 0 END,
  370. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkQty')) AS DECIMAL(18,6)) ELSE 0 END,
  371. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.AssignedPersonnelCount')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.AssignedPersonnelCount')) AS DECIMAL(18,6)) END,
  372. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.AssignedEmployeeID')),
  373. source_biz_key, @BatchId, @Now
  374. FROM mdp_stg_schedule
  375. WHERE source_table='ScheduleResultOpMaster' AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrd')), '') <> ''
  376. ON DUPLICATE KEY UPDATE
  377. work_center=VALUES(work_center), line_code=VALUES(line_code), item_code=VALUES(item_code),
  378. plan_date=VALUES(plan_date), prod_date=VALUES(prod_date), start_time=VALUES(start_time), end_time=VALUES(end_time),
  379. ord_qty=VALUES(ord_qty), comp_qty=VALUES(comp_qty), run_crew=VALUES(run_crew), employee=VALUES(employee),
  380. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  381. """, batchId, now);
  382. }
  383. private IEnumerable<S2MdpSqlCommand> BuildDwdCommands(string batchId, DateTime now)
  384. {
  385. yield return Cmd(
  386. """
  387. INSERT INTO dwd_order_schedule_trans
  388. (tenant_id, factory_id, stat_date, work_order, sales_order_no, item_code, item_name, site_code, prod_line,
  389. status, urgent_flag, qty_ordered, qty_completed, order_date, due_date, release_date, first_plan_date,
  390. last_plan_date, first_start_time, last_end_time, operation_count, scheduled_qty, completed_op_qty,
  391. schedule_cycle_days, schedule_satisfaction_flag, wip_qty, resource_person_count, calc_batch_id, calc_time)
  392. SELECT w.tenant_id, COALESCE(w.factory_id, 1), @StatDate,
  393. w.work_order, w.sales_order_no, w.item_code, w.item_name, w.site_code, w.prod_line,
  394. w.status, w.urgent_flag, w.qty_ordered, w.qty_completed, w.order_date, w.due_date, w.release_date,
  395. MIN(o.plan_date), MAX(o.plan_date), MIN(o.start_time), MAX(o.end_time),
  396. COUNT(o.id), SUM(IFNULL(o.ord_qty, 0)), SUM(IFNULL(o.comp_qty, 0)),
  397. CASE
  398. WHEN COALESCE(MIN(o.plan_date), MAX(o.plan_date)) IS NOT NULL
  399. AND COALESCE(w.release_date, w.order_date) IS NOT NULL
  400. THEN TIMESTAMPDIFF(HOUR, COALESCE(w.release_date, w.order_date), COALESCE(MAX(o.plan_date), MIN(o.plan_date))) / 24
  401. END,
  402. CASE
  403. WHEN w.due_date IS NOT NULL AND COALESCE(MAX(o.plan_date), MIN(o.plan_date)) IS NOT NULL
  404. AND DATE(COALESCE(MAX(o.plan_date), MIN(o.plan_date))) <= DATE(w.due_date) THEN 1
  405. ELSE 0
  406. END,
  407. GREATEST(IFNULL(w.qty_ordered, 0) - IFNULL(w.qty_completed, 0), 0),
  408. SUM(IFNULL(o.run_crew, 0)),
  409. @BatchId, @Now
  410. FROM mdp_std_work_order_schedule w
  411. LEFT JOIN mdp_std_operation_schedule o
  412. ON o.tenant_id=w.tenant_id AND IFNULL(o.work_order,'')=IFNULL(w.work_order,'')
  413. WHERE IFNULL(w.work_order, '') <> ''
  414. GROUP BY w.tenant_id, COALESCE(w.factory_id, 1), w.work_order, w.sales_order_no, w.item_code, w.item_name,
  415. w.site_code, w.prod_line, w.status, w.urgent_flag, w.qty_ordered, w.qty_completed,
  416. w.order_date, w.due_date, w.release_date
  417. ON DUPLICATE KEY UPDATE
  418. sales_order_no=VALUES(sales_order_no), item_code=VALUES(item_code), item_name=VALUES(item_name),
  419. site_code=VALUES(site_code), prod_line=VALUES(prod_line), status=VALUES(status), urgent_flag=VALUES(urgent_flag),
  420. qty_ordered=VALUES(qty_ordered), qty_completed=VALUES(qty_completed), order_date=VALUES(order_date),
  421. due_date=VALUES(due_date), release_date=VALUES(release_date), first_plan_date=VALUES(first_plan_date),
  422. last_plan_date=VALUES(last_plan_date), first_start_time=VALUES(first_start_time), last_end_time=VALUES(last_end_time),
  423. operation_count=VALUES(operation_count), scheduled_qty=VALUES(scheduled_qty), completed_op_qty=VALUES(completed_op_qty),
  424. schedule_cycle_days=VALUES(schedule_cycle_days), schedule_satisfaction_flag=VALUES(schedule_satisfaction_flag),
  425. wip_qty=VALUES(wip_qty), resource_person_count=VALUES(resource_person_count), calc_time=VALUES(calc_time),
  426. update_time=CURRENT_TIMESTAMP
  427. """, batchId, now);
  428. }
  429. private async Task<List<S2KpiCalcRow>> CalculateS2KpiValuesAsync(string batchId, DateTime statDate)
  430. {
  431. var scope = RequireScope();
  432. return await _db.Ado.SqlQueryAsync<S2KpiCalcRow>(MdpSqlScope.InjectTenantFactory(
  433. """
  434. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L1_001' AS MetricCode,
  435. ROUND(AVG(schedule_cycle_days), 4) AS MetricValue
  436. FROM dwd_order_schedule_trans
  437. WHERE calc_batch_id=@BatchId AND schedule_cycle_days IS NOT NULL AND schedule_cycle_days >= 0
  438. GROUP BY tenant_id, factory_id
  439. UNION ALL
  440. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L1_002' AS MetricCode,
  441. ROUND(100 * SUM(schedule_satisfaction_flag) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  442. FROM dwd_order_schedule_trans
  443. WHERE calc_batch_id=@BatchId
  444. GROUP BY tenant_id, factory_id
  445. UNION ALL
  446. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L1_003' AS MetricCode,
  447. ROUND(COUNT(1) / NULLIF(SUM(CASE WHEN IFNULL(resource_person_count, 0) > 0 THEN resource_person_count ELSE 1 END), 0), 4) AS MetricValue
  448. FROM dwd_order_schedule_trans
  449. WHERE calc_batch_id=@BatchId
  450. GROUP BY tenant_id, factory_id
  451. UNION ALL
  452. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L1_004' AS MetricCode,
  453. ROUND(SUM(wip_qty) / NULLIF(SUM(CASE WHEN qty_completed > 0 THEN qty_completed ELSE completed_op_qty END), 0) * 30, 4) AS MetricValue
  454. FROM dwd_order_schedule_trans
  455. WHERE calc_batch_id=@BatchId
  456. GROUP BY tenant_id, factory_id
  457. UNION ALL
  458. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L2_001' AS MetricCode,
  459. ROUND(AVG(schedule_cycle_days), 4) AS MetricValue
  460. FROM dwd_order_schedule_trans
  461. WHERE calc_batch_id=@BatchId AND schedule_cycle_days IS NOT NULL AND schedule_cycle_days >= 0
  462. GROUP BY tenant_id, factory_id
  463. UNION ALL
  464. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L2_002' AS MetricCode,
  465. ROUND(100 * SUM(schedule_satisfaction_flag) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  466. FROM dwd_order_schedule_trans
  467. WHERE calc_batch_id=@BatchId
  468. GROUP BY tenant_id, factory_id
  469. UNION ALL
  470. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L2_003' AS MetricCode,
  471. ROUND(COUNT(1) / NULLIF(SUM(CASE WHEN IFNULL(resource_person_count, 0) > 0 THEN resource_person_count ELSE 1 END), 0), 4) AS MetricValue
  472. FROM dwd_order_schedule_trans
  473. WHERE calc_batch_id=@BatchId
  474. GROUP BY tenant_id, factory_id
  475. UNION ALL
  476. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_001' AS MetricCode,
  477. ROUND(AVG(
  478. CASE
  479. WHEN start_time IS NOT NULL AND end_time IS NOT NULL AND end_time >= start_time
  480. THEN TIMESTAMPDIFF(HOUR, start_time, end_time) / 24
  481. WHEN plan_date IS NOT NULL AND prod_date IS NOT NULL
  482. THEN ABS(TIMESTAMPDIFF(DAY, plan_date, prod_date))
  483. END
  484. ), 4) AS MetricValue
  485. FROM mdp_std_operation_schedule
  486. WHERE IFNULL(work_order, '') <> ''
  487. GROUP BY tenant_id, factory_id
  488. UNION ALL
  489. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_002' AS MetricCode,
  490. ROUND(100 * SUM(CASE WHEN plan_date IS NOT NULL AND (prod_date IS NULL OR prod_date >= plan_date) THEN 1 ELSE 0 END) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  491. FROM mdp_std_operation_schedule
  492. WHERE IFNULL(work_order, '') <> ''
  493. GROUP BY tenant_id, factory_id
  494. UNION ALL
  495. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_003' AS MetricCode,
  496. ROUND(COUNT(1) / NULLIF(SUM(CASE WHEN IFNULL(run_crew, 0) > 0 THEN run_crew ELSE 1 END), 0), 4) AS MetricValue
  497. FROM mdp_std_operation_schedule
  498. WHERE IFNULL(work_order, '') <> ''
  499. GROUP BY tenant_id, factory_id
  500. UNION ALL
  501. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_004' AS MetricCode,
  502. ROUND(AVG(
  503. CASE
  504. WHEN start_time IS NOT NULL AND end_time IS NOT NULL AND end_time >= start_time
  505. THEN TIMESTAMPDIFF(HOUR, start_time, end_time) / 24
  506. WHEN plan_date IS NOT NULL AND prod_date IS NOT NULL
  507. THEN ABS(TIMESTAMPDIFF(DAY, plan_date, prod_date))
  508. END
  509. ), 4) AS MetricValue
  510. FROM mdp_std_operation_schedule
  511. WHERE IFNULL(COALESCE(NULLIF(work_center,''), NULLIF(line_code,'')), '') <> ''
  512. GROUP BY tenant_id, factory_id
  513. UNION ALL
  514. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_005' AS MetricCode,
  515. ROUND(100 * SUM(CASE WHEN plan_date IS NOT NULL AND (prod_date IS NULL OR prod_date >= plan_date) THEN 1 ELSE 0 END) / NULLIF(COUNT(1), 0), 4) AS MetricValue
  516. FROM mdp_std_operation_schedule
  517. WHERE IFNULL(COALESCE(NULLIF(work_center,''), NULLIF(line_code,'')), '') <> ''
  518. GROUP BY tenant_id, factory_id
  519. UNION ALL
  520. SELECT tenant_id AS TenantId, factory_id AS FactoryId, 'S2_L3_006' AS MetricCode,
  521. ROUND(COUNT(DISTINCT COALESCE(NULLIF(work_center,''), NULLIF(line_code,''), CONCAT(work_order, ':', op_no)))
  522. / NULLIF(SUM(CASE WHEN IFNULL(run_crew, 0) > 0 THEN run_crew ELSE 1 END), 0), 4) AS MetricValue
  523. FROM mdp_std_operation_schedule
  524. WHERE IFNULL(COALESCE(NULLIF(work_center,''), NULLIF(line_code,'')), '') <> ''
  525. GROUP BY tenant_id, factory_id
  526. """),
  527. new SugarParameter("@BatchId", batchId),
  528. new SugarParameter("@StatDate", statDate),
  529. new SugarParameter("@TenantId", scope.TenantId),
  530. new SugarParameter("@FactoryId", scope.FactoryId));
  531. }
  532. private async Task<int> UpsertS2KpiValueAsync(S2KpiCalcRow row, DateTime statDate, DateTime now)
  533. {
  534. if (row.MetricValue == null) return 0;
  535. var meta = await _db.Ado.SqlQuerySingleAsync<S2KpiMetaRow>(
  536. """
  537. SELECT MetricLevel, Direction, YellowThreshold, RedThreshold
  538. FROM ado_smart_ops_kpi_master
  539. WHERE TenantId=@TenantId AND ModuleCode='S2' AND MetricCode=@MetricCode AND IsEnabled=1
  540. LIMIT 1
  541. """,
  542. new SugarParameter("@TenantId", row.TenantId),
  543. new SugarParameter("@MetricCode", row.MetricCode));
  544. if (meta == null) return 0;
  545. var table = ResolveKpiValueTable(meta.MetricLevel);
  546. var current = await _db.Ado.SqlQuerySingleAsync<S2KpiValueRow>(
  547. $"""
  548. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  549. FROM {table}
  550. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S2'
  551. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  552. ORDER BY id
  553. LIMIT 1
  554. """,
  555. new SugarParameter("@TenantId", row.TenantId),
  556. new SugarParameter("@FactoryId", row.FactoryId),
  557. new SugarParameter("@MetricCode", row.MetricCode),
  558. new SugarParameter("@BizDate", statDate));
  559. var prior = await _db.Ado.SqlQuerySingleAsync<S2KpiValueRow>(
  560. $"""
  561. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  562. FROM {table}
  563. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S2'
  564. AND metric_code=@MetricCode AND biz_date<@BizDate AND is_deleted=0
  565. ORDER BY biz_date DESC, id DESC
  566. LIMIT 1
  567. """,
  568. new SugarParameter("@TenantId", row.TenantId),
  569. new SugarParameter("@FactoryId", row.FactoryId),
  570. new SugarParameter("@MetricCode", row.MetricCode),
  571. new SugarParameter("@BizDate", statDate));
  572. var actual = Math.Round(row.MetricValue.Value, 4);
  573. var snap = await _kpiTargetResolver.ResolveAsync(row.TenantId, row.FactoryId, row.MetricCode, "S2", statDate);
  574. var target = snap.TargetValue;
  575. var status = AidopS4KpiMerge.AchievementLevel(actual, target, meta.Direction, meta.YellowThreshold, meta.RedThreshold);
  576. var trend = ResolveTrendFlag(actual, prior?.MetricValue);
  577. if (current != null)
  578. {
  579. return await _db.Ado.ExecuteCommandAsync(
  580. $"""
  581. UPDATE {table}
  582. SET metric_value=@MetricValue, target_value=@TargetValue, status_color=@StatusColor, trend_flag=@TrendFlag,
  583. target_config_id=@TargetConfigId, target_source=@TargetSource, target_resolved_at=@TargetResolvedAt,
  584. is_active=1, status='ACTIVE', calc_time=@CalcTime, update_time=@CalcTime
  585. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S2'
  586. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  587. """,
  588. new SugarParameter("@MetricValue", actual),
  589. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  590. new SugarParameter("@StatusColor", status),
  591. new SugarParameter("@TrendFlag", trend),
  592. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  593. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  594. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt),
  595. new SugarParameter("@CalcTime", now),
  596. new SugarParameter("@TenantId", row.TenantId),
  597. new SugarParameter("@FactoryId", row.FactoryId),
  598. new SugarParameter("@MetricCode", row.MetricCode),
  599. new SugarParameter("@BizDate", statDate));
  600. }
  601. var nextId = Yitter.IdGenerator.YitIdHelper.NextId();
  602. return await _db.Ado.ExecuteCommandAsync(
  603. $"""
  604. INSERT INTO {table}
  605. (id, tenant_id, factory_id, status, biz_date, create_time, update_time, is_deleted, is_active,
  606. module_code, metric_code, metric_value, target_value, status_color, trend_flag, calc_time,
  607. target_config_id, target_source, target_resolved_at)
  608. VALUES
  609. (@Id, @TenantId, @FactoryId, 'ACTIVE', @BizDate, @CalcTime, @CalcTime, 0, 1,
  610. 'S2', @MetricCode, @MetricValue, @TargetValue, @StatusColor, @TrendFlag, @CalcTime,
  611. @TargetConfigId, @TargetSource, @TargetResolvedAt)
  612. """,
  613. new SugarParameter("@Id", nextId),
  614. new SugarParameter("@TenantId", row.TenantId),
  615. new SugarParameter("@FactoryId", row.FactoryId),
  616. new SugarParameter("@BizDate", statDate),
  617. new SugarParameter("@CalcTime", now),
  618. new SugarParameter("@MetricCode", row.MetricCode),
  619. new SugarParameter("@MetricValue", actual),
  620. new SugarParameter("@TargetValue", KpiTargetSnapshotSql.ValueOrDbNull(snap)),
  621. new SugarParameter("@StatusColor", status),
  622. new SugarParameter("@TrendFlag", trend),
  623. new SugarParameter("@TargetConfigId", KpiTargetSnapshotSql.ConfigIdOrDbNull(snap)),
  624. new SugarParameter("@TargetSource", KpiTargetSnapshotSql.SourceOrDbNull(snap)),
  625. new SugarParameter("@TargetResolvedAt", snap.ResolvedAt));
  626. }
  627. private async Task<long> InsertSyncLogAsync(long entityId, string entityName, string batchId, int rowsRead)
  628. {
  629. EnsureRunScope();
  630. await _db.Ado.ExecuteCommandAsync(
  631. """
  632. INSERT INTO mdp_sync_log
  633. (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type, sync_start, rows_read, status)
  634. VALUES (@TenantId, @EntityId, 'AIDOPDEV_MYSQL', @EntityName, @BatchId, 'FULL', 'AUTO', NOW(), @RowsRead, 'RUNNING')
  635. """,
  636. new SugarParameter("@TenantId", _runScope.TenantId),
  637. new SugarParameter("@EntityId", entityId),
  638. new SugarParameter("@EntityName", entityName),
  639. new SugarParameter("@BatchId", batchId),
  640. new SugarParameter("@RowsRead", rowsRead));
  641. return await _db.Ado.GetLongAsync(
  642. "SELECT id FROM mdp_sync_log WHERE sync_batch_id=@BatchId AND entity_id=@EntityId ORDER BY id DESC LIMIT 1",
  643. new List<SugarParameter> { new("@BatchId", batchId), new("@EntityId", entityId) });
  644. }
  645. private async Task MarkSyncLogSuccessAsync(long logId, DateTime started, int affected)
  646. {
  647. await _db.Ado.ExecuteCommandAsync(
  648. """
  649. UPDATE mdp_sync_log
  650. SET sync_end=NOW(), duration_ms=@DurationMs, rows_insert=@RowsInsert, rows_update=0, rows_skip=0, rows_error=0, status='SUCCESS'
  651. WHERE id=@Id
  652. """,
  653. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  654. new SugarParameter("@RowsInsert", affected),
  655. new SugarParameter("@Id", logId));
  656. }
  657. private async Task MarkSyncLogFailedAsync(long logId, DateTime started, string message)
  658. {
  659. try
  660. {
  661. await _db.Ado.ExecuteCommandAsync(
  662. """
  663. UPDATE mdp_sync_log
  664. SET sync_end=NOW(), duration_ms=@DurationMs, rows_error=1, status='FAILED', error_msg=@ErrorMsg
  665. WHERE id=@Id
  666. """,
  667. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  668. new SugarParameter("@ErrorMsg", Truncate(message, 1000)),
  669. new SugarParameter("@Id", logId));
  670. }
  671. catch (Exception ex)
  672. {
  673. Console.Error.WriteLine($"[S2MdpSyncTransform] MarkSyncLogFailed write failed (syncLogId={logId}): {ex.Message}");
  674. }
  675. }
  676. private async Task<long> InsertTransformRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  677. {
  678. await _db.Ado.ExecuteCommandAsync(
  679. """
  680. INSERT INTO mdp_transform_run_log
  681. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  682. VALUES (@TenantId, 'S2_MDP_SYNC_TRANSFORM', 'S2 MDP同步与KPI计算', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  683. """,
  684. new SugarParameter("@TenantId", RequireScope().TenantId),
  685. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  686. new SugarParameter("@BatchId", batchId),
  687. new SugarParameter("@StartTime", startedAt));
  688. return await _db.Ado.GetLongAsync(
  689. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  690. new List<SugarParameter> { new("@BatchId", batchId) });
  691. }
  692. private async Task MarkTransformRunSuccessAsync(long runLogId, DateTime startedAt, S2MdpSyncTransformResult result)
  693. {
  694. var finishedAt = DateTime.Now;
  695. await _db.Ado.ExecuteCommandAsync(
  696. """
  697. UPDATE mdp_transform_run_log
  698. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  699. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=@DwdRows,
  700. summary_json=@SummaryJson, update_time=CURRENT_TIMESTAMP
  701. WHERE id=@Id
  702. """,
  703. new SugarParameter("@EndTime", finishedAt),
  704. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  705. new SugarParameter("@StageRows", result.StageRows),
  706. new SugarParameter("@StandardRows", result.StandardRows),
  707. new SugarParameter("@DwdRows", result.DwdRows),
  708. new SugarParameter("@SummaryJson", BuildRunSummaryJson(result)),
  709. new SugarParameter("@Id", runLogId));
  710. }
  711. private async Task MarkTransformRunFailedAsync(long runLogId, DateTime startedAt, string message)
  712. {
  713. try
  714. {
  715. var finishedAt = DateTime.Now;
  716. await _db.Ado.ExecuteCommandAsync(
  717. """
  718. UPDATE mdp_transform_run_log
  719. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  720. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  721. WHERE id=@Id
  722. """,
  723. new SugarParameter("@EndTime", finishedAt),
  724. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  725. new SugarParameter("@ErrorMessage", Truncate(message, 2000)),
  726. new SugarParameter("@Id", runLogId));
  727. }
  728. catch (Exception ex)
  729. {
  730. Console.Error.WriteLine($"[S2MdpSyncTransform] MarkTransformRunFailed write failed (runLogId={runLogId}): {ex.Message}");
  731. }
  732. }
  733. private S2MdpSqlCommand Cmd(string sql, string batchId, DateTime now)
  734. {
  735. var scope = RequireScope();
  736. return new S2MdpSqlCommand(MdpSqlScope.InjectTenantFactory(sql), new[]
  737. {
  738. new SugarParameter("@BatchId", batchId),
  739. new SugarParameter("@Now", now),
  740. new SugarParameter("@StatDate", now.Date),
  741. new SugarParameter("@TenantId", scope.TenantId),
  742. new SugarParameter("@FactoryId", scope.FactoryId)
  743. });
  744. }
  745. private MdpRebuildScope RequireScope()
  746. {
  747. EnsureRunScope();
  748. return _runScope;
  749. }
  750. private void EnsureRunScope()
  751. {
  752. if (_runScope == null)
  753. throw new InvalidOperationException("S2 全量重算必须指定有效 tenantId/factoryId");
  754. }
  755. private static string BuildJsonObjectExpression(IEnumerable<string> columns)
  756. {
  757. var parts = columns.SelectMany(c => new[] { $"'{c.Replace("'", "''")}'", $"s.`{c}`" });
  758. return $"JSON_OBJECT({string.Join(",", parts)})";
  759. }
  760. private static string FindColumn(IEnumerable<string> columns, string expected)
  761. {
  762. return columns.First(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase));
  763. }
  764. private static string NormalizeTriggerType(string? triggerType)
  765. {
  766. return string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  767. }
  768. private static string BuildRunSummaryJson(S2MdpSyncTransformResult result)
  769. {
  770. return $$"""{"batchId":"{{result.BatchId}}","stageRows":{{result.StageRows}},"standardRows":{{result.StandardRows}},"dwdRows":{{result.DwdRows}},"kpiRows":{{result.KpiRows}}}""";
  771. }
  772. private static string ResolveKpiValueTable(int metricLevel)
  773. {
  774. return metricLevel switch
  775. {
  776. 1 => "ado_s9_kpi_value_l1_day",
  777. 2 => "ado_s9_kpi_value_l2_day",
  778. 3 => "ado_s9_kpi_value_l3_day",
  779. 4 => "ado_s9_kpi_value_l4_day",
  780. _ => "ado_s9_kpi_value_l2_day"
  781. };
  782. }
  783. private static decimal DefaultS2Target(string metricCode) => LegacyKpiCodeTargets.GetOrZero(metricCode);
  784. private static string ResolveKpiStatus(decimal actual, decimal target, string? direction, decimal? yellowThreshold, decimal? redThreshold)
  785. {
  786. if (target <= 0) return "gray";
  787. var ratio = actual / target * 100m;
  788. if (string.Equals(direction, "lower_is_better", StringComparison.OrdinalIgnoreCase))
  789. {
  790. if (actual <= target) return "green";
  791. if (ratio <= (yellowThreshold ?? 110m)) return "yellow";
  792. return ratio >= (redThreshold ?? 120m) ? "red" : "yellow";
  793. }
  794. if (actual >= target) return "green";
  795. if (ratio >= (yellowThreshold ?? 95m)) return "yellow";
  796. return ratio <= (redThreshold ?? 80m) ? "red" : "yellow";
  797. }
  798. private static string ResolveTrendFlag(decimal actual, decimal? previous)
  799. {
  800. if (previous == null) return "flat";
  801. if (actual > previous.Value) return "up";
  802. if (actual < previous.Value) return "down";
  803. return "flat";
  804. }
  805. private static string Truncate(string? raw, int maxLength)
  806. {
  807. if (string.IsNullOrEmpty(raw)) return string.Empty;
  808. return raw.Length <= maxLength ? raw : raw[..maxLength];
  809. }
  810. private sealed class S2ColumnRow
  811. {
  812. public string ColumnName { get; set; } = string.Empty;
  813. }
  814. private sealed class S2MdpEntityRow
  815. {
  816. public long Id { get; set; }
  817. public string EntityName { get; set; } = string.Empty;
  818. }
  819. private sealed class S2KpiCalcRow
  820. {
  821. public long TenantId { get; set; }
  822. public long FactoryId { get; set; }
  823. public string MetricCode { get; set; } = string.Empty;
  824. public decimal? MetricValue { get; set; }
  825. }
  826. private sealed class S2KpiMetaRow
  827. {
  828. public int MetricLevel { get; set; }
  829. public string Direction { get; set; } = "higher_is_better";
  830. public decimal? YellowThreshold { get; set; }
  831. public decimal? RedThreshold { get; set; }
  832. }
  833. private sealed class S2KpiValueRow
  834. {
  835. public long Id { get; set; }
  836. public decimal? MetricValue { get; set; }
  837. public decimal? TargetValue { get; set; }
  838. }
  839. }
  840. public sealed class S2MdpSyncTransformResult
  841. {
  842. public long RunLogId { get; set; }
  843. public string BatchId { get; set; } = string.Empty;
  844. public int StageRows { get; set; }
  845. public int StandardRows { get; set; }
  846. public int DwdRows { get; set; }
  847. public int KpiRows { get; set; }
  848. public int AtomicRows { get; set; }
  849. }
  850. internal sealed record S2MdpSqlCommand(string Sql, SugarParameter[] Parameters);
  851. internal sealed record S2MdpEntityConfig(
  852. string EntityCode,
  853. string SourceTable,
  854. string TargetTable,
  855. string SourceRowIdExpression,
  856. string SourceBizKeyExpression)
  857. {
  858. public static readonly IReadOnlyList<S2MdpEntityConfig> All = new List<S2MdpEntityConfig>
  859. {
  860. new("S2_WORK_ORDER_MASTER", "WorkOrdMaster", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrd`,''))"),
  861. new("S2_WORK_ORDER_ROUTING", "WorkOrdRouting", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrd`,''), ':', IFNULL(s.`OP`,''))"),
  862. new("S2_WORK_ORDER_DETAIL", "WorkOrdDetail", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrd`,''), ':', IFNULL(s.`Op`,''), ':', IFNULL(s.`ItemNum`,''))"),
  863. new("S2_PERIOD_SEQUENCE_DET", "PeriodSequenceDet", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrds`,''), ':', IFNULL(s.`Op`,''), ':', IFNULL(s.`Sequence`,''))"),
  864. new("S2_SCHEDULE_RESULT_OP", "ScheduleResultOpMaster", "mdp_stg_schedule", "RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`WorkOrd`,''), ':', IFNULL(s.`Op`,''), ':', IFNULL(s.`WorkDate`,''))")
  865. };
  866. }
  867. internal static class S2MdpDdl
  868. {
  869. public static readonly IReadOnlyList<string> SqlBlocks = new[]
  870. {
  871. """
  872. CREATE TABLE IF NOT EXISTS mdp_stg_schedule (
  873. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  874. tenant_id BIGINT NOT NULL DEFAULT 0,
  875. factory_id VARCHAR(64) NULL,
  876. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  877. source_table VARCHAR(100) NOT NULL,
  878. source_row_id VARCHAR(100) NOT NULL,
  879. source_biz_key VARCHAR(200) NULL,
  880. sync_batch_id VARCHAR(100) NOT NULL,
  881. sync_time DATETIME NOT NULL,
  882. process_status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
  883. process_message VARCHAR(500) NULL,
  884. raw_data JSON NOT NULL,
  885. create_time DATETIME NULL DEFAULT CURRENT_TIMESTAMP,
  886. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  887. UNIQUE KEY uk_mdp_stg_schedule (tenant_id, source_table, source_row_id),
  888. UNIQUE KEY uk_source_key (source_system, source_table, source_biz_key),
  889. KEY idx_mdp_stg_schedule_batch (sync_batch_id),
  890. KEY idx_mdp_stg_schedule_biz (source_biz_key)
  891. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S2生产排程贴源层';
  892. """,
  893. """
  894. CREATE TABLE IF NOT EXISTS mdp_std_work_order_schedule (
  895. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  896. tenant_id BIGINT NOT NULL DEFAULT 0,
  897. factory_id BIGINT NULL,
  898. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  899. work_order VARCHAR(100) NOT NULL,
  900. sales_order_no VARCHAR(100) NULL,
  901. item_code VARCHAR(100) NULL,
  902. item_name VARCHAR(200) NULL,
  903. site_code VARCHAR(50) NULL,
  904. status VARCHAR(50) NULL,
  905. priority DECIMAL(18,6) NULL,
  906. urgent_flag TINYINT NOT NULL DEFAULT 0,
  907. qty_ordered DECIMAL(18,6) NULL,
  908. qty_completed DECIMAL(18,6) NULL,
  909. order_date DATETIME NULL,
  910. due_date DATETIME NULL,
  911. release_date DATETIME NULL,
  912. prod_line VARCHAR(100) NULL,
  913. specification VARCHAR(200) NULL,
  914. lot_serial VARCHAR(100) NULL,
  915. drawing_no VARCHAR(100) NULL,
  916. project VARCHAR(200) NULL,
  917. work_order_type VARCHAR(50) NULL,
  918. labor_variance DECIMAL(18,6) NULL,
  919. source_biz_key VARCHAR(200) NULL,
  920. sync_batch_id VARCHAR(100) NOT NULL,
  921. sync_time DATETIME NOT NULL,
  922. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  923. UNIQUE KEY uk_std_work_order_schedule (tenant_id, work_order),
  924. KEY idx_std_work_order_schedule_batch (sync_batch_id),
  925. KEY idx_std_work_order_schedule_due (tenant_id, due_date)
  926. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S2标准工单排程';
  927. """,
  928. """
  929. CREATE TABLE IF NOT EXISTS mdp_std_operation_schedule (
  930. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  931. tenant_id BIGINT NOT NULL DEFAULT 0,
  932. factory_id BIGINT NULL,
  933. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  934. work_order VARCHAR(100) NOT NULL,
  935. op_no VARCHAR(50) NULL,
  936. work_center VARCHAR(100) NULL,
  937. line_code VARCHAR(100) NULL,
  938. item_code VARCHAR(100) NULL,
  939. plan_date DATETIME NULL,
  940. prod_date DATETIME NULL,
  941. start_time DATETIME NULL,
  942. end_time DATETIME NULL,
  943. ord_qty DECIMAL(18,6) NULL,
  944. comp_qty DECIMAL(18,6) NULL,
  945. run_crew DECIMAL(18,6) NULL,
  946. employee VARCHAR(200) NULL,
  947. source_biz_key VARCHAR(200) NULL,
  948. sync_batch_id VARCHAR(100) NOT NULL,
  949. sync_time DATETIME NOT NULL,
  950. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  951. UNIQUE KEY uk_std_operation_schedule (tenant_id, source_biz_key),
  952. KEY idx_std_operation_schedule_work_order (tenant_id, work_order),
  953. KEY idx_std_operation_schedule_batch (sync_batch_id)
  954. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S2标准工序排程';
  955. """,
  956. """
  957. CREATE TABLE IF NOT EXISTS dwd_order_schedule_trans (
  958. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  959. tenant_id BIGINT NOT NULL DEFAULT 0,
  960. factory_id BIGINT NOT NULL DEFAULT 1,
  961. stat_date DATE NOT NULL,
  962. work_order VARCHAR(100) NOT NULL,
  963. sales_order_no VARCHAR(100) NULL,
  964. item_code VARCHAR(100) NULL,
  965. item_name VARCHAR(200) NULL,
  966. site_code VARCHAR(50) NULL,
  967. prod_line VARCHAR(100) NULL,
  968. status VARCHAR(50) NULL,
  969. urgent_flag TINYINT NOT NULL DEFAULT 0,
  970. qty_ordered DECIMAL(18,6) NULL,
  971. qty_completed DECIMAL(18,6) NULL,
  972. order_date DATETIME NULL,
  973. due_date DATETIME NULL,
  974. release_date DATETIME NULL,
  975. first_plan_date DATETIME NULL,
  976. last_plan_date DATETIME NULL,
  977. first_start_time DATETIME NULL,
  978. last_end_time DATETIME NULL,
  979. operation_count INT NOT NULL DEFAULT 0,
  980. scheduled_qty DECIMAL(18,6) NULL,
  981. completed_op_qty DECIMAL(18,6) NULL,
  982. schedule_cycle_days DECIMAL(18,6) NULL,
  983. schedule_satisfaction_flag TINYINT NOT NULL DEFAULT 0,
  984. wip_qty DECIMAL(18,6) NULL,
  985. resource_person_count DECIMAL(18,6) NULL,
  986. calc_batch_id VARCHAR(100) NOT NULL,
  987. calc_time DATETIME NOT NULL,
  988. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  989. UNIQUE KEY uk_dwd_order_schedule_trans (tenant_id, work_order, calc_batch_id),
  990. KEY idx_dwd_order_schedule_trans_batch (calc_batch_id),
  991. KEY idx_dwd_order_schedule_trans_stat (tenant_id, stat_date),
  992. KEY idx_dwd_order_schedule_trans_order (tenant_id, sales_order_no)
  993. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S2订单工单排程DWD';
  994. """,
  995. """
  996. INSERT INTO mdp_entity
  997. (tenant_id, source_id, entity_code, entity_name, entity_type, source_table_name, target_table_name, sync_mode, batch_size, status, remark, create_time, update_time)
  998. SELECT 0, s.id, v.entity_code, v.entity_name, 'TABLE', v.source_table_name, 'mdp_stg_schedule', 'FULL', 5000, 1, v.remark, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  999. FROM mdp_source s
  1000. JOIN (
  1001. SELECT 'S2_WORK_ORDER_MASTER' AS entity_code, 'S2工单主数据' AS entity_name, 'WorkOrdMaster' AS source_table_name, '工单主数据进入 S2 贴源层' AS remark
  1002. UNION ALL SELECT 'S2_WORK_ORDER_ROUTING', 'S2工单工艺路线', 'WorkOrdRouting', '工单工艺路线进入 S2 贴源层'
  1003. UNION ALL SELECT 'S2_WORK_ORDER_DETAIL', 'S2工单物料明细', 'WorkOrdDetail', '工单物料需求进入 S2 贴源层'
  1004. UNION ALL SELECT 'S2_PERIOD_SEQUENCE_DET', 'S2工序排程计划', 'PeriodSequenceDet', '工序间衔接与排程计划进入 S2 贴源层'
  1005. UNION ALL SELECT 'S2_SCHEDULE_RESULT_OP', 'S2工序排产结果', 'ScheduleResultOpMaster', '工序排产结果进入 S2 贴源层'
  1006. ) v
  1007. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  1008. ON DUPLICATE KEY UPDATE
  1009. source_id=VALUES(source_id), entity_name=VALUES(entity_name), entity_type=VALUES(entity_type),
  1010. source_table_name=VALUES(source_table_name), target_table_name=VALUES(target_table_name),
  1011. sync_mode=VALUES(sync_mode), batch_size=VALUES(batch_size), status=VALUES(status),
  1012. remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  1013. """
  1014. };
  1015. }