S1MdpSyncTransformService.cs 85 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  2. using Admin.NET.Plugin.AiDOP.SmartOps;
  3. namespace Admin.NET.Plugin.AiDOP.Order;
  4. /// <summary>
  5. /// S1 首批 MDP 同步和标准化转换服务。
  6. /// Phase2:源→stg 经统一执行器;std/DWD/KPI 仍读 stg。
  7. /// </summary>
  8. public class S1MdpSyncTransformService : ITransient
  9. {
  10. private const string JobCode = "S1_MDP_SYNC_TRANSFORM";
  11. private readonly ISqlSugarClient _db;
  12. private readonly SmartOpsKpiAtomicBuildService _atomicBuild;
  13. private readonly MdpModuleStagingPuller _stagingPuller;
  14. public S1MdpSyncTransformService(
  15. ISqlSugarClient db,
  16. SmartOpsKpiAtomicBuildService atomicBuild,
  17. MdpModuleStagingPuller stagingPuller)
  18. {
  19. _db = db;
  20. _atomicBuild = atomicBuild;
  21. _stagingPuller = stagingPuller;
  22. }
  23. public async Task<S1MdpSyncTransformResult> RunFullAsync(CancellationToken cancellationToken = default, string triggerType = "AUTO")
  24. {
  25. cancellationToken.ThrowIfCancellationRequested();
  26. var now = DateTime.Now;
  27. var batchId = $"S1_MDP_FULL_{now:yyyyMMddHHmmss}";
  28. var runLogId = await InsertTransformRunLogAsync(batchId, now, triggerType);
  29. var result = new S1MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  30. try
  31. {
  32. await EnsureS1RuntimeObjectsAsync();
  33. result.StageRows = await SyncStagingAsync(batchId, now, cancellationToken);
  34. result.StandardRows = await TransformStandardAsync(batchId, now, cancellationToken);
  35. result.DwdRows = await BuildDwdAsync(batchId, now, cancellationToken);
  36. result.KpiRows = await BuildS1KpiValuesAsync(batchId, now, cancellationToken);
  37. result.AtomicRows = await _atomicBuild.BuildOrderDeliveryDomainForAllDatesAsync(batchId, cancellationToken: cancellationToken);
  38. await MarkTransformRunSuccessAsync(runLogId, now, result);
  39. return result;
  40. }
  41. catch (Exception ex)
  42. {
  43. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  44. throw;
  45. }
  46. }
  47. private async Task EnsureS1RuntimeObjectsAsync()
  48. {
  49. await _db.Ado.ExecuteCommandAsync(
  50. """
  51. CREATE TABLE IF NOT EXISTS dwd_requirement_examine_detail (
  52. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  53. tenant_id BIGINT NOT NULL,
  54. factory_id BIGINT NULL,
  55. stat_date DATE NOT NULL,
  56. row_id BIGINT NOT NULL,
  57. parent_row_id BIGINT NULL,
  58. examine_id BIGINT NULL,
  59. order_entry_id BIGINT NULL,
  60. bill_no VARCHAR(100) NULL,
  61. morder_no VARCHAR(100) NULL,
  62. num VARCHAR(100) NULL,
  63. item_number VARCHAR(100) NULL,
  64. item_name VARCHAR(200) NULL,
  65. bom_number VARCHAR(100) NULL,
  66. model VARCHAR(200) NULL,
  67. kitting_time DATE NULL,
  68. item_type VARCHAR(50) NULL,
  69. erp_cls_name VARCHAR(80) NULL,
  70. qty DECIMAL(18,4) NULL,
  71. wastage DECIMAL(18,4) NULL,
  72. need_count DECIMAL(18,4) NULL,
  73. sqty DECIMAL(18,4) NULL,
  74. use_qty DECIMAL(18,4) NULL,
  75. self_lack_qty DECIMAL(18,4) NULL,
  76. lack_qty DECIMAL(18,4) NULL,
  77. mo_qty DECIMAL(18,4) NULL,
  78. make_qty DECIMAL(18,4) NULL,
  79. purchase_qty DECIMAL(18,4) NULL,
  80. purchase_occupy_qty DECIMAL(18,4) NULL,
  81. satisfy_time DATE NULL,
  82. have_ic_subs VARCHAR(10) NULL,
  83. substitute_code VARCHAR(100) NULL,
  84. create_time DATETIME NULL,
  85. source_system VARCHAR(50) NOT NULL DEFAULT 'AIDOP',
  86. sync_batch_id VARCHAR(100) NOT NULL,
  87. calc_batch_id VARCHAR(100) NOT NULL,
  88. calc_time DATETIME NOT NULL,
  89. update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  90. UNIQUE KEY uk_dwd_req_exam_detail (tenant_id, row_id, calc_batch_id),
  91. KEY idx_req_exam_tenant_batch (tenant_id, calc_batch_id),
  92. KEY idx_req_exam_bill (tenant_id, bill_no),
  93. KEY idx_req_exam_item (tenant_id, item_number)
  94. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='S1需求明细核验DWD';
  95. """);
  96. await _db.Ado.ExecuteCommandAsync(
  97. """
  98. INSERT INTO mdp_entity
  99. (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)
  100. SELECT 0, s.id, 'S1_REQUIREMENT_EXAMINE_RESULT', 'S1需求核验结果主表', 'TABLE',
  101. 'b_examine_result', 'mdp_stg_so', 'FULL', 5000, 1, '需求核验结果主表,进入 S1 贴源层', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  102. FROM mdp_source s
  103. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  104. LIMIT 1
  105. ON DUPLICATE KEY UPDATE
  106. source_id=VALUES(source_id), entity_name=VALUES(entity_name), source_table_name=VALUES(source_table_name),
  107. target_table_name=VALUES(target_table_name), sync_mode=VALUES(sync_mode), status=VALUES(status),
  108. remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  109. """);
  110. await _db.Ado.ExecuteCommandAsync(
  111. """
  112. INSERT INTO mdp_entity
  113. (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)
  114. SELECT 0, s.id, 'S1_REQUIREMENT_EXAMINE_DETAIL', 'S1需求核验BOM明细', 'TABLE',
  115. 'b_bom_child_examine', 'mdp_stg_so', 'FULL', 5000, 1, '需求核验BOM明细,进入 S1 贴源层和 DWD', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  116. FROM mdp_source s
  117. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  118. LIMIT 1
  119. ON DUPLICATE KEY UPDATE
  120. source_id=VALUES(source_id), entity_name=VALUES(entity_name), source_table_name=VALUES(source_table_name),
  121. target_table_name=VALUES(target_table_name), sync_mode=VALUES(sync_mode), status=VALUES(status),
  122. remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  123. """);
  124. await _db.Ado.ExecuteCommandAsync(
  125. """
  126. INSERT INTO mdp_entity
  127. (tenant_id, source_id, entity_code, entity_name, entity_type, source_table_name, target_table_name, sync_mode, incr_column, batch_size, status, remark, create_time, update_time)
  128. SELECT 0, s.id, v.entity_code, v.entity_name, 'TABLE', v.source_table_name, 'mdp_stg_so', v.sync_mode, v.incr_column, 5000, 1, v.remark, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  129. FROM mdp_source s
  130. JOIN (
  131. SELECT 'S1_PRODUCT_DESIGN' AS entity_code, 'S1产品设计主表' AS entity_name, 'ado_product_design' AS source_table_name, 'INCR' AS sync_mode, 'UpdateTime' AS incr_column, '产品设计主表,进入订单标准层' AS remark
  132. UNION ALL SELECT 'S1_PRODUCT_DESIGN_BOM', 'S1产品设计BOM', 'ado_product_design_bom', 'FULL', NULL, '产品设计BOM,表无时间戳列,只能全量抽取'
  133. UNION ALL SELECT 'S1_PRODUCT_DESIGN_ROUTING', 'S1产品设计工艺路线', 'ado_product_design_routing', 'FULL', NULL, '产品设计工艺路线,表无时间戳列,只能全量抽取'
  134. ) v
  135. WHERE s.tenant_id=0 AND s.source_code='AIDOPDEV_MYSQL'
  136. ON DUPLICATE KEY UPDATE
  137. source_id=VALUES(source_id), entity_name=VALUES(entity_name), source_table_name=VALUES(source_table_name),
  138. target_table_name=VALUES(target_table_name), sync_mode=VALUES(sync_mode), incr_column=VALUES(incr_column),
  139. status=VALUES(status), remark=VALUES(remark), update_time=CURRENT_TIMESTAMP;
  140. """);
  141. }
  142. /// <summary>Phase2/3:源→stg 走执行器(DB/API);遗留 SyncOneEntityAsync 不再调用。</summary>
  143. private async Task<int> SyncStagingAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  144. {
  145. _ = now;
  146. var codes = S1MdpEntityConfig.All.Select(x => x.EntityCode);
  147. // 本库源表含 tenant_id;ctx tenantId=0 仅作兜底,Writer 优先取源行值。
  148. return await _stagingPuller.PullEntitiesAsync(
  149. codes, batchId, tenantId: 0, fullRefresh: true, taskCode: "S1_MDP_INBOUND", cancellationToken);
  150. }
  151. /// <summary>
  152. /// 双模式入站:仅抽 stg(可指定 DB 实体码或 *_API),再跑 std/DWD/KPI。
  153. /// </summary>
  154. public async Task<S1MdpSyncTransformResult> RunInboundAsync(
  155. IEnumerable<string>? entityCodes = null,
  156. long tenantId = 0,
  157. bool fullRefresh = false,
  158. CancellationToken cancellationToken = default)
  159. {
  160. cancellationToken.ThrowIfCancellationRequested();
  161. var now = DateTime.Now;
  162. var batchId = $"S1_MDP_IN_{now:yyyyMMddHHmmss}";
  163. var runLogId = await InsertTransformRunLogAsync(batchId, now, "INBOUND");
  164. var result = new S1MdpSyncTransformResult { BatchId = batchId, RunLogId = runLogId };
  165. try
  166. {
  167. await EnsureS1RuntimeObjectsAsync();
  168. var codes = entityCodes?.ToList() ?? S1MdpEntityConfig.All.Select(x => x.EntityCode).ToList();
  169. result.StageRows = await _stagingPuller.PullEntitiesAsync(
  170. codes, batchId, tenantId, fullRefresh, "S1_MDP_INBOUND", cancellationToken);
  171. result.StandardRows = await TransformStandardAsync(batchId, now, cancellationToken);
  172. result.DwdRows = await BuildDwdAsync(batchId, now, cancellationToken);
  173. result.KpiRows = await BuildS1KpiValuesAsync(batchId, now, cancellationToken);
  174. result.AtomicRows = await _atomicBuild.BuildOrderDeliveryDomainForAllDatesAsync(batchId, cancellationToken: cancellationToken);
  175. await MarkTransformRunSuccessAsync(runLogId, now, result);
  176. return result;
  177. }
  178. catch (Exception ex)
  179. {
  180. await MarkTransformRunFailedAsync(runLogId, now, ex.Message);
  181. throw;
  182. }
  183. }
  184. /// <summary>Phase3:已停用。源→stg 改由执行器路径;保留方法供紧急回滚对照,勿再调用。</summary>
  185. [Obsolete("Phase3: use MdpModuleStagingPuller / executors")]
  186. private async Task<int> SyncOneEntityAsync(S1MdpEntityConfig entity, string batchId, DateTime now)
  187. {
  188. var entityRow = await _db.Ado.SqlQuerySingleAsync<S1MdpEntityRow>(
  189. "SELECT id AS Id, entity_name AS EntityName FROM mdp_entity WHERE tenant_id=0 AND entity_code=@EntityCode LIMIT 1",
  190. new SugarParameter("@EntityCode", entity.EntityCode));
  191. if (entityRow == null) throw Oops.Oh($"未找到 MDP 实体配置:{entity.EntityCode}");
  192. var columns = await _db.Ado.SqlQueryAsync<S1ColumnRow>(
  193. """
  194. SELECT COLUMN_NAME AS ColumnName
  195. FROM information_schema.COLUMNS
  196. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@TableName
  197. ORDER BY ORDINAL_POSITION
  198. """,
  199. new SugarParameter("@TableName", entity.SourceTable));
  200. if (columns.Count == 0) throw Oops.Oh($"未找到源表:{entity.SourceTable}");
  201. var names = columns.Select(u => u.ColumnName).ToList();
  202. var tenantExpr = BuildOptionalColumnExpr(names, "tenant_id", "0");
  203. var factoryExpr = BuildOptionalColumnExpr(names, "factory_id", "NULL");
  204. var companyExpr = BuildOptionalColumnExpr(names, "company_id", "NULL");
  205. var sourceRowExpr = names.Any(u => string.Equals(u, entity.SourceRowIdExpression, StringComparison.OrdinalIgnoreCase))
  206. ? $"s.`{FindColumn(names, entity.SourceRowIdExpression)}`"
  207. : entity.SourceRowIdExpression;
  208. var rawDataExpr = BuildJsonObjectExpression(names);
  209. var rowsRead = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM `{entity.SourceTable}`");
  210. var logId = await InsertSyncLogAsync(entityRow.Id, entityRow.EntityName, batchId, rowsRead);
  211. var started = DateTime.Now;
  212. try
  213. {
  214. var affected = await _db.Ado.ExecuteCommandAsync(
  215. $"""
  216. INSERT INTO `{entity.TargetTable}`
  217. (tenant_id, factory_id, company_id, source_system, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time, process_status, raw_data)
  218. SELECT
  219. {tenantExpr},
  220. {factoryExpr},
  221. {companyExpr},
  222. 'AIDOP',
  223. @SourceTable,
  224. CAST({sourceRowExpr} AS CHAR),
  225. CAST(COALESCE({entity.SourceBizKeyExpression}, CAST({sourceRowExpr} AS CHAR)) AS CHAR),
  226. @BatchId,
  227. @Now,
  228. 'PENDING',
  229. {rawDataExpr}
  230. FROM `{entity.SourceTable}` s
  231. ON DUPLICATE KEY UPDATE
  232. tenant_id=VALUES(tenant_id),
  233. factory_id=VALUES(factory_id),
  234. company_id=VALUES(company_id),
  235. source_row_id=VALUES(source_row_id),
  236. sync_batch_id=VALUES(sync_batch_id),
  237. sync_time=VALUES(sync_time),
  238. process_status=VALUES(process_status),
  239. raw_data=VALUES(raw_data),
  240. update_time=CURRENT_TIMESTAMP
  241. """,
  242. new SugarParameter("@SourceTable", entity.SourceTable),
  243. new SugarParameter("@BatchId", batchId),
  244. new SugarParameter("@Now", now));
  245. await MarkSyncLogSuccessAsync(logId, started, affected);
  246. return rowsRead;
  247. }
  248. catch (Exception ex)
  249. {
  250. await MarkSyncLogFailedAsync(logId, started, ex.Message);
  251. throw;
  252. }
  253. }
  254. private async Task<int> TransformStandardAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  255. {
  256. var total = 0;
  257. foreach (var command in BuildStandardCommands(batchId, now))
  258. {
  259. cancellationToken.ThrowIfCancellationRequested();
  260. total += await _db.Ado.ExecuteCommandAsync(command.Sql, command.Parameters);
  261. }
  262. return total;
  263. }
  264. private async Task<int> BuildDwdAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  265. {
  266. var total = 0;
  267. foreach (var command in BuildDwdCommands(batchId, now))
  268. {
  269. cancellationToken.ThrowIfCancellationRequested();
  270. total += await _db.Ado.ExecuteCommandAsync(command.Sql, command.Parameters);
  271. }
  272. return total;
  273. }
  274. private async Task<int> BuildS1KpiValuesAsync(string batchId, DateTime now, CancellationToken cancellationToken)
  275. {
  276. var statDate = now.Date;
  277. var rows = await CalculateS1KpiValuesAsync(batchId, statDate);
  278. var affected = 0;
  279. foreach (var row in rows)
  280. {
  281. cancellationToken.ThrowIfCancellationRequested();
  282. affected += await UpsertS1KpiValueAsync(row, statDate, now);
  283. }
  284. return affected;
  285. }
  286. private async Task<List<S1KpiCalcRow>> CalculateS1KpiValuesAsync(string batchId, DateTime statDate)
  287. {
  288. return await _db.Ado.SqlQueryAsync<S1KpiCalcRow>(
  289. """
  290. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  291. 'S1_L1_001' AS MetricCode,
  292. ROUND(AVG(TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) / 24), 4) AS MetricValue
  293. FROM mdp_std_so
  294. WHERE order_date IS NOT NULL
  295. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  296. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) >= order_date
  297. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  298. UNION ALL
  299. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  300. 'S1_L1_002' AS MetricCode,
  301. ROUND(100 * SUM(CASE WHEN TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) <= 72 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  302. FROM mdp_std_so
  303. WHERE order_date IS NOT NULL
  304. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  305. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  306. UNION ALL
  307. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  308. 'S1_L1_003' AS MetricCode,
  309. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(planner_no, '')), 1), 4) AS MetricValue
  310. FROM mdp_std_so
  311. WHERE order_date IS NOT NULL AND order_date <= @StatDate
  312. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  313. UNION ALL
  314. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  315. 'S1_L1_004' AS MetricCode,
  316. ROUND(AVG(GREATEST(IFNULL(remaining_qty, 0), 0)) / GREATEST(SUM(IFNULL(shipped_qty, 0)), 1) * 30, 4) AS MetricValue
  317. FROM dwd_ship_trans
  318. WHERE calc_batch_id=@BatchId
  319. GROUP BY tenant_id
  320. HAVING SUM(IFNULL(shipped_qty, 0)) > 0
  321. UNION ALL
  322. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  323. 'S1_L2_010' AS MetricCode,
  324. ROUND(AVG(TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) / 24), 4) AS MetricValue
  325. FROM mdp_std_so
  326. WHERE order_date IS NOT NULL
  327. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  328. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) >= order_date
  329. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  330. UNION ALL
  331. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  332. 'S1_L2_011' AS MetricCode,
  333. ROUND(100 * SUM(CASE WHEN TIMESTAMPDIFF(HOUR, order_date, COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date)) <= 72 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  334. FROM mdp_std_so
  335. WHERE order_date IS NOT NULL
  336. AND COALESCE(promised_delivery_date, capacity_date, material_ready_date, plan_delivery_date) IS NOT NULL
  337. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  338. UNION ALL
  339. SELECT tenant_id AS TenantId, COALESCE(NULLIF(factory_id, 0), 1) AS FactoryId,
  340. 'S1_L2_012' AS MetricCode,
  341. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(planner_no, '')), 1), 4) AS MetricValue
  342. FROM mdp_std_so
  343. WHERE order_date IS NOT NULL AND order_date <= @StatDate
  344. GROUP BY tenant_id, COALESCE(NULLIF(factory_id, 0), 1)
  345. UNION ALL
  346. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  347. 'S1_L2_013' AS MetricCode,
  348. ROUND(100 * SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  349. FROM dwd_ship_trans
  350. WHERE calc_batch_id=@BatchId
  351. GROUP BY tenant_id
  352. UNION ALL
  353. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  354. 'S1_L2_014' AS MetricCode,
  355. ROUND(100 * SUM(CASE WHEN shipped_qty >= planned_ship_qty AND planned_ship_qty > 0 THEN 1 ELSE 0 END)
  356. / GREATEST(SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END), 1), 4) AS MetricValue
  357. FROM dwd_ship_trans
  358. WHERE calc_batch_id=@BatchId
  359. GROUP BY tenant_id
  360. HAVING SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) > 0
  361. UNION ALL
  362. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  363. 'S1_L2_015' AS MetricCode,
  364. ROUND(100 * SUM(CASE WHEN linkage_status = 'LINKED' THEN 1 ELSE 0 END)
  365. / GREATEST(SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END), 1), 4) AS MetricValue
  366. FROM dwd_ship_trans
  367. WHERE calc_batch_id=@BatchId
  368. GROUP BY tenant_id
  369. HAVING SUM(CASE WHEN planned_ship_qty > 0 THEN 1 ELSE 0 END) > 0
  370. UNION ALL
  371. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  372. 'S1_L2_001' AS MetricCode,
  373. ROUND(AVG(TIMESTAMPDIFF(HOUR, create_time, update_time) / 24), 4) AS MetricValue
  374. FROM (
  375. SELECT tenant_id,
  376. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f') AS create_time,
  377. COALESCE(
  378. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  379. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  380. ) AS update_time
  381. FROM mdp_stg_so
  382. WHERE source_table = 'ado_contract_review'
  383. ) t
  384. WHERE create_time IS NOT NULL AND update_time IS NOT NULL AND update_time >= create_time
  385. GROUP BY tenant_id
  386. UNION ALL
  387. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  388. 'S1_L2_002' AS MetricCode,
  389. ROUND(100 * SUM(CASE WHEN TIMESTAMPDIFF(HOUR, create_time, update_time) <= 72 THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  390. FROM (
  391. SELECT tenant_id,
  392. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f') AS create_time,
  393. COALESCE(
  394. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.UpdateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  395. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  396. ) AS update_time
  397. FROM mdp_stg_so
  398. WHERE source_table = 'ado_contract_review'
  399. ) t
  400. WHERE create_time IS NOT NULL AND update_time IS NOT NULL AND update_time >= create_time
  401. GROUP BY tenant_id
  402. UNION ALL
  403. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  404. 'S1_L2_003' AS MetricCode,
  405. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(owner_account, '')), 1), 4) AS MetricValue
  406. FROM (
  407. SELECT tenant_id,
  408. COALESCE(
  409. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ResponsibleAccount')),
  410. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))
  411. ) AS owner_account
  412. FROM mdp_stg_so
  413. WHERE source_table = 'ado_contract_review'
  414. ) t
  415. GROUP BY tenant_id
  416. UNION ALL
  417. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  418. 'S1_L2_004' AS MetricCode,
  419. ROUND(AVG(cycle_hours / 24), 4) AS MetricValue
  420. FROM (
  421. SELECT tenant_id,
  422. COALESCE(
  423. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingDesignCycle')) REGEXP '^-?[0-9]+$'
  424. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingDesignCycle')) AS DECIMAL(18,6)) END,
  425. TIMESTAMPDIFF(HOUR,
  426. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualStart')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  427. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualEnd')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  428. )
  429. ) AS cycle_hours
  430. FROM mdp_stg_so
  431. WHERE source_table = 'ado_product_design'
  432. ) t
  433. WHERE cycle_hours IS NOT NULL AND cycle_hours >= 0
  434. GROUP BY tenant_id
  435. UNION ALL
  436. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  437. 'S1_L2_005' AS MetricCode,
  438. ROUND(100 * SUM(CASE WHEN actual_end <= plan_end THEN 1 ELSE 0 END) / COUNT(1), 4) AS MetricValue
  439. FROM (
  440. SELECT tenant_id,
  441. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingPlanEnd')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f') AS plan_end,
  442. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DrawingActualEnd')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f') AS actual_end
  443. FROM mdp_stg_so
  444. WHERE source_table = 'ado_product_design'
  445. ) t
  446. WHERE plan_end IS NOT NULL AND actual_end IS NOT NULL
  447. GROUP BY tenant_id
  448. UNION ALL
  449. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  450. 'S1_L2_006' AS MetricCode,
  451. ROUND(COUNT(1) / GREATEST(COUNT(DISTINCT NULLIF(design_lead, '')), 1), 4) AS MetricValue
  452. FROM (
  453. SELECT tenant_id,
  454. COALESCE(
  455. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DesignLeadAccount')),
  456. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DesignLeadName')),
  457. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CreateUser'))
  458. ) AS design_lead
  459. FROM mdp_stg_so
  460. WHERE source_table = 'ado_product_design'
  461. ) t
  462. GROUP BY tenant_id
  463. UNION ALL
  464. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  465. CASE stage_no
  466. WHEN 1 THEN 'S1_L3_001'
  467. WHEN 2 THEN 'S1_L3_002'
  468. WHEN 3 THEN 'S1_L3_003'
  469. WHEN 4 THEN 'S1_L3_004'
  470. WHEN 5 THEN 'S1_L3_005'
  471. END AS MetricCode,
  472. ROUND(AVG(cycle_hours), 4) AS MetricValue
  473. FROM (
  474. SELECT tenant_id,
  475. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) REGEXP '^-?[0-9]+$'
  476. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) END AS stage_no,
  477. COALESCE(
  478. TIMESTAMPDIFF(HOUR,
  479. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  480. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  481. ),
  482. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  483. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  484. ) AS cycle_hours
  485. FROM mdp_stg_so
  486. WHERE source_table = 'ado_contract_review_flow'
  487. ) t
  488. WHERE stage_no BETWEEN 1 AND 5
  489. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  490. GROUP BY tenant_id, stage_no
  491. UNION ALL
  492. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  493. CASE stage_no
  494. WHEN 1 THEN 'S1_L3_101'
  495. WHEN 2 THEN 'S1_L3_102'
  496. WHEN 3 THEN 'S1_L3_103'
  497. WHEN 4 THEN 'S1_L3_104'
  498. WHEN 5 THEN 'S1_L3_105'
  499. END AS MetricCode,
  500. ROUND(AVG(cycle_hours), 4) AS MetricValue
  501. FROM (
  502. SELECT tenant_id,
  503. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) REGEXP '^-?[0-9]+$'
  504. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) AS SIGNED) END AS stage_no,
  505. COALESCE(
  506. TIMESTAMPDIFF(HOUR,
  507. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  508. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  509. ),
  510. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  511. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  512. ) AS cycle_hours
  513. FROM mdp_stg_so
  514. WHERE source_table = 'ado_contract_review_flow'
  515. ) t
  516. WHERE stage_no BETWEEN 1 AND 5
  517. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  518. GROUP BY tenant_id, stage_no
  519. UNION ALL
  520. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  521. metric_code AS MetricCode,
  522. ROUND(AVG(cycle_hours), 4) AS MetricValue
  523. FROM (
  524. SELECT tenant_id,
  525. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')), ''), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  526. WHEN 'LAW' THEN 'S1_L4_001'
  527. WHEN 'PRE_SALES' THEN 'S1_L4_002'
  528. WHEN 'MPS' THEN 'S1_L4_003'
  529. WHEN 'TEST' THEN 'S1_L4_004'
  530. WHEN '法律事务部' THEN 'S1_L4_001'
  531. WHEN '技术售前组' THEN 'S1_L4_002'
  532. WHEN '综合主计划' THEN 'S1_L4_003'
  533. WHEN '试验站' THEN 'S1_L4_004'
  534. END AS metric_code,
  535. COALESCE(
  536. TIMESTAMPDIFF(HOUR,
  537. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  538. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  539. ),
  540. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  541. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  542. ) AS cycle_hours
  543. FROM mdp_stg_so
  544. WHERE source_table = 'ado_contract_review_flow'
  545. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) = '1'
  546. ) t
  547. WHERE metric_code IS NOT NULL
  548. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  549. GROUP BY tenant_id, metric_code
  550. UNION ALL
  551. SELECT tenant_id AS TenantId, 1 AS FactoryId,
  552. metric_code AS MetricCode,
  553. ROUND(AVG(cycle_hours), 4) AS MetricValue
  554. FROM (
  555. SELECT tenant_id,
  556. CASE COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.DeptNo')), ''), JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Department')))
  557. WHEN 'LAW' THEN 'S1_L4_101'
  558. WHEN 'PRE_SALES' THEN 'S1_L4_102'
  559. WHEN 'MPS' THEN 'S1_L4_103'
  560. WHEN 'TEST' THEN 'S1_L4_104'
  561. WHEN '法律事务部' THEN 'S1_L4_101'
  562. WHEN '技术售前组' THEN 'S1_L4_102'
  563. WHEN '综合主计划' THEN 'S1_L4_103'
  564. WHEN '试验站' THEN 'S1_L4_104'
  565. END AS metric_code,
  566. COALESCE(
  567. TIMESTAMPDIFF(HOUR,
  568. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StartTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  569. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.CompleteTime')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f')
  570. ),
  571. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  572. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ActualDays')) AS DECIMAL(18,6)) * 24 END
  573. ) AS cycle_hours
  574. FROM mdp_stg_so
  575. WHERE source_table = 'ado_contract_review_flow'
  576. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.StageNo')) = '1'
  577. ) t
  578. WHERE metric_code IS NOT NULL
  579. AND cycle_hours IS NOT NULL AND cycle_hours >= 0
  580. GROUP BY tenant_id, metric_code
  581. """,
  582. new SugarParameter("@BatchId", batchId),
  583. new SugarParameter("@StatDate", statDate));
  584. }
  585. private async Task<int> UpsertS1KpiValueAsync(S1KpiCalcRow row, DateTime statDate, DateTime now)
  586. {
  587. var meta = await _db.Ado.SqlQuerySingleAsync<S1KpiMetaRow>(
  588. """
  589. SELECT MetricLevel, Direction, YellowThreshold, RedThreshold
  590. FROM ado_smart_ops_kpi_master
  591. WHERE TenantId=@TenantId AND ModuleCode='S1' AND MetricCode=@MetricCode AND IsEnabled=1
  592. LIMIT 1
  593. """,
  594. new SugarParameter("@TenantId", row.TenantId),
  595. new SugarParameter("@MetricCode", row.MetricCode));
  596. if (meta == null || row.MetricValue == null)
  597. return 0;
  598. var table = ResolveKpiValueTable(meta.MetricLevel);
  599. var current = await _db.Ado.SqlQuerySingleAsync<S1KpiValueRow>(
  600. $"""
  601. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  602. FROM {table}
  603. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  604. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  605. ORDER BY id
  606. LIMIT 1
  607. """,
  608. new SugarParameter("@TenantId", row.TenantId),
  609. new SugarParameter("@FactoryId", row.FactoryId),
  610. new SugarParameter("@MetricCode", row.MetricCode),
  611. new SugarParameter("@BizDate", statDate));
  612. var prior = await _db.Ado.SqlQuerySingleAsync<S1KpiValueRow>(
  613. $"""
  614. SELECT id AS Id, metric_value AS MetricValue, target_value AS TargetValue
  615. FROM {table}
  616. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  617. AND metric_code=@MetricCode AND biz_date<@BizDate AND is_deleted=0
  618. ORDER BY biz_date DESC, id DESC
  619. LIMIT 1
  620. """,
  621. new SugarParameter("@TenantId", row.TenantId),
  622. new SugarParameter("@FactoryId", row.FactoryId),
  623. new SugarParameter("@MetricCode", row.MetricCode),
  624. new SugarParameter("@BizDate", statDate));
  625. var actual = Math.Round(row.MetricValue.Value, 4);
  626. var target = current?.TargetValue ?? prior?.TargetValue ?? DefaultS1Target(row.MetricCode);
  627. var status = ResolveKpiStatus(actual, target, meta.Direction, meta.YellowThreshold, meta.RedThreshold);
  628. var trend = ResolveTrendFlag(actual, prior?.MetricValue);
  629. if (current != null)
  630. {
  631. return await _db.Ado.ExecuteCommandAsync(
  632. $"""
  633. UPDATE {table}
  634. SET metric_value=@MetricValue, target_value=@TargetValue, status_color=@StatusColor, trend_flag=@TrendFlag,
  635. is_active=1, status='ACTIVE', calc_time=@CalcTime, update_time=@CalcTime
  636. WHERE tenant_id=@TenantId AND factory_id=@FactoryId AND module_code='S1'
  637. AND metric_code=@MetricCode AND biz_date=@BizDate AND is_deleted=0
  638. """,
  639. new SugarParameter("@MetricValue", actual),
  640. new SugarParameter("@TargetValue", target),
  641. new SugarParameter("@StatusColor", status),
  642. new SugarParameter("@TrendFlag", trend),
  643. new SugarParameter("@CalcTime", now),
  644. new SugarParameter("@TenantId", row.TenantId),
  645. new SugarParameter("@FactoryId", row.FactoryId),
  646. new SugarParameter("@MetricCode", row.MetricCode),
  647. new SugarParameter("@BizDate", statDate));
  648. }
  649. var nextId = await _db.Ado.GetLongAsync($"SELECT COALESCE(MAX(id), 0) + 1 FROM {table}");
  650. return await _db.Ado.ExecuteCommandAsync(
  651. $"""
  652. INSERT INTO {table}
  653. (id, tenant_id, factory_id, status, biz_date, create_time, update_time, is_deleted, is_active,
  654. module_code, metric_code, metric_value, target_value, status_color, trend_flag, calc_time)
  655. VALUES
  656. (@Id, @TenantId, @FactoryId, 'ACTIVE', @BizDate, @CalcTime, @CalcTime, 0, 1,
  657. 'S1', @MetricCode, @MetricValue, @TargetValue, @StatusColor, @TrendFlag, @CalcTime)
  658. """,
  659. new SugarParameter("@Id", nextId),
  660. new SugarParameter("@TenantId", row.TenantId),
  661. new SugarParameter("@FactoryId", row.FactoryId),
  662. new SugarParameter("@BizDate", statDate),
  663. new SugarParameter("@CalcTime", now),
  664. new SugarParameter("@MetricCode", row.MetricCode),
  665. new SugarParameter("@MetricValue", actual),
  666. new SugarParameter("@TargetValue", target),
  667. new SugarParameter("@StatusColor", status),
  668. new SugarParameter("@TrendFlag", trend));
  669. }
  670. private IEnumerable<S1MdpSqlCommand> BuildStandardCommands(string batchId, DateTime now)
  671. {
  672. yield return Cmd(
  673. """
  674. INSERT INTO mdp_std_so
  675. (tenant_id, factory_id, company_id, source_system, order_id, order_entry_id, order_no, order_line, order_type,
  676. customer_id, customer_no, customer_name, customer_order_no, country, item_code, item_name, item_spec,
  677. map_number, map_name, bom_number, unit, order_qty, delivered_notice_qty, delivered_qty, price, tax_price,
  678. amount, total_amount, order_date, customer_request_date, plan_delivery_date, promised_delivery_date,
  679. capacity_date, material_ready_date, planner_no, planner_name, order_status, review_status, review_stage,
  680. flow_state, progress, urgent, closed, deleted_flag, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  681. SELECT
  682. COALESCE(NULLIF(e.tenant_id, 0), NULLIF(h.tenant_id, 0),
  683. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  684. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tenant_id')) AS SIGNED) END,
  685. 0),
  686. COALESCE(NULLIF(e.factory_id, 0), NULLIF(h.factory_id, 0)),
  687. NULLIF(e.company_id, 0),
  688. 'AIDOP',
  689. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.Id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.Id')) AS SIGNED) END,
  690. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.Id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.Id')) AS SIGNED) END,
  691. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no')), e.source_biz_key),
  692. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.entry_seq')), JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.Id'))) AS CHAR),
  693. CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.order_type')) AS CHAR),
  694. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_id')) AS SIGNED) END,
  695. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_no')),
  696. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.custom_name')),
  697. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.custom_order_bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_from'))),
  698. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.country')),
  699. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.item_number')),
  700. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.item_name')),
  701. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.specification')),
  702. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.map_number')),
  703. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.map_name')),
  704. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bom_number')),
  705. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.unit')),
  706. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.qty')) AS DECIMAL(18,6)) END,
  707. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.deliver_notice_count')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.deliver_notice_count')) AS DECIMAL(18,6)) END,
  708. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.deliver_count')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.deliver_count')) AS DECIMAL(18,6)) END,
  709. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.price')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.price')) AS DECIMAL(18,6)) END,
  710. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tax_price')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.tax_price')) AS DECIMAL(18,6)) END,
  711. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.amount')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.amount')) AS DECIMAL(18,6)) END,
  712. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.total_amount')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.total_amount')) AS DECIMAL(18,6)) END,
  713. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.date')), 'null'), ''),
  714. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.rdate')), 'null'), ''),
  715. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.plan_date')), 'null'), ''),
  716. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.date')), 'null'), ''),
  717. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.sys_capacity_date')), 'null'), ''),
  718. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.sys_material_date')), 'null'), ''),
  719. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.planner_no')),
  720. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.planner_name')),
  721. CASE WHEN JSON_EXTRACT(h.raw_data,'$.closed') IN (1, true) THEN 'CLOSED' ELSE 'OPEN' END,
  722. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.FlowStatus')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.flowstate'))),
  723. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.CurrentDept')), JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.CurrentStage'))),
  724. JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.flowstate')),
  725. JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.progress')),
  726. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(COALESCE(e.raw_data, h.raw_data),'$.urgent')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(COALESCE(e.raw_data, h.raw_data),'$.urgent')) AS SIGNED) END,
  727. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.closed')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.closed')) AS SIGNED) END,
  728. COALESCE(CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(COALESCE(e.raw_data, h.raw_data),'$.IsDeleted')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(COALESCE(e.raw_data, h.raw_data),'$.IsDeleted')) AS SIGNED) END, 0),
  729. e.source_table,
  730. e.source_row_id,
  731. e.source_biz_key,
  732. @BatchId,
  733. @Now
  734. FROM mdp_stg_so e
  735. LEFT JOIN mdp_stg_so h ON h.source_table='crm_seorder'
  736. AND JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.Id')) = JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.seorder_id'))
  737. LEFT JOIN mdp_stg_so r ON r.source_table='ado_contract_review'
  738. AND JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.BillNo')) = COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.contract_no')), JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no')))
  739. WHERE e.source_table='crm_seorderentry'
  740. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(e.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(h.raw_data,'$.bill_no'))), '') <> ''
  741. ON DUPLICATE KEY UPDATE
  742. customer_no=VALUES(customer_no), customer_name=VALUES(customer_name), item_code=VALUES(item_code),
  743. item_name=VALUES(item_name), item_spec=VALUES(item_spec), order_qty=VALUES(order_qty),
  744. delivered_notice_qty=VALUES(delivered_notice_qty), delivered_qty=VALUES(delivered_qty),
  745. plan_delivery_date=VALUES(plan_delivery_date), promised_delivery_date=VALUES(promised_delivery_date),
  746. capacity_date=VALUES(capacity_date), material_ready_date=VALUES(material_ready_date),
  747. order_status=VALUES(order_status), review_status=VALUES(review_status), review_stage=VALUES(review_stage),
  748. flow_state=VALUES(flow_state), progress=VALUES(progress), deleted_flag=VALUES(deleted_flag),
  749. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  750. """, batchId, now);
  751. yield return Cmd(
  752. """
  753. INSERT INTO mdp_std_ship_trans
  754. (tenant_id, factory_id, company_id, source_system, trans_type, plan_id, plan_no, plan_line, order_id, order_entry_id,
  755. order_no, order_line, customer_no, customer_name, country, item_code, item_name, item_spec, qty, plan_qty,
  756. weight, volume, order_date, plan_ship_date, shipping_site, shipping_address, consignee, telephone,
  757. status, confirm_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  758. SELECT
  759. COALESCE(NULLIF(d.tenant_id, 0),
  760. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  761. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END,
  762. 0),
  763. NULLIF(d.factory_id, 0),
  764. NULLIF(d.company_id, 0),
  765. 'AIDOP',
  766. 'SHIP_PLAN',
  767. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id')) AS SIGNED) END,
  768. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.LotSerial')), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id')) AS CHAR)),
  769. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS CHAR),
  770. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.seorder_id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.seorder_id')) AS SIGNED) END,
  771. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sentry_id')) AS SIGNED) END,
  772. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr'))),
  773. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sentry_id')) AS CHAR),
  774. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CustomNo')),
  775. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CustomName')),
  776. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Country')),
  777. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')),
  778. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemName')),
  779. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Specification')),
  780. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Qty')) AS DECIMAL(18,6)) END,
  781. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Qty')) AS DECIMAL(18,6)) END,
  782. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Weight')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Weight')) AS DECIMAL(18,6)) END,
  783. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Volume')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Volume')) AS DECIMAL(18,6)) END,
  784. NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdDate')), 'null'), ''),
  785. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingDate')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.CreateTime'))), 'null'), ''),
  786. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingSite')),
  787. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShippingAddress')),
  788. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Consignee')),
  789. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Telephone')),
  790. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Status')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Status'))),
  791. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.IsConfirm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.IsConfirm'))) AS CHAR),
  792. d.source_table,
  793. d.source_row_id,
  794. d.source_biz_key,
  795. @BatchId,
  796. @Now
  797. FROM mdp_stg_ship_trans d
  798. LEFT JOIN mdp_stg_ship_trans m ON m.source_table='ShippingPlan'
  799. AND JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.RecID')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.plan_id'))
  800. WHERE d.source_table='ShippingPlanDetail'
  801. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr'))), '') <> ''
  802. ON DUPLICATE KEY UPDATE
  803. plan_no=VALUES(plan_no), order_no=VALUES(order_no), customer_no=VALUES(customer_no), customer_name=VALUES(customer_name),
  804. item_code=VALUES(item_code), item_name=VALUES(item_name), qty=VALUES(qty), plan_qty=VALUES(plan_qty),
  805. plan_ship_date=VALUES(plan_ship_date), shipping_site=VALUES(shipping_site), shipping_address=VALUES(shipping_address),
  806. status=VALUES(status), confirm_status=VALUES(confirm_status), sync_batch_id=VALUES(sync_batch_id),
  807. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  808. """, batchId, now);
  809. yield return Cmd(
  810. """
  811. INSERT INTO mdp_std_ship_trans
  812. (tenant_id, factory_id, source_system, trans_type, shipper_rec_id, shipper_no, shipper_line, order_no, order_line,
  813. customer_no, item_code, item_name, qty_to_ship, picking_qty, real_qty, gross_weight, net_weight, volume,
  814. plan_ship_date, actual_ship_date, site, status, confirm_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  815. SELECT
  816. COALESCE(NULLIF(d.tenant_id, 0),
  817. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  818. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END,
  819. 0),
  820. NULLIF(d.factory_id, 0),
  821. 'AIDOP',
  822. 'ASN_SHIPPER',
  823. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ASNBOLShipperRecID')) REGEXP '^-?[0-9]+$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ASNBOLShipperRecID')) AS SIGNED) END,
  824. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Id')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Id'))),
  825. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS CHAR),
  826. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.OrdNbr'))),
  827. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdLine')) AS CHAR),
  828. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.SoldTo')),
  829. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ContainerItem')),
  830. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Descr')),
  831. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyToShip')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyToShip')) AS DECIMAL(18,6)) END,
  832. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.PickingQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.PickingQty')) AS DECIMAL(18,6)) END,
  833. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RealQty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RealQty')) AS DECIMAL(18,6)) END,
  834. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.GrossWeight')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.GrossWeight')) AS DECIMAL(18,6)) END,
  835. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.NetWeight')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.NetWeight')) AS DECIMAL(18,6)) END,
  836. CASE WHEN COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Volume')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Volume'))) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Volume')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Volume'))) AS DECIMAL(18,6)) END,
  837. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ShipDate')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShipDate'))), 'null'), ''),
  838. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ShipDate')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ShipDate'))), 'null'), ''),
  839. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Site')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Site'))),
  840. COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Status')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Status'))),
  841. CAST(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.IsConfirm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.IsConfirm'))) AS CHAR),
  842. d.source_table,
  843. d.source_row_id,
  844. d.source_biz_key,
  845. @BatchId,
  846. @Now
  847. FROM mdp_stg_ship_trans d
  848. LEFT JOIN mdp_stg_ship_trans m ON m.source_table='ASNBOLShipperMaster'
  849. AND JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.RecID')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ASNBOLShipperRecID'))
  850. WHERE d.source_table='ASNBOLShipperDetail'
  851. AND IFNULL(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Id')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.Id'))), '') <> ''
  852. ON DUPLICATE KEY UPDATE
  853. shipper_no=VALUES(shipper_no), order_no=VALUES(order_no), order_line=VALUES(order_line),
  854. customer_no=VALUES(customer_no), item_code=VALUES(item_code), item_name=VALUES(item_name),
  855. qty_to_ship=VALUES(qty_to_ship), picking_qty=VALUES(picking_qty), real_qty=VALUES(real_qty),
  856. actual_ship_date=VALUES(actual_ship_date), site=VALUES(site), status=VALUES(status),
  857. confirm_status=VALUES(confirm_status), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  858. update_time=CURRENT_TIMESTAMP
  859. """, batchId, now);
  860. yield return Cmd(
  861. """
  862. INSERT INTO mdp_std_ship_trans
  863. (tenant_id, factory_id, source_system, trans_type, order_no, customer_no, item_code, item_name, qty,
  864. plan_ship_date, status, linkage_status, source_table, source_row_id, source_biz_key, sync_batch_id, sync_time)
  865. SELECT
  866. COALESCE(NULLIF(d.tenant_id, 0),
  867. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  868. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END,
  869. 0),
  870. NULLIF(d.factory_id, 0),
  871. 'AIDOP',
  872. 'LINKAGE_PLAN',
  873. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')),
  874. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.custom_no')),
  875. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_number')),
  876. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Descr')),
  877. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) AS DECIMAL(18,6)) END,
  878. NULLIF(NULLIF(COALESCE(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sys_capacity_date')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.fystarttime'))), 'null'), ''),
  879. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.type')),
  880. CASE
  881. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.isuse')) = '1' THEN 'LINKED'
  882. ELSE 'INACTIVE'
  883. END,
  884. d.source_table,
  885. d.source_row_id,
  886. d.source_biz_key,
  887. @BatchId,
  888. @Now
  889. FROM mdp_stg_ship_trans d
  890. WHERE d.source_table='LinkagePlan'
  891. AND IFNULL(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bill_no')), '') <> ''
  892. ON DUPLICATE KEY UPDATE
  893. order_no=VALUES(order_no), customer_no=VALUES(customer_no), item_code=VALUES(item_code),
  894. item_name=VALUES(item_name), qty=VALUES(qty), plan_ship_date=VALUES(plan_ship_date),
  895. status=VALUES(status), linkage_status=VALUES(linkage_status), sync_batch_id=VALUES(sync_batch_id),
  896. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  897. """, batchId, now);
  898. }
  899. private IEnumerable<S1MdpSqlCommand> BuildDwdCommands(string batchId, DateTime now)
  900. {
  901. yield return Cmd(
  902. """
  903. INSERT INTO dwd_requirement_examine_detail
  904. (tenant_id, factory_id, stat_date, row_id, parent_row_id, examine_id, order_entry_id, bill_no, morder_no,
  905. num, item_number, item_name, bom_number, model, kitting_time, item_type, erp_cls_name, qty, wastage,
  906. need_count, sqty, use_qty, self_lack_qty, lack_qty, mo_qty, make_qty, purchase_qty, purchase_occupy_qty,
  907. satisfy_time, have_ic_subs, substitute_code, create_time, source_system, sync_batch_id, calc_batch_id, calc_time)
  908. SELECT
  909. COALESCE(NULLIF(d.tenant_id, 0), NULLIF(r.tenant_id, 0), NULLIF(so.tenant_id, 0),
  910. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) REGEXP '^[0-9]+$'
  911. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.tenant_id')) AS SIGNED) END,
  912. 0),
  913. COALESCE(NULLIF(d.factory_id, 0), NULLIF(r.factory_id, 0), NULLIF(so.factory_id, 0)),
  914. @StatDate,
  915. CAST(d.source_row_id AS SIGNED),
  916. CASE
  917. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.num')) = '1' THEN NULL
  918. ELSE p.parent_row_id
  919. END,
  920. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id')) REGEXP '^-?[0-9]+$'
  921. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id')) AS SIGNED) END,
  922. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$'
  923. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) AS SIGNED) END,
  924. COALESCE(so.order_no, JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.bill_no'))),
  925. JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.morder_no')),
  926. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.num')) AS CHAR),
  927. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_number')),
  928. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.item_name')),
  929. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.bom_number')),
  930. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.model')),
  931. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.kitting_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  932. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.type')) = '1' THEN '替代件' ELSE '标准件' END,
  933. COALESCE(
  934. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls_name')),
  935. CASE JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls'))
  936. WHEN '0' THEN '配置类'
  937. WHEN '1' THEN '自制'
  938. WHEN '2' THEN '委外加工'
  939. WHEN '3' THEN '外购'
  940. WHEN '4' THEN '虚拟件'
  941. END
  942. ),
  943. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.qty')) AS DECIMAL(18,4)) END,
  944. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.wastage')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.wastage')) AS DECIMAL(18,4)) END,
  945. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.needCount')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.needCount')) AS DECIMAL(18,4)) END,
  946. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sqty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.sqty')) AS DECIMAL(18,4)) END,
  947. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.use_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.use_qty')) AS DECIMAL(18,4)) END,
  948. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.self_lack_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.self_lack_qty')) AS DECIMAL(18,4)) END,
  949. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.lack_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.lack_qty')) AS DECIMAL(18,4)) END,
  950. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.mo_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.mo_qty')) AS DECIMAL(18,4)) END,
  951. CASE
  952. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) = '0' AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.erp_cls')) = '1'
  953. THEN CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.lack_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.lack_qty')) AS DECIMAL(18,4)) END
  954. WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$'
  955. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.make_qty')) AS DECIMAL(18,4))
  956. END,
  957. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.purchase_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.purchase_qty')) AS DECIMAL(18,4)) END,
  958. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.purchase_occupy_qty')) REGEXP '^-?[0-9]+(\\.[0-9]+)?$' THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.purchase_occupy_qty')) AS DECIMAL(18,4)) END,
  959. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.satisfy_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  960. CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.haveicsubs')) = '1' THEN '是' ELSE '否' END,
  961. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.substitute_code')),
  962. STR_TO_DATE(REPLACE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.create_time')), 'null'), ''), 'T', ' '), '%Y-%m-%d %H:%i:%s.%f'),
  963. 'AIDOP',
  964. d.sync_batch_id,
  965. @BatchId,
  966. @Now
  967. FROM mdp_stg_so d
  968. INNER JOIN mdp_stg_so r ON r.source_table='b_examine_result'
  969. AND r.source_row_id = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id'))
  970. LEFT JOIN mdp_std_so so ON so.tenant_id = COALESCE(r.tenant_id, d.tenant_id)
  971. AND so.order_entry_id = CASE WHEN JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) REGEXP '^-?[0-9]+$'
  972. THEN CAST(JSON_UNQUOTE(JSON_EXTRACT(r.raw_data,'$.sentry_id')) AS SIGNED) END
  973. LEFT JOIN (
  974. SELECT examine_id, MIN(source_row_id) AS parent_row_id
  975. FROM (
  976. SELECT JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.examine_id')) AS examine_id, CAST(source_row_id AS SIGNED) AS source_row_id
  977. FROM mdp_stg_so
  978. WHERE source_table='b_bom_child_examine'
  979. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.num')) = '1'
  980. AND JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.is_use')) IN ('1', 'true', 'True', 'base64:type16:AQ==')
  981. ) x
  982. GROUP BY examine_id
  983. ) p ON p.examine_id = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.examine_id'))
  984. WHERE d.source_table='b_bom_child_examine'
  985. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.is_use')) IN ('1', 'true', 'True', 'base64:type16:AQ==')
  986. ON DUPLICATE KEY UPDATE
  987. parent_row_id=VALUES(parent_row_id), order_entry_id=VALUES(order_entry_id), bill_no=VALUES(bill_no),
  988. morder_no=VALUES(morder_no), num=VALUES(num), item_number=VALUES(item_number), item_name=VALUES(item_name),
  989. bom_number=VALUES(bom_number), model=VALUES(model), kitting_time=VALUES(kitting_time),
  990. item_type=VALUES(item_type), erp_cls_name=VALUES(erp_cls_name), qty=VALUES(qty), wastage=VALUES(wastage),
  991. need_count=VALUES(need_count), sqty=VALUES(sqty), use_qty=VALUES(use_qty), self_lack_qty=VALUES(self_lack_qty),
  992. lack_qty=VALUES(lack_qty), mo_qty=VALUES(mo_qty), make_qty=VALUES(make_qty), purchase_qty=VALUES(purchase_qty),
  993. purchase_occupy_qty=VALUES(purchase_occupy_qty), satisfy_time=VALUES(satisfy_time), have_ic_subs=VALUES(have_ic_subs),
  994. substitute_code=VALUES(substitute_code), create_time=VALUES(create_time), sync_batch_id=VALUES(sync_batch_id),
  995. calc_time=VALUES(calc_time), update_time=CURRENT_TIMESTAMP
  996. """, batchId, now);
  997. yield return Cmd(
  998. """
  999. INSERT INTO dwd_ship_trans
  1000. (tenant_id, factory_id, company_id, stat_date, order_id, order_entry_id, order_no, order_line, customer_no,
  1001. customer_name, country, item_code, item_name, item_spec, order_qty, planned_ship_qty, shipped_qty,
  1002. remaining_qty, order_date, customer_request_date, plan_delivery_date, promised_delivery_date,
  1003. plan_ship_date, actual_ship_date, review_status, order_status, delivery_status, linkage_status, risk_level,
  1004. source_system, source_table, source_row_id, source_biz_key, sync_batch_id, calc_batch_id, calc_time)
  1005. SELECT
  1006. so.tenant_id,
  1007. so.factory_id,
  1008. so.company_id,
  1009. @StatDate,
  1010. so.order_id,
  1011. so.order_entry_id,
  1012. so.order_no,
  1013. IFNULL(so.order_line, ''),
  1014. so.customer_no,
  1015. so.customer_name,
  1016. so.country,
  1017. IFNULL(so.item_code, ''),
  1018. so.item_name,
  1019. so.item_spec,
  1020. IFNULL(so.order_qty, 0),
  1021. IFNULL(p.plan_qty, 0),
  1022. IFNULL(a.real_qty, 0),
  1023. GREATEST(IFNULL(so.order_qty, 0) - IFNULL(a.real_qty, 0), 0),
  1024. so.order_date,
  1025. so.customer_request_date,
  1026. so.plan_delivery_date,
  1027. so.promised_delivery_date,
  1028. p.plan_ship_date,
  1029. a.actual_ship_date,
  1030. so.review_status,
  1031. so.order_status,
  1032. CASE
  1033. WHEN IFNULL(so.order_qty, 0) > 0 AND IFNULL(a.real_qty, 0) >= IFNULL(so.order_qty, 0) THEN 'COMPLETED'
  1034. WHEN COALESCE(p.plan_ship_date, so.promised_delivery_date, so.plan_delivery_date) < @Now THEN 'DELAYED'
  1035. WHEN IFNULL(p.plan_qty, 0) > 0 THEN 'PLANNED'
  1036. ELSE 'OPEN'
  1037. END,
  1038. l.linkage_status,
  1039. CASE
  1040. WHEN COALESCE(p.plan_ship_date, so.promised_delivery_date, so.plan_delivery_date) < @Now
  1041. AND IFNULL(a.real_qty, 0) < IFNULL(so.order_qty, 0) THEN 'HIGH'
  1042. WHEN IFNULL(a.real_qty, 0) < IFNULL(so.order_qty, 0) THEN 'MEDIUM'
  1043. ELSE 'LOW'
  1044. END,
  1045. 'AIDOP',
  1046. so.source_table,
  1047. so.source_row_id,
  1048. so.source_biz_key,
  1049. so.sync_batch_id,
  1050. @BatchId,
  1051. @Now
  1052. FROM mdp_std_so so
  1053. LEFT JOIN (
  1054. SELECT tenant_id, order_no, order_entry_id, IFNULL(order_line, '') AS order_line, IFNULL(item_code, '') AS item_code,
  1055. SUM(IFNULL(plan_qty, IFNULL(qty, 0))) AS plan_qty,
  1056. MIN(plan_ship_date) AS plan_ship_date
  1057. FROM mdp_std_ship_trans
  1058. WHERE trans_type='SHIP_PLAN'
  1059. GROUP BY tenant_id, order_no, order_entry_id, IFNULL(order_line, ''), IFNULL(item_code, '')
  1060. ) p ON so.tenant_id=p.tenant_id
  1061. AND so.order_no=p.order_no
  1062. AND IFNULL(so.item_code, '')=p.item_code
  1063. AND (
  1064. (p.order_entry_id IS NOT NULL AND so.order_entry_id=p.order_entry_id)
  1065. OR (p.order_entry_id IS NULL AND IFNULL(so.order_line, '')=p.order_line)
  1066. )
  1067. LEFT JOIN (
  1068. SELECT tenant_id, order_no, IFNULL(order_line, '') AS order_line, IFNULL(item_code, '') AS item_code,
  1069. SUM(IFNULL(real_qty, IFNULL(qty_to_ship, 0))) AS real_qty,
  1070. MAX(actual_ship_date) AS actual_ship_date
  1071. FROM mdp_std_ship_trans
  1072. WHERE trans_type='ASN_SHIPPER'
  1073. GROUP BY tenant_id, order_no, IFNULL(order_line, ''), IFNULL(item_code, '')
  1074. ) a ON so.tenant_id=a.tenant_id AND so.order_no=a.order_no AND IFNULL(so.order_line, '')=a.order_line AND IFNULL(so.item_code, '')=a.item_code
  1075. LEFT JOIN (
  1076. SELECT tenant_id, order_no, item_code, MAX(linkage_status) AS linkage_status
  1077. FROM mdp_std_ship_trans
  1078. WHERE IFNULL(linkage_status, '') <> ''
  1079. GROUP BY tenant_id, order_no, item_code
  1080. ) l ON so.tenant_id=l.tenant_id AND so.order_no=l.order_no AND IFNULL(so.item_code, '')=IFNULL(l.item_code, '')
  1081. WHERE IFNULL(so.order_no, '') <> ''
  1082. ON DUPLICATE KEY UPDATE
  1083. customer_no=VALUES(customer_no), customer_name=VALUES(customer_name), item_name=VALUES(item_name),
  1084. order_qty=VALUES(order_qty), planned_ship_qty=VALUES(planned_ship_qty), shipped_qty=VALUES(shipped_qty),
  1085. remaining_qty=VALUES(remaining_qty), plan_ship_date=VALUES(plan_ship_date), actual_ship_date=VALUES(actual_ship_date),
  1086. review_status=VALUES(review_status), order_status=VALUES(order_status), delivery_status=VALUES(delivery_status),
  1087. linkage_status=VALUES(linkage_status), risk_level=VALUES(risk_level), sync_batch_id=VALUES(sync_batch_id),
  1088. calc_batch_id=VALUES(calc_batch_id), calc_time=VALUES(calc_time), update_time=CURRENT_TIMESTAMP
  1089. """, batchId, now);
  1090. }
  1091. private async Task<long> InsertSyncLogAsync(long entityId, string entityName, string batchId, int rowsRead)
  1092. {
  1093. await _db.Ado.ExecuteCommandAsync(
  1094. """
  1095. INSERT INTO mdp_sync_log
  1096. (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type, sync_start, rows_read, status)
  1097. VALUES (0, @EntityId, 'AIDOPDEV_MYSQL', @EntityName, @BatchId, 'FULL', 'AUTO', NOW(), @RowsRead, 'RUNNING')
  1098. """,
  1099. new SugarParameter("@EntityId", entityId),
  1100. new SugarParameter("@EntityName", entityName),
  1101. new SugarParameter("@BatchId", batchId),
  1102. new SugarParameter("@RowsRead", rowsRead));
  1103. return await _db.Ado.GetLongAsync(
  1104. "SELECT id FROM mdp_sync_log WHERE sync_batch_id=@BatchId AND entity_id=@EntityId ORDER BY id DESC LIMIT 1",
  1105. new List<SugarParameter>
  1106. {
  1107. new("@BatchId", batchId),
  1108. new("@EntityId", entityId)
  1109. });
  1110. }
  1111. private async Task MarkSyncLogSuccessAsync(long logId, DateTime started, int affected)
  1112. {
  1113. await _db.Ado.ExecuteCommandAsync(
  1114. """
  1115. UPDATE mdp_sync_log
  1116. SET sync_end=NOW(), duration_ms=@DurationMs, rows_insert=@RowsInsert, rows_update=0, rows_skip=0, rows_error=0, status='SUCCESS'
  1117. WHERE id=@Id
  1118. """,
  1119. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  1120. new SugarParameter("@RowsInsert", affected),
  1121. new SugarParameter("@Id", logId));
  1122. }
  1123. private async Task MarkSyncLogFailedAsync(long logId, DateTime started, string message)
  1124. {
  1125. try
  1126. {
  1127. await _db.Ado.ExecuteCommandAsync(
  1128. """
  1129. UPDATE mdp_sync_log
  1130. SET sync_end=NOW(), duration_ms=@DurationMs, rows_error=1, status='FAILED', error_msg=@ErrorMsg
  1131. WHERE id=@Id
  1132. """,
  1133. new SugarParameter("@DurationMs", (int)(DateTime.Now - started).TotalMilliseconds),
  1134. new SugarParameter("@ErrorMsg", Truncate(message, 1000)),
  1135. new SugarParameter("@Id", logId));
  1136. }
  1137. catch (Exception ex)
  1138. {
  1139. // 写库自身失败兜底:避免再抛掩盖原异常;遗留 RUNNING 行可由运维手动清理
  1140. Console.Error.WriteLine($"[S1MdpSyncTransform] MarkSyncLogFailed write failed (syncLogId={logId}): {ex.Message}");
  1141. }
  1142. }
  1143. private async Task<long> InsertTransformRunLogAsync(string batchId, DateTime startedAt, string triggerType)
  1144. {
  1145. await _db.Ado.ExecuteCommandAsync(
  1146. """
  1147. INSERT INTO mdp_transform_run_log
  1148. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  1149. VALUES (0, @JobCode, 'S1 MDP同步与标准化转换', @TriggerType, @BatchId, 'RUNNING', @StartTime)
  1150. """,
  1151. new SugarParameter("@JobCode", JobCode),
  1152. new SugarParameter("@TriggerType", NormalizeTriggerType(triggerType)),
  1153. new SugarParameter("@BatchId", batchId),
  1154. new SugarParameter("@StartTime", startedAt));
  1155. return await _db.Ado.GetLongAsync(
  1156. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId ORDER BY id DESC LIMIT 1",
  1157. new List<SugarParameter> { new("@BatchId", batchId) });
  1158. }
  1159. private async Task MarkTransformRunSuccessAsync(long runLogId, DateTime startedAt, S1MdpSyncTransformResult result)
  1160. {
  1161. var finishedAt = DateTime.Now;
  1162. await _db.Ado.ExecuteCommandAsync(
  1163. """
  1164. UPDATE mdp_transform_run_log
  1165. SET status='SUCCESS', end_time=@EndTime, duration_ms=@DurationMs,
  1166. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=@DwdRows,
  1167. summary_json=@SummaryJson, update_time=CURRENT_TIMESTAMP
  1168. WHERE id=@Id
  1169. """,
  1170. new SugarParameter("@EndTime", finishedAt),
  1171. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  1172. new SugarParameter("@StageRows", result.StageRows),
  1173. new SugarParameter("@StandardRows", result.StandardRows),
  1174. new SugarParameter("@DwdRows", result.DwdRows),
  1175. new SugarParameter("@SummaryJson", BuildRunSummaryJson(result)),
  1176. new SugarParameter("@Id", runLogId));
  1177. }
  1178. private async Task MarkTransformRunFailedAsync(long runLogId, DateTime startedAt, string message)
  1179. {
  1180. try
  1181. {
  1182. var finishedAt = DateTime.Now;
  1183. await _db.Ado.ExecuteCommandAsync(
  1184. """
  1185. UPDATE mdp_transform_run_log
  1186. SET status='FAILED', end_time=@EndTime, duration_ms=@DurationMs,
  1187. error_message=@ErrorMessage, update_time=CURRENT_TIMESTAMP
  1188. WHERE id=@Id
  1189. """,
  1190. new SugarParameter("@EndTime", finishedAt),
  1191. new SugarParameter("@DurationMs", (int)(finishedAt - startedAt).TotalMilliseconds),
  1192. new SugarParameter("@ErrorMessage", Truncate(message, 2000)),
  1193. new SugarParameter("@Id", runLogId));
  1194. }
  1195. catch (Exception ex)
  1196. {
  1197. // 写库自身失败兜底(典型场景:远端 MySQL 瞬断导致 MarkFailed 自身也连不上):
  1198. // 避免再抛二次异常掩盖原错;遗留 RUNNING 行可由运维手动清理。
  1199. Console.Error.WriteLine($"[S1MdpSyncTransform] MarkTransformRunFailed write failed (runLogId={runLogId}): {ex.Message}");
  1200. }
  1201. }
  1202. private static S1MdpSqlCommand Cmd(string sql, string batchId, DateTime now)
  1203. {
  1204. return new S1MdpSqlCommand(sql, new[]
  1205. {
  1206. new SugarParameter("@BatchId", batchId),
  1207. new SugarParameter("@Now", now),
  1208. new SugarParameter("@StatDate", now.Date)
  1209. });
  1210. }
  1211. private static string BuildJsonObjectExpression(IEnumerable<string> columns)
  1212. {
  1213. var parts = columns.SelectMany(c => new[] { $"'{c.Replace("'", "''")}'", $"s.`{c}`" });
  1214. return $"JSON_OBJECT({string.Join(",", parts)})";
  1215. }
  1216. private static string BuildOptionalColumnExpr(IReadOnlyCollection<string> columns, string expected, string fallback)
  1217. {
  1218. return columns.Any(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase))
  1219. ? $"s.`{FindColumn(columns, expected)}`"
  1220. : fallback;
  1221. }
  1222. private static string FindColumn(IEnumerable<string> columns, string expected)
  1223. {
  1224. return columns.First(u => string.Equals(u, expected, StringComparison.OrdinalIgnoreCase));
  1225. }
  1226. private static string NormalizeTriggerType(string? triggerType)
  1227. {
  1228. return string.IsNullOrWhiteSpace(triggerType) ? "AUTO" : triggerType.Trim().ToUpperInvariant();
  1229. }
  1230. private static string BuildRunSummaryJson(S1MdpSyncTransformResult result)
  1231. {
  1232. return $$"""{"batchId":"{{result.BatchId}}","stageRows":{{result.StageRows}},"standardRows":{{result.StandardRows}},"dwdRows":{{result.DwdRows}},"kpiRows":{{result.KpiRows}}}""";
  1233. }
  1234. private static string ResolveKpiValueTable(int metricLevel)
  1235. {
  1236. return metricLevel switch
  1237. {
  1238. 1 => "ado_s9_kpi_value_l1_day",
  1239. 2 => "ado_s9_kpi_value_l2_day",
  1240. 3 => "ado_s9_kpi_value_l3_day",
  1241. 4 => "ado_s9_kpi_value_l4_day",
  1242. _ => "ado_s9_kpi_value_l2_day"
  1243. };
  1244. }
  1245. private static decimal DefaultS1Target(string metricCode)
  1246. {
  1247. return metricCode switch
  1248. {
  1249. "S1_L1_001" or "S1_L2_001" or "S1_L2_004" or "S1_L2_010" => 3m,
  1250. "S1_L1_004" => 60m,
  1251. "S1_L1_002" or "S1_L2_002" or "S1_L2_005" or "S1_L2_011" or "S1_L2_013" or "S1_L2_014" or "S1_L2_015" => 95m,
  1252. "S1_L1_003" or "S1_L2_003" or "S1_L2_006" or "S1_L2_012" => 100m,
  1253. "S1_L3_001" or "S1_L3_101" => 8m,
  1254. "S1_L3_002" or "S1_L3_102" => 12m,
  1255. "S1_L3_003" or "S1_L3_103" => 8m,
  1256. "S1_L3_004" or "S1_L3_104" => 10m,
  1257. "S1_L3_005" or "S1_L3_105" => 2m,
  1258. "S1_L4_001" or "S1_L4_101" => 2m,
  1259. "S1_L4_002" or "S1_L4_102" => 3m,
  1260. "S1_L4_003" or "S1_L4_103" => 1m,
  1261. "S1_L4_004" or "S1_L4_104" => 2m,
  1262. _ => 0m
  1263. };
  1264. }
  1265. private static string ResolveKpiStatus(decimal actual, decimal target, string? direction, decimal? yellowThreshold, decimal? redThreshold)
  1266. {
  1267. if (target <= 0) return "gray";
  1268. var ratio = actual / target * 100m;
  1269. if (string.Equals(direction, "lower_is_better", StringComparison.OrdinalIgnoreCase))
  1270. {
  1271. if (actual <= target) return "green";
  1272. if (ratio <= (yellowThreshold ?? 110m)) return "yellow";
  1273. return ratio >= (redThreshold ?? 120m) ? "red" : "yellow";
  1274. }
  1275. if (actual >= target) return "green";
  1276. if (ratio >= (yellowThreshold ?? 95m)) return "yellow";
  1277. return ratio <= (redThreshold ?? 80m) ? "red" : "yellow";
  1278. }
  1279. private static string ResolveTrendFlag(decimal actual, decimal? previous)
  1280. {
  1281. if (previous == null) return "flat";
  1282. if (actual > previous.Value) return "up";
  1283. if (actual < previous.Value) return "down";
  1284. return "flat";
  1285. }
  1286. private static string Truncate(string? raw, int maxLength)
  1287. {
  1288. if (string.IsNullOrEmpty(raw)) return string.Empty;
  1289. return raw.Length <= maxLength ? raw : raw[..maxLength];
  1290. }
  1291. private sealed class S1ColumnRow
  1292. {
  1293. public string ColumnName { get; set; } = string.Empty;
  1294. }
  1295. private sealed class S1MdpEntityRow
  1296. {
  1297. public long Id { get; set; }
  1298. public string EntityName { get; set; } = string.Empty;
  1299. }
  1300. private sealed class S1KpiCalcRow
  1301. {
  1302. public long TenantId { get; set; }
  1303. public long FactoryId { get; set; }
  1304. public string MetricCode { get; set; } = string.Empty;
  1305. public decimal? MetricValue { get; set; }
  1306. }
  1307. private sealed class S1KpiMetaRow
  1308. {
  1309. public int MetricLevel { get; set; }
  1310. public string Direction { get; set; } = "higher_is_better";
  1311. public decimal? YellowThreshold { get; set; }
  1312. public decimal? RedThreshold { get; set; }
  1313. }
  1314. private sealed class S1KpiValueRow
  1315. {
  1316. public long Id { get; set; }
  1317. public decimal? MetricValue { get; set; }
  1318. public decimal? TargetValue { get; set; }
  1319. }
  1320. }
  1321. public sealed class S1MdpSyncTransformResult
  1322. {
  1323. public long RunLogId { get; set; }
  1324. public string BatchId { get; set; } = string.Empty;
  1325. public int StageRows { get; set; }
  1326. public int StandardRows { get; set; }
  1327. public int DwdRows { get; set; }
  1328. public int KpiRows { get; set; }
  1329. public int AtomicRows { get; set; }
  1330. }
  1331. internal sealed record S1MdpSqlCommand(string Sql, SugarParameter[] Parameters);
  1332. internal sealed record S1MdpEntityConfig(
  1333. string EntityCode,
  1334. string SourceTable,
  1335. string TargetTable,
  1336. string SourceRowIdExpression,
  1337. string SourceBizKeyExpression)
  1338. {
  1339. public static readonly IReadOnlyList<S1MdpEntityConfig> All = new List<S1MdpEntityConfig>
  1340. {
  1341. new("S1_SEORDER", "crm_seorder", "mdp_stg_so", "Id", "COALESCE(s.`bill_no`, CAST(s.`Id` AS CHAR))"),
  1342. new("S1_SEORDER_ENTRY", "crm_seorderentry", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`entry_seq`, CAST(s.`Id` AS CHAR)))"),
  1343. new("S1_SEORDER_CHANGE", "crm_seorder_change", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', CAST(s.`Id` AS CHAR))"),
  1344. new("S1_CONTRACT_REVIEW", "ado_contract_review", "mdp_stg_so", "RecID", "COALESCE(s.`BillNo`, CAST(s.`RecID` AS CHAR))"),
  1345. new("S1_CONTRACT_REVIEW_FLOW", "ado_contract_review_flow", "mdp_stg_so", "RecID", "CONCAT(IFNULL(s.`ReviewBillNo`,''), ':', IFNULL(s.`StageNo`,''), ':', CAST(s.`RecID` AS CHAR))"),
  1346. new("S1_PRODUCT_DESIGN", "ado_product_design", "mdp_stg_so", "Id", "COALESCE(s.`BillNo`, CAST(s.`Id` AS CHAR))"),
  1347. new("S1_PRODUCT_DESIGN_BOM", "ado_product_design_bom", "mdp_stg_so", "Id", "CONCAT(CAST(s.`ProductDesignId` AS CHAR), ':', CAST(s.`Id` AS CHAR))"),
  1348. new("S1_PRODUCT_DESIGN_ROUTING", "ado_product_design_routing", "mdp_stg_so", "Id", "CONCAT(CAST(s.`ProductDesignId` AS CHAR), ':', CAST(s.`Id` AS CHAR))"),
  1349. new("S1_REQUIREMENT_EXAMINE_RESULT", "b_examine_result", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`morder_no`,''), ':', CAST(s.`Id` AS CHAR))"),
  1350. new("S1_REQUIREMENT_EXAMINE_DETAIL", "b_bom_child_examine", "mdp_stg_so", "Id", "CONCAT(IFNULL(s.`examine_id`,''), ':', IFNULL(s.`item_number`,''), ':', CAST(s.`Id` AS CHAR))"),
  1351. new("S1_SHIPPING_PLAN", "ShippingPlan", "mdp_stg_ship_trans", "RecID", "COALESCE(s.`LotSerial`, CAST(s.`RecID` AS CHAR))"),
  1352. new("S1_SHIPPING_PLAN_DETAIL", "ShippingPlanDetail", "mdp_stg_ship_trans", "RecID", "CONCAT(IFNULL(s.`plan_id`,''), ':', IFNULL(s.`OrdNbr`,''), ':', CAST(s.`RecID` AS CHAR))"),
  1353. new("S1_ASN_SHIPPER_MASTER", "ASNBOLShipperMaster", "mdp_stg_ship_trans", "RecID", "COALESCE(s.`Id`, CONCAT(IFNULL(s.`OrdNbr`,''), ':', CAST(s.`RecID` AS CHAR)))"),
  1354. new("S1_ASN_SHIPPER_DETAIL", "ASNBOLShipperDetail", "mdp_stg_ship_trans", "RecID", "CONCAT(IFNULL(s.`Id`,''), ':', IFNULL(s.`Line`, CAST(s.`RecID` AS CHAR)))"),
  1355. new("S1_LINKAGE_PLAN", "LinkagePlan", "mdp_stg_ship_trans", "id", "CONCAT(IFNULL(s.`bill_no`,''), ':', IFNULL(s.`item_number`,''), ':', CAST(s.`id` AS CHAR))")
  1356. };
  1357. }