FEATURED · 精选文章

SeaTunnel Zeta 运行时执行图:在作业 DAG 上叠加节点忙碌度与背压健康状态

发布时间 / 2026/9/19 23:39:43
来源 / 创域科博编辑部
栏目 / 资讯中心
SeaTunnel Zeta 运行时执行图:在作业 DAG 上叠加节点忙碌度与背压健康状态 SeaTunnel Zeta 运行时执行图:在作业 DAG 上叠加节点忙碌度与背压健康状态【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel Zeta 引擎的运行时执行图是运行中作业的统一诊断入口:它复用现有 Job Detail DAG 作为唯一图拓扑,把 Active Master 内存中的实时指标窗口叠加到节点和边上,让使用者无需切换页面就能判断哪个 vertex 在忙碌、哪条带队列的 edge 正在阻塞。本文基于仓库中的设计契约文档 runtime-execution-graph.md,结合引擎服务端源码,完整讲解该功能的问题边界、数据模型、REST 契约、刷新与保留机制,以及大 DAG 场景下的降级策略;读完后你可以掌握如何在现有可观测链路上定位作业 slowdown 的起点,并理解各字段的精确计算口径。问题背景:为什么要做运行时执行图SeaTunnel 已经暴露了作业拓扑、vertex 指标、edge 指标、checkpoint 历史、异常和日志。但运行中的作业出现变慢或阻塞时,使用者仍然需要一条统一路径回答这些问题:哪个 vertex 现在忙、空闲或等待;哪条带队列的 edge 正在阻塞或逐渐填满;当前瓶颈更接近 Source 读取、Transform 处理、Sink 写入、checkpoint 或 commit,还是外部系统;在进入表格和日志前,先判断 slowdown 大概率从哪里开始。设计契约对此的定位非常明确:运行时执行图应当作为诊断入口,而不是变成第二套监控系统。该设计刻意控制了四个边界:复用 Zeta 作业 DAG 作为唯一图拓扑;复用 Active Master 内存中的实时指标窗口展示节点和边的健康状态;复用现有指标 API、checkpoint REST 接口和 UI 详情面板做 drill-down;第一版(V1)不引入持久化指标快照,也不提供长期回放。现有基础:V1 组合的五个既有契约V1 运行时执行图不新增任何数据链路,它组合的是仓库中已经落地的 Zeta 契约:范围现有契约拓扑Job Detail 已经渲染JobDAGInfo拓扑,对应源码 JobDAGInfo。Vertex 指标/metrics/realtime/jobs/{jobId}/vertices返回最近窗口内 Source、Transform、Sink 的 bucket。Edge 指标/metrics/realtime/jobs/{jobId}/edges返回带queueId与targetVertexId的队列 edge bucket。Checkpoint 状态/jobs/checkpoints/:jobId和/jobs/checkpoints/history/:jobId暴露 checkpoint 概览与历史。异常与日志Job Detail、Exception 和 Log API 已经提供失败上下文。这些 REST 端点由 RealtimeMetricsServlet 提供,路由挂载在RestConstant.REST_URL_REALTIME_METRICS /metrics/realtime之下(见 RestConstant.java)。源码中可以确认窗口参数的处理:默认查询窗口为 3 分钟(DEFAULT_WINDOW_MS),最大 10 分钟(MAX_WINDOW_MS),超过上限的请求会被钳制到最大值,非法windowMs参数回退到默认值。数据生产端是 Active Master 上的 RealtimeMetricsService。该类由 master 选举生命周期负责启动与停止:一个单线程后台 collector 定期从 worker 拉取目标指标前缀,聚合成有界的内存 bucket,再向 REST 层暴露 best-effort 快照;shutdown 时会清掉所有缓存作业状态,避免 standby 或已停止节点提供陈旧数据。这正印证了设计文档中Active Master 只保留短窗口内存数据的表述。目标与非目标目标——运行时执行图要做到六件事:在 DAG 上直接展示节点健康状态;在 DAG 上直接展示队列 edge 健康状态;不切换页面也能看出最热 vertex 和最阻塞 edge;点击图上的节点或边,可以进入现有指标表格或详情抽屉;在图附近展示 checkpoint 与任务错误上下文,但不把它们变成 graph-only 契约;对大 DAG 提供可预期、成本更低的降级展示。非目标同样重要,划出了 V1 的边界:分布式 tracing;任意时间范围的历史回放;持久化运行时指标快照;connector 专属图组件;与 realtime observability 并行的新后端指标模型;替代人工判断的自动根因结论。运行时图数据模型拓扑:不发明第二套形状图拓扑直接来自当前作业 DAG,包含:jobId;vertexId;vertex 类型,例如 Source、Transform 或 Sink;vertex 之间的有向边;可与targetVertexId关联的 edge 元数据。设计契约强调:运行时执行图不能发明另一套拓扑。如果未来执行 DAG 发生变化,图应跟随引擎 DAG 契约,而不是维护自己的形状。拓扑载体是引擎核心的 JobDAGInfo,其中 vertex/edge 的结构定义在 LogicalEdge 与 Edge 中。节点运行状态:按 vertex 类型选取主视觉信号每个图节点通过vertexId合并最近一个 realtime vertex point。不同 vertex 类型使用不同的主视觉信号和辅助字段:Vertex 类型主要视觉信号辅助字段SourcesourceReadRatio和sourceIdleRatiosourceReadNs、sourceIdleNs、subtaskCountTransformtransformBusyRatiotransformProcessNsPerRecord、transformRecordsIn、transformRecordsOut、subtaskCountSinksinkBusyRatiosinkWriteNsPerRecord、sinkRecordsIn、sinkPrepareCommitNs、sinkCommitNs、sinkAbortNs、subtaskCount节点主颜色应代表最符合该 vertex 类型的指标:Source 用读取或空闲比例,Transform 用 transform busy ratio,Sink 用 sink busy ratio。这些字段的精确计算口径可以在RealtimeMetricsService的vertexPoint方法中得到印证:sourceReadRatio/sourceIdleRatio 对应纳秒增量 / (bucketMs转纳秒 ×subtaskCount),并钳制在 [0, 1];transformBusyRatiotransformProcessNs同口径比例;transformProcessNsPerRecordtransformProcessNs / transformRecordsIn;sinkBusyRatio注意只基于sinkWriteNs;sinkWriteNsPerRecordsinkWriteNs / sinkRecordsIn,commit 阶段耗时单独由sinkPrepareCommitNs、sinkCommitNs、sinkAbortNs字段承载。也就是说,一个 Sink 节点即使sinkBusyRatio不高,只要sinkCommitNs或sinkPrepareCommitNs显著,说明瓶颈可能在提交链路而不是写入本身——这正是文档要求把 commit 阶段字段作为辅助信号的原因。边运行状态:只有带队列的 edge 才有背压指标V1 只有带队列的 edge 才能暴露背压指标。每条图 edge 优先通过 REST 返回的targetVertexId合并最近一个 realtime edge point,必要时再从queueId解码。各字段含义:字段含义queueIdrealtime 聚合使用的稳定队列指标标识。targetVertexId用来把队列指标映射回 DAG edge 的下游 vertex。bpRatio当前 bucket 内生产端等待队列容量的时间占比。queueFillRatio最近一次采样到的队列填充比例。queueSize最近一次采样到的队列大小。queueCapacity最近一次采样到的队列容量。视觉映射规则:edge 颜色应代表bpRatio,edge 宽度应代表queueFillRatio;详情抽屉应展示原始字段和最近 bucket 序列。源码印证了两点关键实现。其一,bpRatio的计算(edgePoint方法)为:emitBlockedNs / (bucketMs 转纳秒 × subtaskCount),其中emitBlockedNs是 bucket 内累计的队列 put 阻塞纳秒(来自INTERMEDIATE_QUEUE_PUT_BLOCKED_NANOS计数的差分累加);queueFillRatio queueSize / queueCapacity,且当采样不一致出现queueSize queueCapacity时会被钳制到 capacity,避免渲染出超过 1 的填充比例。其二,targetVertexId的解码逻辑在decodeQueueTargetVertexId中:负偶数 queueId(异步边界队列,由 PhysicalPlanGenerator 生成,形式为-2 * actionId)解码为abs/2;负奇数(sink split 队列,形式为-(2 * actionId 1))解码为(abs - 1) / 2;非负 queueId 原样返回。这样 UI 就能把背压指标直接挂回 DAG 的下游 vertexId。采集与聚合链路从源码结构看,collector 的工作方式是:每 5 秒(POLL_INTERVAL_MS 5000)拉取一次 worker 指标,拉取超时 3 秒;只拉取 13 个目标前缀:三个队列指标(put blocked、size、capacity)加 Source/Transform/Sink 的 10 个计数器,避免全量指标回流 master;每个作业维护一个JobStore,内含按bucketMs对齐的 bucket 队列(默认 bucket 5 秒),计数器采用与上次快照的差分累加,因此 counter reset 或采样缺失时增量按 0 处理(Math.max(0, delta));超过retentionMinutes的 bucket 被逐出,全程只在内存中。作业级开关与参数由 ObservabilityConfig 统一解析,前缀为engine.observability.,可从该源码确认的默认值与约束:配置键默认值约束engine.observability.enabledfalse显式配置async_boundaries或split_sink_io且未显式enabledfalse时自动开启engine.observability.bucket_ms5000最小钳制为 1000msengine.observability.retention_minutes3最小 1、最大 10 分钟engine.observability.split_sink_iofalse为每个 sink 拆分独立队列阶段以暴露 sink 前背压engine.observability.edge_buffer_capacity0(用引擎默认)非负engine.observability.edge_overrides空单项 capacity 上限 100000,超限钳制并告警稳定契约与 Best-Effort 信号运行时图需要区分稳定标识字段和有诊断价值但来自采样的信号,这是设计契约中最容易被忽略、却决定 UI 表达方式的部分。稳定契约字段(可被用于对齐与匹配):jobIdvertexIdqueueIdtargetVertexIdbucketMsfromMstoMspoint 时间戳tssubtaskCountBest-effort 诊断字段(用于着色与趋势判断):busy ratio、idle ratio;单条记录耗时估算;queue size、queue fill ratio、producer wait ratio;checkpoint 与错误摘要 badge。设计文档要求:UI 必须把 best-effort 字段表达为实时诊断信号,而不是审计口径的统计值,因为在恢复、rescale、counter reset 或采样延迟附近它们可能波动。这一点与源码中的防御性处理一致:差分累加取非负、queueSize钳制、ratio 钳制到 [0,1],都是在承认采样信号不保证精确的前提下保证渲染不越界。刷新、保留与成本V1 应保持现有刷新和保留模型:worker counter 由 Active Master 收集;收集结果只保留短窗口内存数据;REST 查询窗口默认 3 分钟;REST 查询窗口最大 10 分钟;UI 可以比 collector 更频繁刷新,但必须能接受连续刷新拿到相同 bucket;默认不把 runtime graph 数据写入磁盘。这样可以让运行时图的成本跟随现有 realtime observability,而不是新增一条轮询或持久化链路。对应源码事实:collector 周期 5 秒(见POLL_INTERVAL_MS),toEdgesResponse/toVerticesResponse会把请求窗口与保留窗口取较小值(effectiveWindowMs min(requestedWindowMs, retentionMs)),REST 层再做一次 10 分钟上限钳制;当作业未开启 observability 时,接口返回enabledfalse与空的edges/vertices列表而不是报错,这保证了 UI 在指标未开启时有确定的降级行为。Checkpoint 与错误上下文:邻近但不混入Checkpoint 和错误信号应显示在图附近,但 V1 中保持独立契约:checkpoint 概览与历史继续来自 checkpoint REST 接口;作业异常与日志继续来自现有作业详情 API;图上可以展示小型状态 badge 或入口链接,但 drill-down 应打开现有详情面板。这样做的动机是避免在生命周期与保留规则尚未明确前,把 checkpoint 数据混入 realtime metrics endpoint。也就是说,checkpoint badge 是入口,它的完整语义(哪些 checkpoint 成功/失败、历史趋势)仍然由既有契约负责。大 DAG 降级策略大 DAG 可能难以阅读,也会带来较高重绘成本。V1 应当降级,而不是强行在每个元素上渲染全部信号。推荐行为:保留拓扑渲染能力;展示最热 vertices 与最阻塞 edges 的摘要健康表;DAG 较大时限制自动 fit 和动画;保持选中节点或边后的 drill-down 能力;不为了 UI 复杂度提高 master 采集频率。设计契约还要求:具体实现 PR 应记录选择的大 DAG 阈值,并将其保持为UI 渲染规则,而不是后端采样规则——这保证后端数据口径不被 UI 性能诉求污染。V1 交付计划与验收要求V1 交付计划共 7 步:保持 Active Master realtime REST endpoint 作为 vertex 和 edge 运行状态来源;根据匹配vertexId的最新 vertex point 渲染节点颜色;根据匹配targetVertexId的最新 edge point 渲染 edge 颜色和宽度;在现有详情抽屉中展示最近 point 序列和原始字段;将 checkpoint 与错误信息作为邻近上下文入口,而不是向 realtime metrics 中嵌入新字段;对过大的 DAG 提供降级视图,列出最热 vertices 与最阻塞 edges;同步更新 REST、Web UI、运维文档的英文与中文说明。实现 PR 至少应覆盖的验收检查:后端 realtime edge 和 vertex 响应测试覆盖当前字段与targetVertexId映射;UI 测试覆盖节点染色、edge 染色、edge 宽度和指标未开启时的行为;大 DAG 降级使用确定性的合成 DAG 做测试;文档明确 realtime window 仅为内存、best effort;V1 不引入新的持久化指标表、文件或写入路径。仓库中已经存在可对照的测试入口:RealtimeMetricsRestIT 验证 REST 端点契约,JobInfoDagStabilityRestIT 验证 Job Detail DAG 的稳定性,ObservabilityConfigTest(测试路径)则覆盖配置解析语义。小结:一张图,一条链路,一套契约运行时执行图的设计核心可以概括为三句话:图只有一份(Zeta 作业 DAG),数据只走一条链路(Active Master 短窗口内存指标 既有 REST),信号只分两级(稳定标识字段 best-effort 诊断字段)。诊断工作流因此非常直接:打开 Job Detail DAG → 看节点主色定位最忙 vertex → 看 edge 颜色/宽度定位最阻塞队列 → 点击 node/edge 进入详情抽屉查看原始 bucket 序列 → 结合邻近的 checkpoint 与错误 badge 判断是否需要进入日志。对使用者而言,这套约束换来的收益是:新增诊断能力的同时,不改变既有指标的生命周期、保留规则和成本模型。相关文档实时可观测性忙碌度与背压Web UIRESTful API V2【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻