S1MdpSyncTransformService.cs 83 KB

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