MdpMonitorService.cs 36 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776
  1. namespace Admin.NET.Plugin.AiDOP.Order;
  2. /// <summary>
  3. /// 数据中台统一 MDP 运行监控。
  4. ///
  5. /// 【租户安全边界】需认证访问(已移除类级 <c>[AllowAnonymous]</c>);租户一律经
  6. /// <c>AidopTenantScope.ResolveOrThrow</c> 从认证后 JWT 解析,无有效租户即拒绝,
  7. /// 不读前端 tenantId、无默认回退(原实现:类级匿名 + 直接取 <c>_userManager.TenantId</c>,
  8. /// 匿名时退化为 <c>tenant_id = 0</c> 平台行)。租户行与平台行可见性语义(<c>BuildMdpRunLogTenantWhere</c>)不变。
  9. /// </summary>
  10. [ApiDescriptionSettings(Order = 322, Description = "统一MDP运行监控")]
  11. [Route("api/DataPlatform")]
  12. [NonUnify]
  13. public class MdpMonitorService : IDynamicApiController, ITransient
  14. {
  15. private static readonly Dictionary<string, string> ModuleJobCodes = new(StringComparer.OrdinalIgnoreCase)
  16. {
  17. ["S1"] = "S1_MDP_SYNC_TRANSFORM",
  18. ["S2"] = "S2_MDP_SYNC_TRANSFORM",
  19. ["S3"] = "S3_MDP_SYNC_TRANSFORM",
  20. ["S4"] = "S4_MDP_SYNC_TRANSFORM"
  21. };
  22. private static readonly Dictionary<string, MdpJobCatalogItem> JobCatalog = new(StringComparer.OrdinalIgnoreCase)
  23. {
  24. ["S1_MDP_SYNC_TRANSFORM"] = new(
  25. "S1_MDP_SYNC_TRANSFORM",
  26. "order_delivery",
  27. "订单交付域",
  28. ["S1", "S2", "S3", "S4", "S7", "S9"],
  29. "GLOBAL_DOMAIN",
  30. "订单交付域 MDP 同步",
  31. "job_s1_mdp_sync_transform"),
  32. ["S2_MDP_SYNC_TRANSFORM"] = new(
  33. "S2_MDP_SYNC_TRANSFORM",
  34. "work_schedule",
  35. "工单排程域",
  36. ["S2", "S3", "S5", "S6", "S8", "S9"],
  37. "GLOBAL_DOMAIN",
  38. "工单排程域 MDP 同步",
  39. "job_s2_mdp_sync_transform"),
  40. ["S3_MDP_SYNC_TRANSFORM"] = new(
  41. "S3_MDP_SYNC_TRANSFORM",
  42. "supply_purchase",
  43. "供应采购域",
  44. ["S3", "S4", "S5", "S8", "S9"],
  45. "GLOBAL_DOMAIN",
  46. "供应采购域 MDP 同步",
  47. "job_s3_mdp_sync_transform"),
  48. ["S4_MDP_SYNC_TRANSFORM"] = new(
  49. "S4_MDP_SYNC_TRANSFORM",
  50. "purchase_execution",
  51. "采购执行域",
  52. ["S4", "S5", "S8", "S9"],
  53. "GLOBAL_DOMAIN",
  54. "采购执行域 MDP 同步",
  55. "job_s4_mdp_sync_transform")
  56. };
  57. private readonly ISqlSugarClient _db;
  58. private readonly UserManager _userManager;
  59. public MdpMonitorService(ISqlSugarClient db, UserManager userManager)
  60. {
  61. _db = db;
  62. _userManager = userManager;
  63. }
  64. [DisplayName("MDP模块选项")]
  65. [HttpGet("mdp-monitor/modules")]
  66. public object GetModules() => BuildCatalogResponse();
  67. [DisplayName("MDP任务目录")]
  68. [HttpGet("mdp-monitor/catalog")]
  69. public object GetCatalog() => BuildCatalogResponse();
  70. [DisplayName("MDP最近运行状态")]
  71. [HttpGet("mdp-monitor/latest")]
  72. public async Task<object> GetLatest([FromQuery] MdpMonitorQueryInput input)
  73. {
  74. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  75. var (whereSql, pars) = BuildWhere(input, tenantId);
  76. var row = await _db.Ado.SqlQuerySingleAsync<MdpMonitorRunLogRow>(
  77. $"""
  78. {SelectColumnsSql()}
  79. FROM mdp_transform_run_log
  80. WHERE {whereSql}
  81. ORDER BY start_time DESC, id DESC
  82. LIMIT 1
  83. """,
  84. pars)
  85. ?? new MdpMonitorRunLogRow();
  86. AttachCatalogInfo(row);
  87. return row;
  88. }
  89. [DisplayName("MDP运行日志列表")]
  90. [HttpGet("mdp-monitor/list")]
  91. public async Task<object> GetList([FromQuery] MdpMonitorListInput input)
  92. {
  93. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  94. var page = input.Page <= 0 ? 1 : input.Page;
  95. var pageSize = input.PageSize <= 0 ? 10 : input.PageSize;
  96. var offset = (page - 1) * pageSize;
  97. var (whereSql, pars) = BuildWhere(input, tenantId);
  98. var total = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM mdp_transform_run_log WHERE {whereSql}", pars);
  99. var list = await _db.Ado.SqlQueryAsync<MdpMonitorRunLogRow>(
  100. $"""
  101. {SelectColumnsSql()}
  102. FROM mdp_transform_run_log
  103. WHERE {whereSql}
  104. ORDER BY start_time DESC, id DESC
  105. LIMIT {pageSize} OFFSET {offset}
  106. """,
  107. pars);
  108. foreach (var row in list)
  109. AttachCatalogInfo(row);
  110. return new { total, page, pageSize, list };
  111. }
  112. [DisplayName("MDP运行日志详情")]
  113. [HttpGet("mdp-monitor/detail/{id}")]
  114. public async Task<object> GetDetail(long id, [FromQuery] MdpMonitorQueryInput input)
  115. {
  116. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  117. var (whereSql, pars) = BuildWhere(input, tenantId);
  118. pars.Add(new SugarParameter("@Id", id));
  119. var row = await _db.Ado.SqlQuerySingleAsync<MdpMonitorRunLogRow>(
  120. $"""
  121. {SelectColumnsSql()}
  122. FROM mdp_transform_run_log
  123. WHERE id=@Id AND {whereSql}
  124. LIMIT 1
  125. """,
  126. pars);
  127. if (row == null)
  128. throw Oops.Oh("运行日志不存在");
  129. AttachCatalogInfo(row);
  130. return row;
  131. }
  132. [DisplayName("MDP同步链路详情")]
  133. [HttpGet("mdp-monitor/lineage")]
  134. public async Task<object> GetLineage([FromQuery] MdpMonitorLineageInput input)
  135. {
  136. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  137. var moduleCode = ResolveModuleCode(input.ModuleCode, input.JobCode);
  138. if (string.IsNullOrWhiteSpace(moduleCode))
  139. throw Oops.Oh("请选择 MDP 模块");
  140. var jobCode = ResolveJobCode(moduleCode, input.JobCode);
  141. var entityPrefix = $"{moduleCode}_%";
  142. var entities = await _db.Ado.SqlQueryAsync<MdpLineageEntityRow>(
  143. """
  144. SELECT e.id AS Id, e.entity_code AS EntityCode, e.entity_name AS EntityName,
  145. e.entity_type AS EntityType, s.source_code AS SourceCode,
  146. s.source_name AS SourceName, s.source_type AS SourceType,
  147. s.db_type AS SourceDbType, s.db_host AS SourceDbHost,
  148. s.db_port AS SourceDbPort, s.db_name AS SourceDbName,
  149. e.source_table_name AS SourceTableName, e.source_api_path AS SourceApiPath,
  150. s.db_type AS TargetDbType, s.db_host AS TargetDbHost,
  151. s.db_port AS TargetDbPort, s.db_name AS TargetDbName,
  152. e.target_table_name AS TargetTableName, e.sync_mode AS SyncMode,
  153. e.incr_column AS IncrColumn, e.status AS Status
  154. FROM mdp_entity e
  155. LEFT JOIN mdp_source s ON s.id = e.source_id
  156. WHERE e.entity_code LIKE @EntityPrefix
  157. AND (e.tenant_id = @TenantId OR e.tenant_id = 0)
  158. ORDER BY e.entity_code
  159. """,
  160. new SugarParameter("@EntityPrefix", entityPrefix),
  161. new SugarParameter("@TenantId", tenantId));
  162. if (entities.Count == 0)
  163. {
  164. var emptyOutput = new MdpLineageOutput
  165. {
  166. ModuleCode = moduleCode,
  167. JobCode = jobCode,
  168. BatchId = input.BatchId,
  169. Stages = BuildStageDescriptions(jobCode, moduleCode),
  170. Entities = new List<MdpLineageEntityRow>()
  171. };
  172. AttachLineageCatalogInfo(emptyOutput);
  173. return emptyOutput;
  174. }
  175. var entityIds = string.Join(",", entities.Select(u => u.Id));
  176. var mappings = await _db.Ado.SqlQueryAsync<MdpLineageFieldMappingRow>(
  177. $"""
  178. SELECT entity_id AS EntityId, source_field AS SourceField, target_field AS TargetField,
  179. field_type AS FieldType, transform_script AS TransformScript,
  180. const_value AS ConstValue, lookup_table AS LookupTable,
  181. is_required AS IsRequired, default_value AS DefaultValue, sort_order AS SortOrder
  182. FROM mdp_field_mapping
  183. WHERE entity_id IN ({entityIds})
  184. ORDER BY entity_id, sort_order, target_field
  185. """);
  186. var mappingsByEntity = mappings.GroupBy(u => u.EntityId).ToDictionary(u => u.Key, u => u.ToList());
  187. Dictionary<long, MdpLineageSyncLogRow> syncLogsByEntity = new();
  188. if (!string.IsNullOrWhiteSpace(input.BatchId))
  189. {
  190. var syncLogs = await _db.Ado.SqlQueryAsync<MdpLineageSyncLogRow>(
  191. """
  192. SELECT entity_id AS EntityId, entity_name AS EntityName, status AS Status,
  193. rows_read AS RowsRead, rows_insert AS RowsInsert, rows_update AS RowsUpdate,
  194. rows_skip AS RowsSkip, rows_error AS RowsError,
  195. sync_start AS SyncStart, sync_end AS SyncEnd, duration_ms AS DurationMs,
  196. error_msg AS ErrorMsg
  197. FROM mdp_sync_log
  198. WHERE sync_batch_id = @BatchId
  199. AND (tenant_id = @TenantId OR tenant_id = 0)
  200. AND entity_id IN (
  201. SELECT id FROM mdp_entity
  202. WHERE entity_code LIKE @EntityPrefix
  203. AND (tenant_id = @TenantId OR tenant_id = 0)
  204. )
  205. ORDER BY sync_start, id
  206. """,
  207. new SugarParameter("@BatchId", input.BatchId.Trim()),
  208. new SugarParameter("@EntityPrefix", entityPrefix),
  209. new SugarParameter("@TenantId", tenantId));
  210. syncLogsByEntity = syncLogs
  211. .GroupBy(u => u.EntityId)
  212. .ToDictionary(u => u.Key, u => u.OrderByDescending(x => x.SyncStart).First());
  213. }
  214. var batchId = input.BatchId?.Trim();
  215. foreach (var entity in entities)
  216. {
  217. entity.SourceFullName = BuildObjectFullName(entity.SourceDbType, entity.SourceDbHost, entity.SourceDbPort, entity.SourceDbName, entity.SourceTableName ?? entity.SourceApiPath);
  218. entity.TargetFullName = BuildObjectFullName(entity.TargetDbType, entity.TargetDbHost, entity.TargetDbPort, entity.TargetDbName, entity.TargetTableName);
  219. if (mappingsByEntity.TryGetValue(entity.Id, out var entityMappings) && entityMappings.Count > 0)
  220. {
  221. foreach (var mapping in entityMappings)
  222. {
  223. mapping.MappingSource = "CONFIG";
  224. mapping.IsFallback = false;
  225. }
  226. entity.FieldMappings = entityMappings;
  227. entity.FieldMappingCount = entityMappings.Count;
  228. }
  229. else
  230. {
  231. entity.FieldMappings = BuildFallbackFieldMappings(entity, batchId);
  232. entity.FieldMappingCount = entity.FieldMappings.Count;
  233. }
  234. if (syncLogsByEntity.TryGetValue(entity.Id, out var syncLog))
  235. entity.SyncLog = syncLog;
  236. }
  237. var output = new MdpLineageOutput
  238. {
  239. ModuleCode = moduleCode,
  240. JobCode = jobCode,
  241. BatchId = input.BatchId,
  242. Stages = BuildStageDescriptions(jobCode, moduleCode),
  243. Entities = entities
  244. };
  245. AttachLineageCatalogInfo(output);
  246. return output;
  247. }
  248. private static (string WhereSql, List<SugarParameter> Parameters) BuildWhere(MdpMonitorQueryInput input, long tenantId)
  249. {
  250. var where = new List<string> { BuildMdpRunLogTenantWhere(tenantId) };
  251. var pars = new List<SugarParameter> { new("@TenantId", tenantId) };
  252. var (jobFilterSql, jobFilterPars, noMatch) = BuildJobCodeFilter(input);
  253. if (noMatch)
  254. where.Add("1=0");
  255. else
  256. {
  257. where.Add(jobFilterSql);
  258. pars.AddRange(jobFilterPars);
  259. }
  260. if (!string.IsNullOrWhiteSpace(input.BatchId))
  261. {
  262. where.Add("batch_id LIKE @BatchId");
  263. pars.Add(new SugarParameter("@BatchId", $"%{input.BatchId.Trim()}%"));
  264. }
  265. if (!string.IsNullOrWhiteSpace(input.Status))
  266. {
  267. where.Add("status=@Status");
  268. pars.Add(new SugarParameter("@Status", input.Status.Trim().ToUpperInvariant()));
  269. }
  270. if (input.StartTime.HasValue)
  271. {
  272. where.Add("start_time >= @StartTime");
  273. pars.Add(new SugarParameter("@StartTime", input.StartTime.Value));
  274. }
  275. if (input.EndTime.HasValue)
  276. {
  277. where.Add("start_time <= @EndTime");
  278. pars.Add(new SugarParameter("@EndTime", input.EndTime.Value));
  279. }
  280. return (string.Join(" AND ", where), pars);
  281. }
  282. /// <summary>
  283. /// MDP 转换任务当前以 tenant_id=0 写入运行日志;登录租户查询时需兼容这类全局任务记录。
  284. /// </summary>
  285. internal static string BuildMdpRunLogTenantWhere(long tenantId) =>
  286. tenantId > 0 ? "(tenant_id = @TenantId OR tenant_id = 0)" : "tenant_id = 0";
  287. private static object BuildCatalogResponse()
  288. {
  289. return JobCatalog.Values
  290. .OrderBy(u => u.JobCode, StringComparer.OrdinalIgnoreCase)
  291. .Select(u => new
  292. {
  293. jobCode = u.JobCode,
  294. displayName = u.DisplayName,
  295. businessDomainCode = u.BusinessDomainCode,
  296. businessDomainName = u.BusinessDomainName,
  297. consumerModules = u.ConsumerModules,
  298. scopeType = u.ScopeType,
  299. moduleCode = ResolveModuleCodeFromJobCode(u.JobCode),
  300. scheduleJobId = u.ScheduleJobId
  301. })
  302. .ToList();
  303. }
  304. private static (string Sql, List<SugarParameter> Parameters, bool NoMatch) BuildJobCodeFilter(MdpMonitorQueryInput input)
  305. {
  306. var pars = new List<SugarParameter>();
  307. if (!string.IsNullOrWhiteSpace(input.JobCode))
  308. {
  309. pars.Add(new SugarParameter("@JobCode", input.JobCode.Trim().ToUpperInvariant()));
  310. return ("job_code=@JobCode", pars, false);
  311. }
  312. var hasDomainFilter = !string.IsNullOrWhiteSpace(input.BusinessDomainCode);
  313. var consumerModule = !string.IsNullOrWhiteSpace(input.ConsumerModule)
  314. ? input.ConsumerModule.Trim()
  315. : !string.IsNullOrWhiteSpace(input.ModuleCode) ? input.ModuleCode.Trim() : null;
  316. var hasConsumerFilter = !string.IsNullOrWhiteSpace(consumerModule);
  317. if (!hasDomainFilter && !hasConsumerFilter)
  318. return ("IFNULL(job_code, '') LIKE '%MDP%'", pars, false);
  319. var allowed = ResolveAllowedJobCodes(input.BusinessDomainCode, consumerModule);
  320. if (allowed.Count == 0)
  321. return (string.Empty, pars, true);
  322. if (allowed.Count == 1)
  323. {
  324. pars.Add(new SugarParameter("@JobCode", allowed.First()));
  325. return ("job_code=@JobCode", pars, false);
  326. }
  327. var inParts = new List<string>();
  328. var index = 0;
  329. foreach (var code in allowed.OrderBy(u => u, StringComparer.OrdinalIgnoreCase))
  330. {
  331. var paramName = $"@JobCode{index++}";
  332. inParts.Add(paramName);
  333. pars.Add(new SugarParameter(paramName, code));
  334. }
  335. return ($"job_code IN ({string.Join(", ", inParts)})", pars, false);
  336. }
  337. private static HashSet<string> ResolveAllowedJobCodes(string? businessDomainCode, string? consumerModule)
  338. {
  339. IEnumerable<MdpJobCatalogItem> items = JobCatalog.Values;
  340. if (!string.IsNullOrWhiteSpace(businessDomainCode))
  341. {
  342. var domain = businessDomainCode.Trim();
  343. items = items.Where(u => string.Equals(u.BusinessDomainCode, domain, StringComparison.OrdinalIgnoreCase));
  344. }
  345. if (!string.IsNullOrWhiteSpace(consumerModule))
  346. {
  347. var module = consumerModule.Trim();
  348. items = items.Where(u => u.ConsumerModules.Contains(module, StringComparer.OrdinalIgnoreCase));
  349. }
  350. return items.Select(u => u.JobCode).ToHashSet(StringComparer.OrdinalIgnoreCase);
  351. }
  352. private static void AttachCatalogInfo(MdpMonitorRunLogRow row)
  353. {
  354. if (row == null || string.IsNullOrWhiteSpace(row.JobCode))
  355. return;
  356. if (!JobCatalog.TryGetValue(row.JobCode.Trim(), out var item))
  357. return;
  358. row.BusinessDomainCode = item.BusinessDomainCode;
  359. row.BusinessDomainName = item.BusinessDomainName;
  360. row.ConsumerModules = string.Join(",", item.ConsumerModules);
  361. row.ScopeType = item.ScopeType;
  362. row.DisplayName = item.DisplayName;
  363. row.ScheduleJobId = item.ScheduleJobId;
  364. }
  365. private static void AttachLineageCatalogInfo(MdpLineageOutput output)
  366. {
  367. if (output == null || string.IsNullOrWhiteSpace(output.JobCode))
  368. return;
  369. if (!JobCatalog.TryGetValue(output.JobCode.Trim(), out var item))
  370. return;
  371. output.BusinessDomainCode = item.BusinessDomainCode;
  372. output.BusinessDomainName = item.BusinessDomainName;
  373. output.ConsumerModules = string.Join(",", item.ConsumerModules);
  374. output.ScopeType = item.ScopeType;
  375. output.DisplayName = item.DisplayName;
  376. output.ScheduleJobId = item.ScheduleJobId;
  377. }
  378. private static string FormatConsumerModulesLabel(MdpJobCatalogItem item) =>
  379. string.Join("、", item.ConsumerModules);
  380. private static List<MdpLineageStageRow> BuildStageDescriptions(string? jobCode, string? moduleCode)
  381. {
  382. if (!string.IsNullOrWhiteSpace(jobCode) && JobCatalog.TryGetValue(jobCode.Trim(), out var catalogItem))
  383. {
  384. return catalogItem.BusinessDomainCode switch
  385. {
  386. "order_delivery" => BuildOrderDeliveryStages(catalogItem),
  387. "work_schedule" => BuildWorkScheduleStages(catalogItem),
  388. "supply_purchase" => BuildSupplyPurchaseStages(catalogItem),
  389. "purchase_execution" => BuildPurchaseExecutionStages(catalogItem),
  390. _ => BuildGenericStages(moduleCode)
  391. };
  392. }
  393. return BuildGenericStages(moduleCode);
  394. }
  395. private static List<MdpLineageStageRow> BuildOrderDeliveryStages(MdpJobCatalogItem item)
  396. {
  397. var consumers = FormatConsumerModulesLabel(item);
  398. return new List<MdpLineageStageRow>
  399. {
  400. new() { StageCode = "STAGING", StageName = "订单交付域 · 贴源同步", Layer = "mdp_stg", Description = $"按 mdp_entity 登记订单交付域源对象抽取数据,保留 raw_data JSON 便于追溯;产出供 {consumers} 消费。", InputObjects = "旧系统 / 当前库源对象", OutputObjects = "mdp_stg_so, mdp_stg_ship_trans", Execution = "S1MdpSyncTransformService.SyncStagingAsync" },
  401. new() { StageCode = "STANDARD", StageName = "订单交付域 · 标准层转换", Layer = "mdp_std", Description = "解析贴源 raw_data,做字段标准化、租户兜底和幂等写入。", InputObjects = "mdp_stg_so, mdp_stg_ship_trans", OutputObjects = "mdp_std_so, mdp_std_ship_trans", Execution = "S1MdpSyncTransformService.BuildStandardCommands" },
  402. new() { StageCode = "DWD", StageName = "订单交付域 · DWD宽表", Layer = "dwd", Description = $"沉淀订单交付事实,供 {consumers} 看板与诊断读取。", InputObjects = "mdp_std_so, mdp_std_ship_trans", OutputObjects = "dwd_ship_trans", Execution = "S1MdpSyncTransformService.BuildDwdAsync" },
  403. new() { StageCode = "KPI", StageName = "指标写入", Layer = "ado_s9", Description = $"计算订单交付域 L1 指标并写入统一指标值表,供 {consumers} 消费。", InputObjects = "mdp_std_so, dwd_ship_trans", OutputObjects = "ado_s9_kpi_value_l1_day", Execution = "S1MdpSyncTransformService.BuildS1KpiValuesAsync" }
  404. };
  405. }
  406. private static List<MdpLineageStageRow> BuildWorkScheduleStages(MdpJobCatalogItem item)
  407. {
  408. var consumers = FormatConsumerModulesLabel(item);
  409. return new List<MdpLineageStageRow>
  410. {
  411. new() { StageCode = "STAGING", StageName = "工单排程域 · 贴源同步", Layer = "mdp_stg", Description = $"按 mdp_entity 登记工单排程域源对象抽取数据;产出供 {consumers} 消费。", InputObjects = "工单 / 工序 / 排程源对象", OutputObjects = "mdp_stg_*", Execution = "S2MdpSyncTransformService.SyncStagingAsync" },
  412. new() { StageCode = "STANDARD", StageName = "工单排程域 · 标准层转换", Layer = "mdp_std", Description = "将工单、工序、排程等对象标准化。", InputObjects = "mdp_stg_*", OutputObjects = "mdp_std_*", Execution = "S2MdpSyncTransformService.BuildStandardCommands" },
  413. new() { StageCode = "DWD", StageName = "工单排程域 · DWD宽表", Layer = "dwd", Description = $"生成制造执行与排程分析宽表,供 {consumers} 读取。", InputObjects = "mdp_std_*", OutputObjects = "dwd_*", Execution = "S2MdpSyncTransformService.BuildDwdAsync" },
  414. new() { StageCode = "KPI", StageName = "指标写入", Layer = "ado_s9", Description = $"写入工单排程域 KPI,供 {consumers} 消费。", InputObjects = "mdp_std_* / dwd_*", OutputObjects = "ado_s9_kpi_value_*", Execution = "S2MdpSyncTransformService.BuildS2KpiValuesAsync" }
  415. };
  416. }
  417. private static List<MdpLineageStageRow> BuildSupplyPurchaseStages(MdpJobCatalogItem item)
  418. {
  419. var consumers = FormatConsumerModulesLabel(item);
  420. return new List<MdpLineageStageRow>
  421. {
  422. new() { StageCode = "STAGING", StageName = "供应采购域 · 贴源同步", Layer = "mdp_stg", Description = $"按 mdp_entity 登记供应采购域源对象抽取数据;产出供 {consumers} 消费。", InputObjects = "供应 / 物料 / 采购源对象", OutputObjects = "mdp_stg_*", Execution = "S3MdpSyncTransformService.SyncStagingAsync" },
  423. new() { StageCode = "STANDARD", StageName = "供应采购域 · 标准层转换", Layer = "mdp_std", Description = "将供应、物料、采购、交货计划等对象标准化。", InputObjects = "mdp_stg_*", OutputObjects = "mdp_std_*", Execution = "S3MdpSyncTransformService.BuildStandardCommands" },
  424. new() { StageCode = "DWD", StageName = "供应采购域 · DWD宽表", Layer = "dwd", Description = $"生成供应交付、齐套、风险等分析宽表,供 {consumers} 读取。", InputObjects = "mdp_std_*", OutputObjects = "dwd_supplier_delivery / dwd_material_readiness 等", Execution = "S3MdpSyncTransformService.BuildDwdAsync" },
  425. new() { StageCode = "KPI", StageName = "指标写入", Layer = "ado_s9", Description = $"写入供应采购域指标,供 {consumers} 消费。", InputObjects = "mdp_std_* / dwd_*", OutputObjects = "ado_s9_kpi_value_*", Execution = "S3MdpSyncTransformService.BuildS3KpiValuesAsync" }
  426. };
  427. }
  428. private static List<MdpLineageStageRow> BuildPurchaseExecutionStages(MdpJobCatalogItem item)
  429. {
  430. var consumers = FormatConsumerModulesLabel(item);
  431. return new List<MdpLineageStageRow>
  432. {
  433. new() { StageCode = "STAGING", StageName = "采购执行域 · 贴源同步", Layer = "mdp_stg_s4_*", Description = $"同步采购执行域 IQC/发货/退货/欠料事实,供 {consumers} 消费;共享采购主链仍由供应采购域维护。", InputObjects = "PurOrdRctDetail / scm_shdzb / srm_polist_ds / dwd_material_shortage", OutputObjects = "mdp_stg_s4_iqc / mdp_stg_s4_shipment / mdp_stg_s4_return / mdp_stg_s4_shortage", Execution = "S4MdpSyncTransformService.SyncStagingAsync" },
  434. new() { StageCode = "STANDARD", StageName = "采购执行域 · 标准层转换", Layer = "mdp_std_s4_*", Description = "将采购执行域贴源对象标准化,并统计供应采购域共享标准层行数。", InputObjects = "mdp_stg_s4_*", OutputObjects = "mdp_std_s4_iqc / mdp_std_s4_shipment / mdp_std_s4_return / mdp_std_s4_shortage", Execution = "S4MdpSyncTransformService.BuildStandardCommands" },
  435. new() { StageCode = "DWD", StageName = "采购执行域 · DWD宽表", Layer = "dwd", Description = $"写入采购执行分析宽表,供 {consumers} 读取。", InputObjects = "dwd_supplier_delivery / mdp_std_s4_*", OutputObjects = "dwd_s4_purchase_execution / dwd_po_trans / dwd_qc_trans", Execution = "S4MdpSyncTransformService.BuildDwdAsync" },
  436. new() { StageCode = "KPI", StageName = "指标写入", Layer = "ado_s9", Description = $"写入采购执行域 L1/L2/L3 指标,供 {consumers} 消费。", InputObjects = "mdp_std_delivery_schedule / dwd_supplier_delivery / dwd_s4_purchase_execution", OutputObjects = "ado_s9_kpi_value_l1/l2/l3_day", Execution = "S4MdpSyncTransformService.BuildS4KpiValuesAsync" }
  437. };
  438. }
  439. /// <summary>
  440. /// 贴源同步在代码中维护、但未写入 mdp_field_mapping 的实体主键/业务键提示(仅监控展示兜底)。
  441. /// </summary>
  442. private static readonly Dictionary<string, (string SourceRowId, string SourceBizKeyExpr)> KnownStagingEntityKeys =
  443. new(StringComparer.OrdinalIgnoreCase)
  444. {
  445. ["S4_IQC_RECEIPT"] = ("RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`Receiver`,''), ':', IFNULL(s.`Line`,''))"),
  446. ["S4_SHIPMENT_EXEC"] = ("id", "CONCAT(IFNULL(s.`glid`,''), ':', IFNULL(s.`id`,''))"),
  447. ["S4_RETURN_EXEC"] = ("Id", "s.`dsnum`"),
  448. ["S4_SHORTAGE_EXEC"] = ("id", "CONCAT(IFNULL(s.`work_order`,''), ':', IFNULL(s.`component_item_code`,''))")
  449. };
  450. private static List<MdpLineageFieldMappingRow> BuildFallbackFieldMappings(MdpLineageEntityRow entity, string? batchId)
  451. {
  452. var sourceRowId = entity.IncrColumn;
  453. string? sourceBizKeyExpr = null;
  454. if (!string.IsNullOrWhiteSpace(entity.EntityCode) &&
  455. KnownStagingEntityKeys.TryGetValue(entity.EntityCode, out var known))
  456. {
  457. sourceRowId ??= known.SourceRowId;
  458. sourceBizKeyExpr = known.SourceBizKeyExpr;
  459. }
  460. var fallbackNote = "当前实体未配置逐字段映射,贴源同步保留源行 raw_data";
  461. var mappings = new List<MdpLineageFieldMappingRow>();
  462. var sort = 10;
  463. void Add(string sourceField, string targetField, string fieldType, string? transformScript, string? constValue = null, bool isRequired = false)
  464. {
  465. mappings.Add(new MdpLineageFieldMappingRow
  466. {
  467. EntityId = entity.Id,
  468. SourceField = sourceField,
  469. TargetField = targetField,
  470. FieldType = fieldType,
  471. TransformScript = transformScript,
  472. ConstValue = constValue,
  473. IsRequired = isRequired,
  474. SortOrder = sort,
  475. MappingSource = "FALLBACK",
  476. IsFallback = true
  477. });
  478. sort += 10;
  479. }
  480. Add("tenant_id", "tenant_id", "DIRECT", "当前租户或全局租户兜底");
  481. Add($"CONST:{entity.SourceCode ?? "AIDOP"}", "source_system", "CONST", "来源数据源编码", entity.SourceCode ?? "AIDOP", isRequired: true);
  482. Add($"CONST:{entity.SourceTableName ?? entity.SourceApiPath ?? "--"}", "source_table", "CONST", "源表名", entity.SourceTableName ?? entity.SourceApiPath, isRequired: true);
  483. if (!string.IsNullOrWhiteSpace(sourceRowId))
  484. Add(sourceRowId, "source_row_id", "DIRECT", "来自 mdp_entity.incr_column 或实体主键配置");
  485. else
  486. Add("--", "source_row_id", "DIRECT", fallbackNote);
  487. if (!string.IsNullOrWhiteSpace(sourceBizKeyExpr))
  488. Add(sourceBizKeyExpr, "source_biz_key", "EXPR", "来自实体业务键表达式", isRequired: true);
  489. else
  490. Add("--", "source_biz_key", "EXPR", fallbackNote);
  491. Add("*", "raw_data", "JSON", "源行整行 JSON");
  492. Add("CONST", "sync_batch_id", "CONST", "当前同步批次", string.IsNullOrWhiteSpace(batchId) ? "当前同步批次" : batchId, isRequired: true);
  493. Add("NOW()", "sync_time", "CONST", "同步时间", isRequired: true);
  494. return mappings;
  495. }
  496. private static List<MdpLineageStageRow> BuildGenericStages(string? moduleCode)
  497. {
  498. if (string.Equals(moduleCode, "S1", StringComparison.OrdinalIgnoreCase))
  499. return BuildOrderDeliveryStages(JobCatalog["S1_MDP_SYNC_TRANSFORM"]);
  500. if (string.Equals(moduleCode, "S2", StringComparison.OrdinalIgnoreCase))
  501. return BuildWorkScheduleStages(JobCatalog["S2_MDP_SYNC_TRANSFORM"]);
  502. if (string.Equals(moduleCode, "S3", StringComparison.OrdinalIgnoreCase))
  503. return BuildSupplyPurchaseStages(JobCatalog["S3_MDP_SYNC_TRANSFORM"]);
  504. if (string.Equals(moduleCode, "S4", StringComparison.OrdinalIgnoreCase))
  505. return BuildPurchaseExecutionStages(JobCatalog["S4_MDP_SYNC_TRANSFORM"]);
  506. return new List<MdpLineageStageRow>
  507. {
  508. new() { StageCode = "STAGING", StageName = "贴源同步", Layer = "mdp_stg", Description = "按 mdp_entity 登记源对象抽取数据。", InputObjects = "源对象", OutputObjects = "mdp_stg_*", Execution = "MDP 同步服务" },
  509. new() { StageCode = "STANDARD", StageName = "标准层转换", Layer = "mdp_std", Description = "标准层/DWD/KPI 当前由后端 Service 承载。", InputObjects = "mdp_stg_*", OutputObjects = "mdp_std_* / dwd_* / 指标表", Execution = "MDP 转换服务" }
  510. };
  511. }
  512. private static string? ResolveModuleCodeFromJobCode(string? jobCode)
  513. {
  514. if (string.IsNullOrWhiteSpace(jobCode))
  515. return null;
  516. return ModuleJobCodes.FirstOrDefault(u =>
  517. string.Equals(u.Value, jobCode.Trim(), StringComparison.OrdinalIgnoreCase)).Key;
  518. }
  519. private static string? ResolveJobCode(string? moduleCode, string? jobCode)
  520. {
  521. if (!string.IsNullOrWhiteSpace(jobCode))
  522. return jobCode.Trim().ToUpperInvariant();
  523. if (string.IsNullOrWhiteSpace(moduleCode))
  524. return null;
  525. return ModuleJobCodes.TryGetValue(moduleCode.Trim(), out var mapped) ? mapped : null;
  526. }
  527. private static string? ResolveModuleCode(string? moduleCode, string? jobCode)
  528. {
  529. if (!string.IsNullOrWhiteSpace(moduleCode))
  530. return moduleCode.Trim().ToUpperInvariant();
  531. if (string.IsNullOrWhiteSpace(jobCode))
  532. return null;
  533. var normalizedJobCode = jobCode.Trim();
  534. return ModuleJobCodes.FirstOrDefault(u => string.Equals(u.Value, normalizedJobCode, StringComparison.OrdinalIgnoreCase)).Key;
  535. }
  536. private static string? BuildObjectFullName(string? dbType, string? host, int? port, string? dbName, string? objectName)
  537. {
  538. if (string.IsNullOrWhiteSpace(objectName))
  539. return null;
  540. var databaseObject = string.IsNullOrWhiteSpace(dbName) ? objectName : $"{dbName}.{objectName}";
  541. var hostPart = string.IsNullOrWhiteSpace(host) ? null : port.HasValue ? $"{host}:{port}" : host;
  542. return string.Join(" / ", new[] { dbType, hostPart, databaseObject }.Where(u => !string.IsNullOrWhiteSpace(u)));
  543. }
  544. private static string SelectColumnsSql()
  545. {
  546. return """
  547. SELECT id AS Id, tenant_id AS TenantId, job_code AS JobCode, job_name AS JobName, trigger_type AS TriggerType,
  548. batch_id AS BatchId, status AS Status, start_time AS StartTime, end_time AS EndTime, duration_ms AS DurationMs,
  549. stage_rows AS StageRows, standard_rows AS StandardRows, dwd_rows AS DwdRows,
  550. error_message AS ErrorMessage, summary_json AS SummaryJson, create_time AS CreateTime, update_time AS UpdateTime
  551. """;
  552. }
  553. }
  554. public sealed record MdpJobCatalogItem(
  555. string JobCode,
  556. string BusinessDomainCode,
  557. string BusinessDomainName,
  558. string[] ConsumerModules,
  559. string ScopeType,
  560. string DisplayName,
  561. string? ScheduleJobId);
  562. public class MdpMonitorQueryInput
  563. {
  564. public string? BusinessDomainCode { get; set; }
  565. public string? ConsumerModule { get; set; }
  566. public string? ModuleCode { get; set; }
  567. public string? JobCode { get; set; }
  568. public string? BatchId { get; set; }
  569. public string? Status { get; set; }
  570. public DateTime? StartTime { get; set; }
  571. public DateTime? EndTime { get; set; }
  572. }
  573. public sealed class MdpMonitorListInput : MdpMonitorQueryInput
  574. {
  575. public int Page { get; set; } = 1;
  576. public int PageSize { get; set; } = 10;
  577. }
  578. public sealed class MdpMonitorLineageInput : MdpMonitorQueryInput
  579. {
  580. }
  581. public sealed class MdpMonitorRunLogRow
  582. {
  583. public long Id { get; set; }
  584. public long TenantId { get; set; }
  585. public string? JobCode { get; set; }
  586. public string? JobName { get; set; }
  587. public string? TriggerType { get; set; }
  588. public string? BatchId { get; set; }
  589. public string? Status { get; set; }
  590. public DateTime? StartTime { get; set; }
  591. public DateTime? EndTime { get; set; }
  592. public int? DurationMs { get; set; }
  593. public int? StageRows { get; set; }
  594. public int? StandardRows { get; set; }
  595. public int? DwdRows { get; set; }
  596. public string? ErrorMessage { get; set; }
  597. public string? SummaryJson { get; set; }
  598. public DateTime? CreateTime { get; set; }
  599. public DateTime? UpdateTime { get; set; }
  600. public string? BusinessDomainCode { get; set; }
  601. public string? BusinessDomainName { get; set; }
  602. public string? ConsumerModules { get; set; }
  603. public string? ScopeType { get; set; }
  604. public string? DisplayName { get; set; }
  605. public string? ScheduleJobId { get; set; }
  606. }
  607. public sealed class MdpLineageOutput
  608. {
  609. public string? ModuleCode { get; set; }
  610. public string? JobCode { get; set; }
  611. public string? BatchId { get; set; }
  612. public string? BusinessDomainCode { get; set; }
  613. public string? BusinessDomainName { get; set; }
  614. public string? ConsumerModules { get; set; }
  615. public string? ScopeType { get; set; }
  616. public string? DisplayName { get; set; }
  617. public string? ScheduleJobId { get; set; }
  618. public List<MdpLineageStageRow> Stages { get; set; } = new();
  619. public List<MdpLineageEntityRow> Entities { get; set; } = new();
  620. }
  621. public sealed class MdpLineageStageRow
  622. {
  623. public string? StageCode { get; set; }
  624. public string? StageName { get; set; }
  625. public string? Layer { get; set; }
  626. public string? Description { get; set; }
  627. public string? InputObjects { get; set; }
  628. public string? OutputObjects { get; set; }
  629. public string? Execution { get; set; }
  630. }
  631. public sealed class MdpLineageEntityRow
  632. {
  633. public long Id { get; set; }
  634. public string? EntityCode { get; set; }
  635. public string? EntityName { get; set; }
  636. public string? EntityType { get; set; }
  637. public string? SourceCode { get; set; }
  638. public string? SourceName { get; set; }
  639. public string? SourceType { get; set; }
  640. public string? SourceDbType { get; set; }
  641. public string? SourceDbHost { get; set; }
  642. public int? SourceDbPort { get; set; }
  643. public string? SourceDbName { get; set; }
  644. public string? SourceTableName { get; set; }
  645. public string? SourceApiPath { get; set; }
  646. public string? SourceFullName { get; set; }
  647. public string? TargetDbType { get; set; }
  648. public string? TargetDbHost { get; set; }
  649. public int? TargetDbPort { get; set; }
  650. public string? TargetDbName { get; set; }
  651. public string? TargetTableName { get; set; }
  652. public string? TargetFullName { get; set; }
  653. public string? SyncMode { get; set; }
  654. public string? IncrColumn { get; set; }
  655. public int? Status { get; set; }
  656. public int FieldMappingCount { get; set; }
  657. public List<MdpLineageFieldMappingRow> FieldMappings { get; set; } = new();
  658. public MdpLineageSyncLogRow? SyncLog { get; set; }
  659. }
  660. public sealed class MdpLineageFieldMappingRow
  661. {
  662. public long EntityId { get; set; }
  663. public string? SourceField { get; set; }
  664. public string? TargetField { get; set; }
  665. public string? FieldType { get; set; }
  666. public string? TransformScript { get; set; }
  667. public string? ConstValue { get; set; }
  668. public string? LookupTable { get; set; }
  669. public bool IsRequired { get; set; }
  670. public string? DefaultValue { get; set; }
  671. public int SortOrder { get; set; }
  672. public string? MappingSource { get; set; }
  673. public bool IsFallback { get; set; }
  674. }
  675. public sealed class MdpLineageSyncLogRow
  676. {
  677. public long EntityId { get; set; }
  678. public string? EntityName { get; set; }
  679. public string? Status { get; set; }
  680. public long? RowsRead { get; set; }
  681. public long? RowsInsert { get; set; }
  682. public long? RowsUpdate { get; set; }
  683. public long? RowsSkip { get; set; }
  684. public long? RowsError { get; set; }
  685. public DateTime? SyncStart { get; set; }
  686. public DateTime? SyncEnd { get; set; }
  687. public int? DurationMs { get; set; }
  688. public string? ErrorMsg { get; set; }
  689. }