MdpMonitorService.cs 50 KB

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