MdpMonitorService.cs 46 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018
  1. using Admin.NET.Core.Service;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform;
  3. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  4. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  5. using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
  6. namespace Admin.NET.Plugin.AiDOP.Order;
  7. /// <summary>
  8. /// 数据中台统一 MDP 运行监控。
  9. ///
  10. /// 【租户安全边界】需认证访问(已移除类级 <c>[AllowAnonymous]</c>);租户一律经
  11. /// <c>AidopTenantScope.ResolveOrThrow</c> 从认证后 JWT 解析,无有效租户即拒绝,
  12. /// 不读前端 tenantId、无默认回退(原实现:类级匿名 + 直接取 <c>_userManager.TenantId</c>,
  13. /// 匿名时退化为 <c>tenant_id = 0</c> 平台行)。租户行与平台行可见性语义(<c>BuildMdpRunLogTenantWhere</c>)不变。
  14. /// </summary>
  15. [ApiDescriptionSettings(Order = 322, Description = "统一MDP运行监控")]
  16. [Route("api/DataPlatform")]
  17. [NonUnify]
  18. public class MdpMonitorService : IDynamicApiController, ITransient
  19. {
  20. private static readonly Dictionary<string, string> ModuleJobCodes = new(StringComparer.OrdinalIgnoreCase)
  21. {
  22. ["S1"] = "S1_MDP_SYNC_TRANSFORM",
  23. ["S2"] = "S2_MDP_SYNC_TRANSFORM",
  24. ["S3"] = "S3_MDP_SYNC_TRANSFORM",
  25. ["S4"] = "S4_MDP_SYNC_TRANSFORM"
  26. };
  27. private static readonly Dictionary<string, MdpJobCatalogItem> JobCatalog = new(StringComparer.OrdinalIgnoreCase)
  28. {
  29. ["S1_MDP_SYNC_TRANSFORM"] = new(
  30. "S1_MDP_SYNC_TRANSFORM",
  31. "order_delivery",
  32. "订单交付域",
  33. ["S1", "S2", "S3", "S4", "S7", "S9"],
  34. "GLOBAL_DOMAIN",
  35. "订单交付域 MDP 同步",
  36. "job_s1_mdp_sync_transform"),
  37. ["S2_MDP_SYNC_TRANSFORM"] = new(
  38. "S2_MDP_SYNC_TRANSFORM",
  39. "work_schedule",
  40. "工单排程域",
  41. ["S2", "S3", "S5", "S6", "S8", "S9"],
  42. "GLOBAL_DOMAIN",
  43. "工单排程域 MDP 同步",
  44. "job_s2_mdp_sync_transform"),
  45. ["S3_MDP_SYNC_TRANSFORM"] = new(
  46. "S3_MDP_SYNC_TRANSFORM",
  47. "supply_purchase",
  48. "供应采购域",
  49. ["S3", "S4", "S5", "S8", "S9"],
  50. "GLOBAL_DOMAIN",
  51. "供应采购域 MDP 同步",
  52. "job_s3_mdp_sync_transform"),
  53. ["S4_MDP_SYNC_TRANSFORM"] = new(
  54. "S4_MDP_SYNC_TRANSFORM",
  55. "purchase_execution",
  56. "采购执行域",
  57. ["S4", "S5", "S8", "S9"],
  58. "GLOBAL_DOMAIN",
  59. "采购执行域 MDP 同步",
  60. "job_s4_mdp_sync_transform")
  61. };
  62. private readonly ISqlSugarClient _db;
  63. private readonly UserManager _userManager;
  64. private readonly MdpOutboundGate _outboundGate;
  65. private readonly DataPlatform.MdpRebuild.IModuleRebuildCapability _rebuildCapability;
  66. private readonly IEtlInstanceStore _etlInstances;
  67. private readonly DataPlatform.MdpRebuild.ModuleRebuildService _rebuild;
  68. private readonly DataPlatform.MdpRebuild.IModuleRebuildJobStore _rebuildJobs;
  69. public MdpMonitorService(
  70. ISqlSugarClient db,
  71. UserManager userManager,
  72. MdpOutboundGate outboundGate,
  73. DataPlatform.MdpRebuild.IModuleRebuildCapability rebuildCapability,
  74. IEtlInstanceStore etlInstances,
  75. DataPlatform.MdpRebuild.ModuleRebuildService rebuild,
  76. DataPlatform.MdpRebuild.IModuleRebuildJobStore rebuildJobs)
  77. {
  78. _db = db;
  79. _userManager = userManager;
  80. _outboundGate = outboundGate;
  81. _rebuildCapability = rebuildCapability;
  82. _etlInstances = etlInstances;
  83. _rebuild = rebuild;
  84. _rebuildJobs = rebuildJobs;
  85. }
  86. /// <summary>超管接口的统一门。前端的 v-if 只是不渲染,鉴权只认这里。</summary>
  87. private void RequireSuperAdmin()
  88. {
  89. if (!_userManager.SuperAdmin) throw Oops.Oh(ErrorCodeEnum.SA001);
  90. }
  91. [DisplayName("超管:租户列表")]
  92. [HttpGet("mdp-monitor/admin/tenants")]
  93. public async Task<object> GetTenantsAsync()
  94. {
  95. RequireSuperAdmin();
  96. return await _db.Queryable<SysTenant>()
  97. .LeftJoin<SysOrg>((u, a) => u.OrgId == a.Id).ClearFilter()
  98. .Where(u => u.Status == StatusEnum.Enable)
  99. .Select((u, a) => new
  100. {
  101. tenantId = u.Id,
  102. label = SqlFunc.HasValue(u.Title) ? $"{u.Title}-{a.Name}" : a.Name
  103. })
  104. .ToListAsync();
  105. }
  106. [DisplayName("超管:跨租户运行日志")]
  107. [HttpGet("mdp-monitor/admin/list")]
  108. public async Task<object> GetAdminListAsync([FromQuery] long? tenantId, [FromQuery] int page = 1, [FromQuery] int pageSize = 20)
  109. {
  110. RequireSuperAdmin();
  111. var size = Math.Clamp(pageSize, 1, 200);
  112. var offset = Math.Max(0, page - 1) * size;
  113. var where = tenantId is > 0 ? BuildMdpRunLogTenantWhere(tenantId.Value) : "1=1";
  114. var pars = new List<SugarParameter>();
  115. if (tenantId is > 0) pars.Add(new SugarParameter("@TenantId", tenantId.Value));
  116. var total = await _db.Ado.GetIntAsync($"SELECT COUNT(*) FROM mdp_transform_run_log WHERE {where}", pars);
  117. pars.Add(new SugarParameter("@Limit", size));
  118. pars.Add(new SugarParameter("@Offset", offset));
  119. var list = await _db.Ado.SqlQueryAsync<dynamic>(
  120. $"SELECT * FROM mdp_transform_run_log WHERE {where} ORDER BY id DESC LIMIT @Limit OFFSET @Offset", pars);
  121. // 重算任务与 run log 是两张表,取消只能作用于前者。单独列出仍可取消的。
  122. var jobsQuery = _db.Queryable<AdoModuleDashboardRebuildJob>()
  123. .Where(x => x.Status == "QUEUED" || x.Status == "RUNNING");
  124. if (tenantId is > 0)
  125. jobsQuery = jobsQuery.Where(x => x.TenantId == tenantId.Value);
  126. var activeJobs = await jobsQuery
  127. .OrderBy(x => x.Id, OrderByType.Desc)
  128. .Select(x => new { x.Id, x.ModuleCode, x.TenantId, x.Status, x.TriggerType, x.CurrentStage, x.CancelRequestedFlag, x.SubmittedAt })
  129. .Take(50)
  130. .ToListAsync();
  131. return new { total, page, pageSize = size, list, activeJobs };
  132. }
  133. public sealed class AssignRunnerInput
  134. {
  135. public string InstanceId { get; set; } = string.Empty;
  136. }
  137. [DisplayName("超管:指派执行机")]
  138. [HttpPost("mdp-monitor/admin/assign-runner")]
  139. public async Task<object> AssignRunnerAsync([FromBody] AssignRunnerInput input)
  140. {
  141. RequireSuperAdmin();
  142. if (string.IsNullOrWhiteSpace(input?.InstanceId))
  143. throw Oops.Oh("必须指定实例");
  144. var ok = await _etlInstances.AssignRunnerAsync(input.InstanceId.Trim());
  145. if (!ok) throw Oops.Oh("该实例不存在或已退出");
  146. return new { ok = true };
  147. }
  148. [DisplayName("超管:取消执行机指派")]
  149. [HttpPost("mdp-monitor/admin/clear-runner")]
  150. public async Task<object> ClearRunnerAsync()
  151. {
  152. RequireSuperAdmin();
  153. await _etlInstances.ClearRunnerAsync();
  154. return new { ok = true };
  155. }
  156. public sealed class AdminRebuildInput
  157. {
  158. public string ModuleCode { get; set; } = string.Empty;
  159. public long TenantId { get; set; }
  160. public long FactoryId { get; set; } = 1;
  161. }
  162. [DisplayName("超管:指定租户重算")]
  163. [HttpPost("mdp-monitor/admin/rebuild")]
  164. public async Task<object> RebuildAsync([FromBody] AdminRebuildInput input)
  165. {
  166. RequireSuperAdmin();
  167. if (input == null || input.TenantId <= 0 || string.IsNullOrWhiteSpace(input.ModuleCode))
  168. throw Oops.Oh("必须指定模块与租户");
  169. // requestedBy 传超管自己:IsManualTrigger 据此放行,不受 AUTO 冷却窗口限制。
  170. var (status, body) = await _rebuild.EnqueueAsync(
  171. input.ModuleCode, input.TenantId, input.FactoryId, _userManager.UserId, "MANUAL");
  172. return new { status, body };
  173. }
  174. public sealed class AdminCancelInput
  175. {
  176. public long JobId { get; set; }
  177. }
  178. [DisplayName("超管:取消重算")]
  179. [HttpPost("mdp-monitor/admin/cancel")]
  180. public async Task<object> CancelRebuildAsync([FromBody] AdminCancelInput input)
  181. {
  182. RequireSuperAdmin();
  183. if (input == null || input.JobId <= 0)
  184. throw Oops.Oh("必须指定任务");
  185. var ok = await _rebuildJobs.RequestCancelAsync(input.JobId);
  186. return new { ok };
  187. }
  188. [DisplayName("出站推送总开关告警")]
  189. [HttpGet("mdp-monitor/outbound-gate")]
  190. public async Task<object> GetOutboundGate()
  191. {
  192. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  193. var pending = await _db.Queryable<MdpOutbox>()
  194. .Where(x => x.TenantId == tenantId && x.Status == 0)
  195. .CountAsync();
  196. var enabled = _outboundGate.IsEnabled;
  197. return new
  198. {
  199. enabled,
  200. alert = !enabled,
  201. code = enabled ? "OUTBOUND_ON" : "OUTBOUND_DISABLED",
  202. message = enabled ? null : MdpOutboundGate.DisabledMessage,
  203. refusedCount = _outboundGate.RefusedCount,
  204. pendingCount = pending
  205. };
  206. }
  207. [DisplayName("发运对账差异")]
  208. [HttpGet("mdp-monitor/ship-recon")]
  209. public async Task<object> GetShipRecon()
  210. {
  211. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  212. var table = await _db.Ado.GetIntAsync(
  213. "SELECT COUNT(*) FROM information_schema.TABLES WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME='mdp_ship_recon_diff'");
  214. if (table == 0)
  215. return new { alert = false, diffCount = 0, message = "对账表尚未建立" };
  216. var n = await _db.Ado.GetIntAsync(
  217. "SELECT COUNT(*) FROM mdp_ship_recon_diff WHERE tenant_id=@t AND ABS(diff_qty)>0.0001",
  218. new SugarParameter("@t", tenantId));
  219. return new
  220. {
  221. alert = n > 0,
  222. diffCount = n,
  223. message = n > 0 ? $"销售发运对账有 {n} 条数量差异" : null
  224. };
  225. }
  226. [DisplayName("MDP模块选项")]
  227. [HttpGet("mdp-monitor/modules")]
  228. public object GetModules() => BuildCatalogResponse();
  229. [DisplayName("MDP任务目录")]
  230. [HttpGet("mdp-monitor/catalog")]
  231. public object GetCatalog() => BuildCatalogResponse();
  232. /// <summary>
  233. /// 本实例是否为 ETL 执行机、全库存活实例清单,以及全库重算并发的实时占用。
  234. ///
  235. /// <para><b>为什么要暴露</b>:不是执行机的表现是「所有定时作业静默不跑」——
  236. /// 页面上只能看到「最近执行」时间不再前进,看不出原因是本实例没被指派、还是作业跑挂了。
  237. /// 没有这个接口,未指派就是一个不可观测的状态。</para>
  238. ///
  239. /// <para><b>指派模型</b>:执行机由 <c>ado_etl_instance.is_runner</c> 人工指派,
  240. /// 不再由配置决定;<c>assignmentKnown=false</c> 表示注册器尚未心跳成功或已过期,
  241. /// 此时 <c>isRunner</c> 取自配置回退,页面必须把这种情况与「确实未被指派」区分开,
  242. /// 否则运维会误以为刚点的指派已经生效。</para>
  243. ///
  244. /// <para>只读进程内状态、实例注册表与一条 COUNT。实例是部署事实不属于任何租户,故不解析租户。</para>
  245. /// </summary>
  246. [DisplayName("ETL执行机状态")]
  247. [HttpGet("mdp-monitor/etl-runner")]
  248. public async Task<object> GetEtlRunnerStatus()
  249. {
  250. var runningScopes = await _db.Ado.GetIntAsync(
  251. "SELECT COUNT(*) FROM ado_module_dashboard_rebuild_job WHERE status='RUNNING'");
  252. var heartbeatCutoff = DateTime.Now - AidopRunnerState.FreshnessWindow;
  253. var instances = await _db.Queryable<AdoEtlInstance>()
  254. .OrderBy(x => x.MachineName)
  255. .OrderBy(x => x.ProcessId)
  256. .Select(x => new
  257. {
  258. x.InstanceId,
  259. x.MachineName,
  260. x.ProcessId,
  261. x.AppVersion,
  262. x.IsRunner,
  263. x.StartedAt,
  264. x.LastHeartbeatAt
  265. })
  266. .ToListAsync();
  267. return new
  268. {
  269. instanceId = AidopInstanceIdentity.InstanceId,
  270. isRunner = AidopJobGate.IsRunner,
  271. // 领取 AUTO / BOOTSTRAP 的资格与执行机身份同源,页面上单列一项避免误读为两回事
  272. claimAuto = AidopJobGate.IsRunner,
  273. assignmentKnown = AidopJobGate.IsAssignmentKnown,
  274. configKey = AidopJobGate.EnabledKey,
  275. envName = AidopJobGate.EnvEnabledName,
  276. maxParallelScopes = _rebuildCapability.MaxParallelScopes,
  277. globalMaxParallelScopes = _rebuildCapability.GlobalMaxParallelScopes,
  278. runningScopes,
  279. instances = instances.Select(x => new
  280. {
  281. x.InstanceId,
  282. x.MachineName,
  283. x.ProcessId,
  284. x.AppVersion,
  285. x.IsRunner,
  286. x.StartedAt,
  287. x.LastHeartbeatAt,
  288. alive = x.LastHeartbeatAt >= heartbeatCutoff,
  289. self = x.InstanceId == AidopInstanceIdentity.InstanceId
  290. })
  291. };
  292. }
  293. [DisplayName("MDP最近运行状态")]
  294. [HttpGet("mdp-monitor/latest")]
  295. public async Task<object> GetLatest([FromQuery] MdpMonitorQueryInput input)
  296. {
  297. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  298. var (whereSql, pars) = BuildWhere(input, tenantId);
  299. var row = await _db.Ado.SqlQuerySingleAsync<MdpMonitorRunLogRow>(
  300. $"""
  301. {SelectColumnsSql()}
  302. FROM mdp_transform_run_log
  303. WHERE {whereSql}
  304. ORDER BY start_time DESC, id DESC
  305. LIMIT 1
  306. """,
  307. pars)
  308. ?? new MdpMonitorRunLogRow();
  309. AttachCatalogInfo(row);
  310. return row;
  311. }
  312. [DisplayName("MDP运行日志列表")]
  313. [HttpGet("mdp-monitor/list")]
  314. public async Task<object> GetList([FromQuery] MdpMonitorListInput input)
  315. {
  316. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  317. var page = input.Page <= 0 ? 1 : input.Page;
  318. var pageSize = input.PageSize <= 0 ? 10 : input.PageSize;
  319. var offset = (page - 1) * pageSize;
  320. var (whereSql, pars) = BuildWhere(input, tenantId);
  321. var total = await _db.Ado.GetIntAsync($"SELECT COUNT(1) FROM mdp_transform_run_log WHERE {whereSql}", pars);
  322. var list = await _db.Ado.SqlQueryAsync<MdpMonitorRunLogRow>(
  323. $"""
  324. {SelectColumnsSql()}
  325. FROM mdp_transform_run_log
  326. WHERE {whereSql}
  327. ORDER BY start_time DESC, id DESC
  328. LIMIT {pageSize} OFFSET {offset}
  329. """,
  330. pars);
  331. foreach (var row in list)
  332. AttachCatalogInfo(row);
  333. return new { total, page, pageSize, list };
  334. }
  335. [DisplayName("MDP运行日志详情")]
  336. [HttpGet("mdp-monitor/detail/{id}")]
  337. public async Task<object> GetDetail(long id, [FromQuery] MdpMonitorQueryInput input)
  338. {
  339. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  340. var (whereSql, pars) = BuildWhere(input, tenantId);
  341. pars.Add(new SugarParameter("@Id", id));
  342. var row = await _db.Ado.SqlQuerySingleAsync<MdpMonitorRunLogRow>(
  343. $"""
  344. {SelectColumnsSql()}
  345. FROM mdp_transform_run_log
  346. WHERE id=@Id AND {whereSql}
  347. LIMIT 1
  348. """,
  349. pars);
  350. if (row == null)
  351. throw Oops.Oh("运行日志不存在");
  352. AttachCatalogInfo(row);
  353. return row;
  354. }
  355. [DisplayName("MDP同步链路详情")]
  356. [HttpGet("mdp-monitor/lineage")]
  357. public async Task<object> GetLineage([FromQuery] MdpMonitorLineageInput input)
  358. {
  359. var tenantId = AidopTenantScope.ResolveOrThrow(_userManager);
  360. var moduleCode = ResolveModuleCode(input.ModuleCode, input.JobCode);
  361. if (string.IsNullOrWhiteSpace(moduleCode))
  362. throw Oops.Oh("请选择 MDP 模块");
  363. var jobCode = ResolveJobCode(moduleCode, input.JobCode);
  364. var entityPrefix = $"{moduleCode}_%";
  365. var entities = await _db.Ado.SqlQueryAsync<MdpLineageEntityRow>(
  366. """
  367. SELECT e.id AS Id, e.entity_code AS EntityCode, e.entity_name AS EntityName,
  368. e.entity_type AS EntityType, s.source_code AS SourceCode,
  369. s.source_name AS SourceName, s.source_type AS SourceType,
  370. s.db_type AS SourceDbType, s.db_host AS SourceDbHost,
  371. s.db_port AS SourceDbPort, s.db_name AS SourceDbName,
  372. e.source_table_name AS SourceTableName, e.source_api_path AS SourceApiPath,
  373. s.db_type AS TargetDbType, s.db_host AS TargetDbHost,
  374. s.db_port AS TargetDbPort, s.db_name AS TargetDbName,
  375. e.target_table_name AS TargetTableName, e.sync_mode AS SyncMode,
  376. e.incr_column AS IncrColumn, e.status AS Status
  377. FROM mdp_entity e
  378. LEFT JOIN mdp_source s ON s.id = e.source_id
  379. WHERE e.entity_code LIKE @EntityPrefix
  380. AND (e.tenant_id = @TenantId OR e.tenant_id = 0)
  381. ORDER BY e.entity_code
  382. """,
  383. new SugarParameter("@EntityPrefix", entityPrefix),
  384. new SugarParameter("@TenantId", tenantId));
  385. if (entities.Count == 0)
  386. {
  387. var emptyOutput = new MdpLineageOutput
  388. {
  389. ModuleCode = moduleCode,
  390. JobCode = jobCode,
  391. BatchId = input.BatchId,
  392. Stages = BuildStageDescriptions(jobCode, moduleCode),
  393. Entities = new List<MdpLineageEntityRow>()
  394. };
  395. AttachLineageCatalogInfo(emptyOutput);
  396. return emptyOutput;
  397. }
  398. var entityIds = string.Join(",", entities.Select(u => u.Id));
  399. var mappings = await _db.Ado.SqlQueryAsync<MdpLineageFieldMappingRow>(
  400. $"""
  401. SELECT entity_id AS EntityId, source_field AS SourceField, target_field AS TargetField,
  402. field_type AS FieldType, transform_script AS TransformScript,
  403. const_value AS ConstValue, lookup_table AS LookupTable,
  404. is_required AS IsRequired, default_value AS DefaultValue, sort_order AS SortOrder
  405. FROM mdp_field_mapping
  406. WHERE entity_id IN ({entityIds})
  407. ORDER BY entity_id, sort_order, target_field
  408. """);
  409. var mappingsByEntity = mappings.GroupBy(u => u.EntityId).ToDictionary(u => u.Key, u => u.ToList());
  410. Dictionary<long, MdpLineageSyncLogRow> syncLogsByEntity = new();
  411. if (!string.IsNullOrWhiteSpace(input.BatchId))
  412. {
  413. var syncLogs = await _db.Ado.SqlQueryAsync<MdpLineageSyncLogRow>(
  414. """
  415. SELECT entity_id AS EntityId, entity_name AS EntityName, status AS Status,
  416. rows_read AS RowsRead, rows_insert AS RowsInsert, rows_update AS RowsUpdate,
  417. rows_skip AS RowsSkip, rows_error AS RowsError,
  418. sync_start AS SyncStart, sync_end AS SyncEnd, duration_ms AS DurationMs,
  419. error_msg AS ErrorMsg
  420. FROM mdp_sync_log
  421. WHERE sync_batch_id = @BatchId
  422. AND (tenant_id = @TenantId OR tenant_id = 0)
  423. AND entity_id IN (
  424. SELECT id FROM mdp_entity
  425. WHERE entity_code LIKE @EntityPrefix
  426. AND (tenant_id = @TenantId OR tenant_id = 0)
  427. )
  428. ORDER BY sync_start, id
  429. """,
  430. new SugarParameter("@BatchId", input.BatchId.Trim()),
  431. new SugarParameter("@EntityPrefix", entityPrefix),
  432. new SugarParameter("@TenantId", tenantId));
  433. syncLogsByEntity = syncLogs
  434. .GroupBy(u => u.EntityId)
  435. .ToDictionary(u => u.Key, u => u.OrderByDescending(x => x.SyncStart).First());
  436. }
  437. var batchId = input.BatchId?.Trim();
  438. foreach (var entity in entities)
  439. {
  440. entity.SourceFullName = BuildObjectFullName(entity.SourceDbType, entity.SourceDbHost, entity.SourceDbPort, entity.SourceDbName, entity.SourceTableName ?? entity.SourceApiPath);
  441. entity.TargetFullName = BuildObjectFullName(entity.TargetDbType, entity.TargetDbHost, entity.TargetDbPort, entity.TargetDbName, entity.TargetTableName);
  442. if (mappingsByEntity.TryGetValue(entity.Id, out var entityMappings) && entityMappings.Count > 0)
  443. {
  444. foreach (var mapping in entityMappings)
  445. {
  446. mapping.MappingSource = "CONFIG";
  447. mapping.IsFallback = false;
  448. }
  449. entity.FieldMappings = entityMappings;
  450. entity.FieldMappingCount = entityMappings.Count;
  451. }
  452. else
  453. {
  454. entity.FieldMappings = BuildFallbackFieldMappings(entity, batchId);
  455. entity.FieldMappingCount = entity.FieldMappings.Count;
  456. }
  457. if (syncLogsByEntity.TryGetValue(entity.Id, out var syncLog))
  458. entity.SyncLog = syncLog;
  459. }
  460. var output = new MdpLineageOutput
  461. {
  462. ModuleCode = moduleCode,
  463. JobCode = jobCode,
  464. BatchId = input.BatchId,
  465. Stages = BuildStageDescriptions(jobCode, moduleCode),
  466. Entities = entities
  467. };
  468. AttachLineageCatalogInfo(output);
  469. return output;
  470. }
  471. private static (string WhereSql, List<SugarParameter> Parameters) BuildWhere(MdpMonitorQueryInput input, long tenantId)
  472. {
  473. var where = new List<string> { BuildMdpRunLogTenantWhere(tenantId) };
  474. var pars = new List<SugarParameter> { new("@TenantId", tenantId) };
  475. var (jobFilterSql, jobFilterPars, noMatch) = BuildJobCodeFilter(input);
  476. if (noMatch)
  477. where.Add("1=0");
  478. else
  479. {
  480. where.Add(jobFilterSql);
  481. pars.AddRange(jobFilterPars);
  482. }
  483. if (!string.IsNullOrWhiteSpace(input.BatchId))
  484. {
  485. where.Add("batch_id LIKE @BatchId");
  486. pars.Add(new SugarParameter("@BatchId", $"%{input.BatchId.Trim()}%"));
  487. }
  488. if (!string.IsNullOrWhiteSpace(input.Status))
  489. {
  490. where.Add("status=@Status");
  491. pars.Add(new SugarParameter("@Status", input.Status.Trim().ToUpperInvariant()));
  492. }
  493. if (input.StartTime.HasValue)
  494. {
  495. where.Add("start_time >= @StartTime");
  496. pars.Add(new SugarParameter("@StartTime", input.StartTime.Value));
  497. }
  498. if (input.EndTime.HasValue)
  499. {
  500. where.Add("start_time <= @EndTime");
  501. pars.Add(new SugarParameter("@EndTime", input.EndTime.Value));
  502. }
  503. return (string.Join(" AND ", where), pars);
  504. }
  505. /// <summary>
  506. /// MDP 转换任务当前以 tenant_id=0 写入运行日志;登录租户查询时需兼容这类全局任务记录。
  507. /// </summary>
  508. internal static string BuildMdpRunLogTenantWhere(long tenantId) =>
  509. tenantId > 0 ? "(tenant_id = @TenantId OR tenant_id = 0)" : "tenant_id = 0";
  510. private static object BuildCatalogResponse()
  511. {
  512. return JobCatalog.Values
  513. .OrderBy(u => u.JobCode, StringComparer.OrdinalIgnoreCase)
  514. .Select(u => new
  515. {
  516. jobCode = u.JobCode,
  517. displayName = u.DisplayName,
  518. businessDomainCode = u.BusinessDomainCode,
  519. businessDomainName = u.BusinessDomainName,
  520. consumerModules = u.ConsumerModules,
  521. scopeType = u.ScopeType,
  522. moduleCode = ResolveModuleCodeFromJobCode(u.JobCode),
  523. scheduleJobId = u.ScheduleJobId
  524. })
  525. .ToList();
  526. }
  527. private static (string Sql, List<SugarParameter> Parameters, bool NoMatch) BuildJobCodeFilter(MdpMonitorQueryInput input)
  528. {
  529. var pars = new List<SugarParameter>();
  530. if (!string.IsNullOrWhiteSpace(input.JobCode))
  531. {
  532. pars.Add(new SugarParameter("@JobCode", input.JobCode.Trim().ToUpperInvariant()));
  533. return ("job_code=@JobCode", pars, false);
  534. }
  535. var hasDomainFilter = !string.IsNullOrWhiteSpace(input.BusinessDomainCode);
  536. var consumerModule = !string.IsNullOrWhiteSpace(input.ConsumerModule)
  537. ? input.ConsumerModule.Trim()
  538. : !string.IsNullOrWhiteSpace(input.ModuleCode) ? input.ModuleCode.Trim() : null;
  539. var hasConsumerFilter = !string.IsNullOrWhiteSpace(consumerModule);
  540. if (!hasDomainFilter && !hasConsumerFilter)
  541. return ("IFNULL(job_code, '') LIKE '%MDP%'", pars, false);
  542. var allowed = ResolveAllowedJobCodes(input.BusinessDomainCode, consumerModule);
  543. if (allowed.Count == 0)
  544. return (string.Empty, pars, true);
  545. if (allowed.Count == 1)
  546. {
  547. pars.Add(new SugarParameter("@JobCode", allowed.First()));
  548. return ("job_code=@JobCode", pars, false);
  549. }
  550. var inParts = new List<string>();
  551. var index = 0;
  552. foreach (var code in allowed.OrderBy(u => u, StringComparer.OrdinalIgnoreCase))
  553. {
  554. var paramName = $"@JobCode{index++}";
  555. inParts.Add(paramName);
  556. pars.Add(new SugarParameter(paramName, code));
  557. }
  558. return ($"job_code IN ({string.Join(", ", inParts)})", pars, false);
  559. }
  560. private static HashSet<string> ResolveAllowedJobCodes(string? businessDomainCode, string? consumerModule)
  561. {
  562. IEnumerable<MdpJobCatalogItem> items = JobCatalog.Values;
  563. if (!string.IsNullOrWhiteSpace(businessDomainCode))
  564. {
  565. var domain = businessDomainCode.Trim();
  566. items = items.Where(u => string.Equals(u.BusinessDomainCode, domain, StringComparison.OrdinalIgnoreCase));
  567. }
  568. if (!string.IsNullOrWhiteSpace(consumerModule))
  569. {
  570. var module = consumerModule.Trim();
  571. items = items.Where(u => u.ConsumerModules.Contains(module, StringComparer.OrdinalIgnoreCase));
  572. }
  573. return items.Select(u => u.JobCode).ToHashSet(StringComparer.OrdinalIgnoreCase);
  574. }
  575. private static void AttachCatalogInfo(MdpMonitorRunLogRow row)
  576. {
  577. if (row == null || string.IsNullOrWhiteSpace(row.JobCode))
  578. return;
  579. if (!JobCatalog.TryGetValue(row.JobCode.Trim(), out var item))
  580. return;
  581. row.BusinessDomainCode = item.BusinessDomainCode;
  582. row.BusinessDomainName = item.BusinessDomainName;
  583. row.ConsumerModules = string.Join(",", item.ConsumerModules);
  584. row.ScopeType = item.ScopeType;
  585. row.DisplayName = item.DisplayName;
  586. row.ScheduleJobId = item.ScheduleJobId;
  587. }
  588. private static void AttachLineageCatalogInfo(MdpLineageOutput output)
  589. {
  590. if (output == null || string.IsNullOrWhiteSpace(output.JobCode))
  591. return;
  592. if (!JobCatalog.TryGetValue(output.JobCode.Trim(), out var item))
  593. return;
  594. output.BusinessDomainCode = item.BusinessDomainCode;
  595. output.BusinessDomainName = item.BusinessDomainName;
  596. output.ConsumerModules = string.Join(",", item.ConsumerModules);
  597. output.ScopeType = item.ScopeType;
  598. output.DisplayName = item.DisplayName;
  599. output.ScheduleJobId = item.ScheduleJobId;
  600. }
  601. private static string FormatConsumerModulesLabel(MdpJobCatalogItem item) =>
  602. string.Join("、", item.ConsumerModules);
  603. private static List<MdpLineageStageRow> BuildStageDescriptions(string? jobCode, string? moduleCode)
  604. {
  605. if (!string.IsNullOrWhiteSpace(jobCode) && JobCatalog.TryGetValue(jobCode.Trim(), out var catalogItem))
  606. {
  607. return catalogItem.BusinessDomainCode switch
  608. {
  609. "order_delivery" => BuildOrderDeliveryStages(catalogItem),
  610. "work_schedule" => BuildWorkScheduleStages(catalogItem),
  611. "supply_purchase" => BuildSupplyPurchaseStages(catalogItem),
  612. "purchase_execution" => BuildPurchaseExecutionStages(catalogItem),
  613. _ => BuildGenericStages(moduleCode)
  614. };
  615. }
  616. return BuildGenericStages(moduleCode);
  617. }
  618. private static List<MdpLineageStageRow> BuildOrderDeliveryStages(MdpJobCatalogItem item)
  619. {
  620. var consumers = FormatConsumerModulesLabel(item);
  621. return new List<MdpLineageStageRow>
  622. {
  623. 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" },
  624. 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" },
  625. new() { StageCode = "DWD", StageName = "订单交付域 · DWD宽表", Layer = "dwd", Description = $"沉淀订单交付事实,供 {consumers} 看板与诊断读取。", InputObjects = "mdp_std_so, mdp_std_ship_trans", OutputObjects = "dwd_ship_trans", Execution = "S1MdpSyncTransformService.BuildDwdAsync" },
  626. 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" }
  627. };
  628. }
  629. private static List<MdpLineageStageRow> BuildWorkScheduleStages(MdpJobCatalogItem item)
  630. {
  631. var consumers = FormatConsumerModulesLabel(item);
  632. return new List<MdpLineageStageRow>
  633. {
  634. new() { StageCode = "STAGING", StageName = "工单排程域 · 贴源同步", Layer = "mdp_stg", Description = $"按 mdp_entity 登记工单排程域源对象抽取数据;产出供 {consumers} 消费。", InputObjects = "工单 / 工序 / 排程源对象", OutputObjects = "mdp_stg_*", Execution = "S2MdpSyncTransformService.SyncStagingAsync" },
  635. new() { StageCode = "STANDARD", StageName = "工单排程域 · 标准层转换", Layer = "mdp_std", Description = "将工单、工序、排程等对象标准化。", InputObjects = "mdp_stg_*", OutputObjects = "mdp_std_*", Execution = "S2MdpSyncTransformService.BuildStandardCommands" },
  636. new() { StageCode = "DWD", StageName = "工单排程域 · DWD宽表", Layer = "dwd", Description = $"生成制造执行与排程分析宽表,供 {consumers} 读取。", InputObjects = "mdp_std_*", OutputObjects = "dwd_*", Execution = "S2MdpSyncTransformService.BuildDwdAsync" },
  637. new() { StageCode = "KPI", StageName = "指标写入", Layer = "ado_s9", Description = $"写入工单排程域 KPI,供 {consumers} 消费。", InputObjects = "mdp_std_* / dwd_*", OutputObjects = "ado_s9_kpi_value_*", Execution = "S2MdpSyncTransformService.BuildS2KpiValuesAsync" }
  638. };
  639. }
  640. private static List<MdpLineageStageRow> BuildSupplyPurchaseStages(MdpJobCatalogItem item)
  641. {
  642. var consumers = FormatConsumerModulesLabel(item);
  643. return new List<MdpLineageStageRow>
  644. {
  645. new() { StageCode = "STAGING", StageName = "供应采购域 · 贴源同步", Layer = "mdp_stg", Description = $"按 mdp_entity 登记供应采购域源对象抽取数据;产出供 {consumers} 消费。", InputObjects = "供应 / 物料 / 采购源对象", OutputObjects = "mdp_stg_*", Execution = "S3MdpSyncTransformService.SyncStagingAsync" },
  646. new() { StageCode = "STANDARD", StageName = "供应采购域 · 标准层转换", Layer = "mdp_std", Description = "将供应、物料、采购、交货计划等对象标准化。", InputObjects = "mdp_stg_*", OutputObjects = "mdp_std_*", Execution = "S3MdpSyncTransformService.BuildStandardCommands" },
  647. new() { StageCode = "DWD", StageName = "供应采购域 · DWD宽表", Layer = "dwd", Description = $"生成供应交付、齐套、风险等分析宽表,供 {consumers} 读取。", InputObjects = "mdp_std_*", OutputObjects = "dwd_supplier_delivery / dwd_material_readiness 等", Execution = "S3MdpSyncTransformService.BuildDwdAsync" },
  648. new() { StageCode = "KPI", StageName = "指标写入", Layer = "ado_s9", Description = $"写入供应采购域指标,供 {consumers} 消费。", InputObjects = "mdp_std_* / dwd_*", OutputObjects = "ado_s9_kpi_value_*", Execution = "S3MdpSyncTransformService.BuildS3KpiValuesAsync" }
  649. };
  650. }
  651. private static List<MdpLineageStageRow> BuildPurchaseExecutionStages(MdpJobCatalogItem item)
  652. {
  653. var consumers = FormatConsumerModulesLabel(item);
  654. return new List<MdpLineageStageRow>
  655. {
  656. 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" },
  657. 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" },
  658. 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" },
  659. 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" }
  660. };
  661. }
  662. /// <summary>
  663. /// 贴源同步在代码中维护、但未写入 mdp_field_mapping 的实体主键/业务键提示(仅监控展示兜底)。
  664. /// </summary>
  665. private static readonly Dictionary<string, (string SourceRowId, string SourceBizKeyExpr)> KnownStagingEntityKeys =
  666. new(StringComparer.OrdinalIgnoreCase)
  667. {
  668. ["S4_IQC_RECEIPT"] = ("RecID", "CONCAT(IFNULL(s.`Domain`,''), ':', IFNULL(s.`Receiver`,''), ':', IFNULL(s.`Line`,''))"),
  669. ["S4_SHIPMENT_EXEC"] = ("id", "CONCAT(IFNULL(s.`glid`,''), ':', IFNULL(s.`id`,''))"),
  670. ["S4_RETURN_EXEC"] = ("Id", "s.`dsnum`"),
  671. ["S4_SHORTAGE_EXEC"] = ("id", "CONCAT(IFNULL(s.`work_order`,''), ':', IFNULL(s.`component_item_code`,''))")
  672. };
  673. private static List<MdpLineageFieldMappingRow> BuildFallbackFieldMappings(MdpLineageEntityRow entity, string? batchId)
  674. {
  675. var sourceRowId = entity.IncrColumn;
  676. string? sourceBizKeyExpr = null;
  677. if (!string.IsNullOrWhiteSpace(entity.EntityCode) &&
  678. KnownStagingEntityKeys.TryGetValue(entity.EntityCode, out var known))
  679. {
  680. sourceRowId ??= known.SourceRowId;
  681. sourceBizKeyExpr = known.SourceBizKeyExpr;
  682. }
  683. var fallbackNote = "当前实体未配置逐字段映射,贴源同步保留源行 raw_data";
  684. var mappings = new List<MdpLineageFieldMappingRow>();
  685. var sort = 10;
  686. void Add(string sourceField, string targetField, string fieldType, string? transformScript, string? constValue = null, bool isRequired = false)
  687. {
  688. mappings.Add(new MdpLineageFieldMappingRow
  689. {
  690. EntityId = entity.Id,
  691. SourceField = sourceField,
  692. TargetField = targetField,
  693. FieldType = fieldType,
  694. TransformScript = transformScript,
  695. ConstValue = constValue,
  696. IsRequired = isRequired,
  697. SortOrder = sort,
  698. MappingSource = "FALLBACK",
  699. IsFallback = true
  700. });
  701. sort += 10;
  702. }
  703. Add("tenant_id", "tenant_id", "DIRECT", "当前租户或全局租户兜底");
  704. Add($"CONST:{entity.SourceCode ?? MdpSourceIdentity.Native}", "source_system", "CONST", "来源数据源编码", entity.SourceCode ?? MdpSourceIdentity.Native, isRequired: true);
  705. Add($"CONST:{entity.SourceTableName ?? entity.SourceApiPath ?? "--"}", "source_table", "CONST", "源表名", entity.SourceTableName ?? entity.SourceApiPath, isRequired: true);
  706. if (!string.IsNullOrWhiteSpace(sourceRowId))
  707. Add(sourceRowId, "source_row_id", "DIRECT", "来自 mdp_entity.incr_column 或实体主键配置");
  708. else
  709. Add("--", "source_row_id", "DIRECT", fallbackNote);
  710. if (!string.IsNullOrWhiteSpace(sourceBizKeyExpr))
  711. Add(sourceBizKeyExpr, "source_biz_key", "EXPR", "来自实体业务键表达式", isRequired: true);
  712. else
  713. Add("--", "source_biz_key", "EXPR", fallbackNote);
  714. Add("*", "raw_data", "JSON", "源行整行 JSON");
  715. Add("CONST", "sync_batch_id", "CONST", "当前同步批次", string.IsNullOrWhiteSpace(batchId) ? "当前同步批次" : batchId, isRequired: true);
  716. Add("NOW()", "sync_time", "CONST", "同步时间", isRequired: true);
  717. return mappings;
  718. }
  719. private static List<MdpLineageStageRow> BuildGenericStages(string? moduleCode)
  720. {
  721. if (string.Equals(moduleCode, "S1", StringComparison.OrdinalIgnoreCase))
  722. return BuildOrderDeliveryStages(JobCatalog["S1_MDP_SYNC_TRANSFORM"]);
  723. if (string.Equals(moduleCode, "S2", StringComparison.OrdinalIgnoreCase))
  724. return BuildWorkScheduleStages(JobCatalog["S2_MDP_SYNC_TRANSFORM"]);
  725. if (string.Equals(moduleCode, "S3", StringComparison.OrdinalIgnoreCase))
  726. return BuildSupplyPurchaseStages(JobCatalog["S3_MDP_SYNC_TRANSFORM"]);
  727. if (string.Equals(moduleCode, "S4", StringComparison.OrdinalIgnoreCase))
  728. return BuildPurchaseExecutionStages(JobCatalog["S4_MDP_SYNC_TRANSFORM"]);
  729. return new List<MdpLineageStageRow>
  730. {
  731. new() { StageCode = "STAGING", StageName = "贴源同步", Layer = "mdp_stg", Description = "按 mdp_entity 登记源对象抽取数据。", InputObjects = "源对象", OutputObjects = "mdp_stg_*", Execution = "MDP 同步服务" },
  732. new() { StageCode = "STANDARD", StageName = "标准层转换", Layer = "mdp_std", Description = "标准层/DWD/KPI 当前由后端 Service 承载。", InputObjects = "mdp_stg_*", OutputObjects = "mdp_std_* / dwd_* / 指标表", Execution = "MDP 转换服务" }
  733. };
  734. }
  735. private static string? ResolveModuleCodeFromJobCode(string? jobCode)
  736. {
  737. if (string.IsNullOrWhiteSpace(jobCode))
  738. return null;
  739. return ModuleJobCodes.FirstOrDefault(u =>
  740. string.Equals(u.Value, jobCode.Trim(), StringComparison.OrdinalIgnoreCase)).Key;
  741. }
  742. private static string? ResolveJobCode(string? moduleCode, string? jobCode)
  743. {
  744. if (!string.IsNullOrWhiteSpace(jobCode))
  745. return jobCode.Trim().ToUpperInvariant();
  746. if (string.IsNullOrWhiteSpace(moduleCode))
  747. return null;
  748. return ModuleJobCodes.TryGetValue(moduleCode.Trim(), out var mapped) ? mapped : null;
  749. }
  750. private static string? ResolveModuleCode(string? moduleCode, string? jobCode)
  751. {
  752. if (!string.IsNullOrWhiteSpace(moduleCode))
  753. return moduleCode.Trim().ToUpperInvariant();
  754. if (string.IsNullOrWhiteSpace(jobCode))
  755. return null;
  756. var normalizedJobCode = jobCode.Trim();
  757. return ModuleJobCodes.FirstOrDefault(u => string.Equals(u.Value, normalizedJobCode, StringComparison.OrdinalIgnoreCase)).Key;
  758. }
  759. private static string? BuildObjectFullName(string? dbType, string? host, int? port, string? dbName, string? objectName)
  760. {
  761. if (string.IsNullOrWhiteSpace(objectName))
  762. return null;
  763. var databaseObject = string.IsNullOrWhiteSpace(dbName) ? objectName : $"{dbName}.{objectName}";
  764. var hostPart = string.IsNullOrWhiteSpace(host) ? null : port.HasValue ? $"{host}:{port}" : host;
  765. return string.Join(" / ", new[] { dbType, hostPart, databaseObject }.Where(u => !string.IsNullOrWhiteSpace(u)));
  766. }
  767. private static string SelectColumnsSql()
  768. {
  769. return """
  770. SELECT id AS Id, tenant_id AS TenantId, job_code AS JobCode, job_name AS JobName, trigger_type AS TriggerType,
  771. batch_id AS BatchId, status AS Status, start_time AS StartTime, end_time AS EndTime, duration_ms AS DurationMs,
  772. stage_rows AS StageRows, standard_rows AS StandardRows, dwd_rows AS DwdRows,
  773. error_message AS ErrorMessage, summary_json AS SummaryJson, create_time AS CreateTime, update_time AS UpdateTime
  774. """;
  775. }
  776. }
  777. public sealed record MdpJobCatalogItem(
  778. string JobCode,
  779. string BusinessDomainCode,
  780. string BusinessDomainName,
  781. string[] ConsumerModules,
  782. string ScopeType,
  783. string DisplayName,
  784. string? ScheduleJobId);
  785. public class MdpMonitorQueryInput
  786. {
  787. public string? BusinessDomainCode { get; set; }
  788. public string? ConsumerModule { get; set; }
  789. public string? ModuleCode { get; set; }
  790. public string? JobCode { get; set; }
  791. public string? BatchId { get; set; }
  792. public string? Status { get; set; }
  793. public DateTime? StartTime { get; set; }
  794. public DateTime? EndTime { get; set; }
  795. }
  796. public sealed class MdpMonitorListInput : MdpMonitorQueryInput
  797. {
  798. public int Page { get; set; } = 1;
  799. public int PageSize { get; set; } = 10;
  800. }
  801. public sealed class MdpMonitorLineageInput : MdpMonitorQueryInput
  802. {
  803. }
  804. public sealed class MdpMonitorRunLogRow
  805. {
  806. public long Id { get; set; }
  807. public long TenantId { get; set; }
  808. public string? JobCode { get; set; }
  809. public string? JobName { get; set; }
  810. public string? TriggerType { get; set; }
  811. public string? BatchId { get; set; }
  812. public string? Status { get; set; }
  813. public DateTime? StartTime { get; set; }
  814. public DateTime? EndTime { get; set; }
  815. public int? DurationMs { get; set; }
  816. public int? StageRows { get; set; }
  817. public int? StandardRows { get; set; }
  818. public int? DwdRows { get; set; }
  819. public string? ErrorMessage { get; set; }
  820. public string? SummaryJson { get; set; }
  821. public DateTime? CreateTime { get; set; }
  822. public DateTime? UpdateTime { get; set; }
  823. public string? BusinessDomainCode { get; set; }
  824. public string? BusinessDomainName { get; set; }
  825. public string? ConsumerModules { get; set; }
  826. public string? ScopeType { get; set; }
  827. public string? DisplayName { get; set; }
  828. public string? ScheduleJobId { get; set; }
  829. }
  830. public sealed class MdpLineageOutput
  831. {
  832. public string? ModuleCode { get; set; }
  833. public string? JobCode { get; set; }
  834. public string? BatchId { get; set; }
  835. public string? BusinessDomainCode { get; set; }
  836. public string? BusinessDomainName { get; set; }
  837. public string? ConsumerModules { get; set; }
  838. public string? ScopeType { get; set; }
  839. public string? DisplayName { get; set; }
  840. public string? ScheduleJobId { get; set; }
  841. public List<MdpLineageStageRow> Stages { get; set; } = new();
  842. public List<MdpLineageEntityRow> Entities { get; set; } = new();
  843. }
  844. public sealed class MdpLineageStageRow
  845. {
  846. public string? StageCode { get; set; }
  847. public string? StageName { get; set; }
  848. public string? Layer { get; set; }
  849. public string? Description { get; set; }
  850. public string? InputObjects { get; set; }
  851. public string? OutputObjects { get; set; }
  852. public string? Execution { get; set; }
  853. }
  854. public sealed class MdpLineageEntityRow
  855. {
  856. public long Id { get; set; }
  857. public string? EntityCode { get; set; }
  858. public string? EntityName { get; set; }
  859. public string? EntityType { get; set; }
  860. public string? SourceCode { get; set; }
  861. public string? SourceName { get; set; }
  862. public string? SourceType { get; set; }
  863. public string? SourceDbType { get; set; }
  864. public string? SourceDbHost { get; set; }
  865. public int? SourceDbPort { get; set; }
  866. public string? SourceDbName { get; set; }
  867. public string? SourceTableName { get; set; }
  868. public string? SourceApiPath { get; set; }
  869. public string? SourceFullName { get; set; }
  870. public string? TargetDbType { get; set; }
  871. public string? TargetDbHost { get; set; }
  872. public int? TargetDbPort { get; set; }
  873. public string? TargetDbName { get; set; }
  874. public string? TargetTableName { get; set; }
  875. public string? TargetFullName { get; set; }
  876. public string? SyncMode { get; set; }
  877. public string? IncrColumn { get; set; }
  878. public int? Status { get; set; }
  879. public int FieldMappingCount { get; set; }
  880. public List<MdpLineageFieldMappingRow> FieldMappings { get; set; } = new();
  881. public MdpLineageSyncLogRow? SyncLog { get; set; }
  882. }
  883. public sealed class MdpLineageFieldMappingRow
  884. {
  885. public long EntityId { get; set; }
  886. public string? SourceField { get; set; }
  887. public string? TargetField { get; set; }
  888. public string? FieldType { get; set; }
  889. public string? TransformScript { get; set; }
  890. public string? ConstValue { get; set; }
  891. public string? LookupTable { get; set; }
  892. public bool IsRequired { get; set; }
  893. public string? DefaultValue { get; set; }
  894. public int SortOrder { get; set; }
  895. public string? MappingSource { get; set; }
  896. public bool IsFallback { get; set; }
  897. }
  898. public sealed class MdpLineageSyncLogRow
  899. {
  900. public long EntityId { get; set; }
  901. public string? EntityName { get; set; }
  902. public string? Status { get; set; }
  903. public long? RowsRead { get; set; }
  904. public long? RowsInsert { get; set; }
  905. public long? RowsUpdate { get; set; }
  906. public long? RowsSkip { get; set; }
  907. public long? RowsError { get; set; }
  908. public DateTime? SyncStart { get; set; }
  909. public DateTime? SyncEnd { get; set; }
  910. public int? DurationMs { get; set; }
  911. public string? ErrorMsg { get; set; }
  912. }