SeaTunnel Zeta 引擎 Telemetry 监控接入指南:Prometheus 指标导出、配置与 Grafana 可视化
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文基于 SeaTunnel 开源仓库的 Zeta 引擎 Telemetry 文档 展开,系统讲解如何通过 Prometheus 协议导出 SeaTunnel Zeta 引擎的集群指标,并接入 Prometheus 与 Grafana 构建可观测告警体系。读完本文,你将掌握seatunnel.yaml中 telemetry 配置的完整写法、各指标类别的含义与典型 PromQL 查询、以及从抓取配置到仪表盘导入的端到端实战流程。
一、Telemetry 概述:为什么需要指标导出
SeaTunnel Zeta 引擎(seatunnel-engine模块)是基于 Hazelcast 构建的分布式计算引擎。在多节点集群环境下,仅靠日志难以回答"当前集群有几个节点在运行""哪个节点的执行线程池已经饱和""checkpoint 积压了多少"这类运维问题。Telemetry 机制通过 Prometheus 格式的指标暴露接口,让 SeaTunnel 集群可以被 Prometheus、Grafana 等标准监控平台无缝纳管,从而显著提升集群的监控与告警能力。
从实现层面看,指标导出由 ExportsInstanceInitializer 完成:每个 Hazelcast 节点启动时会初始化一个 PrometheusCollectorRegistry,注册NodeMetricExports、ClusterMetricExports、JobMetricExports、JobThreadPoolStatusExports、ReportMetricsOperationExports、RequestSlotOperationExports、EngineStateStoreMetricExports、EngineStateStoreLogicalMetricExports等 Collector,并调用DefaultExports.initialize()注入 JVM 热点指标。随后 MetricsServlet 会根据请求路径选择 Prometheus 文本格式(/metrics)或 OpenMetrics 格式(/openmetrics)序列化输出。
配置开关:seatunnel.yaml
Telemetry 配置位于集群配置文件seatunnel.yaml(参考 config/seatunnel.yaml 与仓库内置模板 seatunnel-engine/seatunnel-engine-common/src/main/resources/seatunnel.yaml),默认情况下指标导出是关闭的:
seatunnel: engine: telemetry: metric: enabled: false # Whether open metrics export开启指标导出只需将enabled置为true:
seatunnel: engine: telemetry: metric: enabled: true # Whether open metrics export配置项的解析入口位于 YamlSeaTunnelDomConfigProcessor 的parseTelemetryConfig/parseTelemetryMetricConfig方法,对应选项定义在 ServerConfigOptions 中:
telemetry.metric.enabled:布尔类型,默认false,控制是否开放指标导出;telemetry.metric:TelemetryMetricConfig类型,承载上述开关;- 同级的
telemetry.logs.scheduled-deletion-enable默认true,控制历史作业日志的定时清理(与指标无关,但同属 telemetry 配置域)。
指标抓取地址
开启后,通过 REST 服务暴露两条抓取路径(端点常量定义见 RestConstant):
- Prometheus 文本格式:
http://{instanceHost}:5801/hazelcast/rest/instance/metrics - OpenMetrics 文本格式:
http://{instanceHost}:5801/hazelcast/rest/instance/openmetrics
其中5801是 Hazelcast 成员通信端口(见 config/hazelcast.yaml 中的port配置),{instanceHost}替换为 SeaTunnel 集群任一节点的地址即可。
二、通用标签约定
所有指标都带有一个相同的标签名cluster,其值取自hazelcast.cluster-name配置(见 config/hazelcast.yaml 中的cluster-name: seatunnel)。在 AbstractCollector 中,clusterLabelNames()会把cluster作为首个标签名,labelValues()会把集群名作为首个标签值。这意味着:
- 同一套 Prometheus 中可以安全地纳管多个 SeaTunnel 集群,靠
cluster标签区分; - 编写告警规则或面板时,可按
cluster分组、过滤,避免集群间数据互相污染。
此外,凡是带address标签的指标,其取值格式为IP:port,例如127.0.0.1:5801。
三、指标类别详解
当前可用的指标分为以下类别。每类指标均对应一个独立的 Collector 实现,下面结合源码说明其数据来源与用途。
3.1 节点指标(Node Metrics)
由 NodeMetricExports 与 ClusterMetricExports 提供,反映集群整体与单节点运行状态:
| MetricName | Type | Labels | DESCRIPTION |
|---|---|---|---|
| cluster_info | Gauge | hazelcastVersion,Hazelcast 版本;master,SeaTunnel master 地址 | 集群信息 |
| cluster_time | Gauge | hazelcastVersion,Hazelcast 版本 | 集群启动时间 |
| node_count | Gauge | - | 集群节点总数 |
| node_state | Gauge | address,服务实例地址,例如127.0.0.1:5801 | SeaTunnel 节点是否在线 |
| hazelcast_executor_executedCount | Gauge | type,执行器类型:asyncclientclientBlockingclientQueryiooffloadablescheduledsystem | 集群节点上 Hazelcast 执行器的已执行任务数 |
| hazelcast_executor_isShutdown | Gauge | type,同上 | 执行器是否已关闭 |
| hazelcast_executor_isTerminated | Gauge | type,同上 | 执行器是否已终止 |
| hazelcast_executor_maxPoolSize | Gauge | type,同上 | 执行器最大线程池大小 |
| hazelcast_executor_poolSize | Gauge | type,同上 | 执行器当前线程池大小 |
| hazelcast_executor_queueRemainingCapacity | Gauge | type,同上 | 执行器队列剩余容量 |
| hazelcast_executor_queueSize | Gauge | type,同上 | 执行器队列大小 |
| hazelcast_partition_partitionCount | Gauge | - | 集群节点的分区总数 |
| hazelcast_partition_activePartition | Gauge | - | 集群节点的活跃分区数 |
| hazelcast_partition_isClusterSafe | Gauge | - | 分区集群是否安全 |
| hazelcast_partition_isLocalMemberSafe | Gauge | - | 分区本地成员是否安全 |
从源码看,cluster_info在导出时会快照一次 master 地址以避免 master 选举期间的竞态,若 master 地址暂不可用则跳过该条指标;node_count直接取集群成员列表大小(ClusterMetricExports.java#L62-L95)。hazelcast_executor_*系列则逐一读取async、client、clientBlocking、clientQuery、io、offloadable、scheduled、system八类执行器的 JMX MBean(NodeMetricExports.java#L48-L132),适合用来判断哪类执行器出现队列积压或线程池饱和。
典型查询示例:
# 集群节点总数 node_count{cluster="seatunnel"} # 按执行器类型查看队列积压 hazelcast_executor_queueSize{cluster="seatunnel"}3.2 引擎状态存储指标(Engine State Store Metrics)
Zeta 引擎的状态存储(checkpoint、running job metrics、finished job 记录等)当前底层基于 Hazelcast IMap,因此backend标签取值恒为hazelcast。此类指标反映状态存储的基础规模与本机资源占用,由每个节点各自上报,带address标签。若要监控某个状态存储的全局条目总数,需要在 Prometheus 中对engine_state_store_local_owned_entries做聚合。
| MetricName | Type | Labels | DESCRIPTION |
|---|---|---|---|
| engine_state_store_local_owned_entries | Gauge | address,实例地址;store,状态存储名;backend,状态存储后端 | 该节点上状态存储的本地自有条目数 |
| engine_state_store_local_backup_entries | Gauge | address;store;backend | 该节点上状态存储的本地备份条目数 |
| engine_state_store_local_heap_cost_bytes | Gauge | address;store;backend | 后端暴露时的本地堆占用字节数 |
# 按状态存储聚合总条目数 sum by (cluster, store, backend) (engine_state_store_local_owned_entries)3.3 引擎状态存储逻辑指标(Engine State Store Logical Metrics)
逻辑指标暴露的是面向业务语义的计数,仅由活跃 master 导出(worker 节点抓不到)。指标名保持后端无关(backend-neutral),但当前实现仍统一带backend="hazelcast"标签。
| MetricName | Type | Labels | DESCRIPTION |
|---|---|---|---|
| engine_state_store_running_job_metrics_task_contexts | Gauge | backend | engine_runningJobMetrics中当前存储的任务指标上下文总数 |
| engine_state_store_running_job_metrics_active_partition_keys | Gauge | backend | engine_runningJobMetrics中当前活跃的顶层分区桶数 |
| engine_state_store_checkpoint_monitor_jobs | Gauge | backend | engine_checkpoint_monitor当前追踪的作业数 |
| engine_state_store_checkpoint_monitor_in_progress_checkpoints | Gauge | backend | engine_checkpoint_monitor中进行中的 checkpoint 数 |
| engine_state_store_checkpoint_monitor_retained_history_entries | Gauge | backend | engine_checkpoint_monitor中保留的 checkpoint 历史条目数 |
| engine_state_store_finished_job_records | Gauge | store,finished job 存储名;backend | finished job 存储中的当前记录数 |
| engine_state_store_finished_job_cleanup_total | Counter | store;backend | 从 finished job 存储过期清理事件累计数 |
| engine_state_store_connector_jar_tracked_jars | Gauge | backend | engine_connectorJarRefCounters中当前追踪的 connector jar 数 |
| engine_state_store_connector_jar_total_references | Gauge | backend | engine_connectorJarRefCounters中 connector jar 引用计数之和 |
逻辑指标与上面的本地存储指标互补,选择原则是:
- 需要理解各节点上的数据分布与内存占用时,用
engine_state_store_local_*; - 需要理解引擎语义(如 checkpoint 积压、finished job 保留量、connector jar 复用情况)时,用
engine_state_store_*逻辑指标。
官方给出的 PromQL 示例:
# Total state store entries by store sum by (cluster, store, backend) (engine_state_store_local_owned_entries) # Current checkpoint backlog engine_state_store_checkpoint_monitor_in_progress_checkpoints{backend="hazelcast"} # Finished-job cleanup growth over the last 15 minutes increase(engine_state_store_finished_job_cleanup_total{backend="hazelcast"}[15m]) # Connector jar reference pressure engine_state_store_connector_jar_total_references{backend="hazelcast"}对应的 Grafana 面板建议:
State Store Total Entries
sum by (store) (engine_state_store_local_owned_entries{backend="hazelcast"})Checkpoint In-Progress Count
engine_state_store_checkpoint_monitor_in_progress_checkpoints{backend="hazelcast"}Finished Job Cleanup Rate
sum by (store) (rate(engine_state_store_finished_job_cleanup_total{backend="hazelcast"}[5m]))Connector Jar Reference Count
engine_state_store_connector_jar_total_references{backend="hazelcast"}3.4 线程池状态(Thread Pool Status)
由 JobThreadPoolStatusExports 提供,反映 coordinator 作业执行器(缓存线程池)的健康度。仅由活跃 master 导出,抓取 worker 节点端点不会返回这些指标。
| MetricName | Type | Labels | DESCRIPTION |
|---|---|---|---|
| job_thread_pool_activeCount | Gauge | address | coordinator 作业执行器缓存线程池活跃线程数 |
| job_thread_pool_corePoolSize | Gauge | address | 核心线程池大小 |
| job_thread_pool_maximumPoolSize | Gauge | address | 最大线程池大小 |
| job_thread_pool_poolSize | Gauge | address | 当前线程池大小 |
| job_thread_pool_queueTaskCount | Gauge | address | 队列中的任务数 |
| job_thread_pool_completedTask_total | Counter | address | 已完成任务数 |
| job_thread_pool_task_total | Counter | address | 已提交任务总数 |
| job_thread_pool_rejection_total | Counter | address | 被拒绝任务数 |
实战中,job_thread_pool_rejection_total一旦持续增长,说明 coordinator 的执行线程池已经无法承接任务提交,通常需要检查作业并发度与资源分配。
3.5 指标上报操作(Report Metrics Operation)
ReportMetricsOperation是 worker 向 master 上报任务指标快照的 RPC 操作。每个指标桶(bucket)的写入与删除最多尝试 10 次条件更新,若竞争持续存在,操作会以Failed to update metrics partition ... after 10 concurrent modifications失败。失败的 worker 上报会被记录日志并计入report_metrics_operation_total{result="failure"},后续的定时上报可在任务上下文保留期间重试;指标删除失败时,待清理的 pipeline 记录会被保留。该上限约束的是冲突重试次数而非网络延迟——单次 Hazelcast 调用仍使用各自配置的超时时间。指标格式以及 checkpoint/savepoint 状态均不受影响。
| MetricName | Type | Labels | DESCRIPTION |
|---|---|---|---|
| report_metrics_operation_total | Counter | address,worker 实例地址;result,取值successfailureinterrupted | worker 发送的ReportMetricsOperation调用总数 |
| report_metrics_operation_last_payload_task_count | Gauge | address | 最近一次上报负载中包含的任务指标数 |
| report_metrics_operation_last_invocation_latency_ms | Gauge | address | 最近一次上报延迟(毫秒),含本地指标采集与 worker 到 master 的调用耗时 |
| report_metrics_operation_max_invocation_latency_ms | Gauge | address | 自 worker 启动以来观测到的最大上报延迟(毫秒) |
# 上报失败率(近 5 分钟) sum by (cluster) (rate(report_metrics_operation_total{result="failure"}[5m])) / sum by (cluster) (rate(report_metrics_operation_total[5m]))3.6 请求槽位操作(Request Slot Operation)
RequestSlotOperation是 master 侧的槽位分配 RPC 路径:当作业需要执行资源时,活跃 master 会选择候选 worker 并发送请求以预留槽位。这类指标用于帮助运维区分槽位分配 RPC 缓慢/失败与普遍性资源不足。仅由活跃 master 导出,且是聚合信号,不包含 job、worker、slot 或 resource-profile 标签。
| MetricName | Type | Labels | DESCRIPTION |
|---|---|---|---|
| request_slot_operation_total | Counter | address,master 实例地址;result,取值successno_slotfailure | master 向 worker 发送的RequestSlotOperation调用总数 |
| request_slot_operation_last_invocation_latency_ms | Gauge | address | 最近一次 master 侧调用延迟(毫秒) |
| request_slot_operation_max_invocation_latency_ms | Gauge | address | 自 master 启动以来观测到的最大调用延迟(毫秒) |
result标签含义如下:
success:worker 返回了分配的槽位;no_slot:请求到达 worker 并正常完成,但 worker 没有返回合适的槽位。该值持续上升可能意味着 master 的 worker 资源视图与 worker 当前槽位状态出现偏差,或者并发分配在预检查与请求执行之间消耗了槽位;failure:master 到 worker 的调用失败或操作异常完成。
# 槽位分配无槽位率 sum by (cluster) (rate(request_slot_operation_total{result="no_slot"}[5m]))3.7 作业信息明细(Job Info Detail)
由 JobMetricExports 提供,仅由活跃 master 导出。它只按状态报告聚合计数,没有 per-job 标签,因此不能针对某个具体作业 ID 或名称告警,只能用于集群维度总量告警,例如job_count{type="failed"}。
| MetricName | Type | Labels | DESCRIPTION |
|---|---|---|---|
| job_count | Gauge | type,作业类型:canceledcancellingcreatedfailedfailingfinishedrunningscheduled | SeaTunnel 集群全部作业计数 |
# 集群运行中作业数 job_count{type="running"} # 集群失败作业数(可配置告警阈值) job_count{type="failed"}3.8 JVM 指标
JVM 指标由 Prometheus Java 客户端的DefaultExports(标准 hotspot collectors)提供,在 ExportsInstanceInitializer 中通过DefaultExports.initialize()一次性初始化,覆盖线程、类加载、内存池、GC、进程等维度:
| MetricName | Type | Labels | DESCRIPTION |
|---|---|---|---|
| jvm_threads_current / daemon / peak / started_total / deadlocked / deadlocked_monitor | Gauge / Counter | - /state | 当前、守护、峰值、累计启动线程数;死锁检测(含仅等待对象监视器的死锁周期) |
| jvm_threads_state | Gauge | state:NEWTERMINATEDRUNNABLEBLOCKEDWAITINGTIMED_WAITINGUNKNOWN | 按线程状态计数 |
| jvm_classes_currently_loaded / loaded_total / unloaded_total | Gauge / Counter | - | 当前已加载类数、累计加载/卸载类数 |
| jvm_memory_pool_allocated_bytes_total | Counter | pool:Code CachePS Eden SpacePS Old GenPS Survivor SpaceCompressed Class SpaceMetaspace | 内存池累计分配字节(仅在 GC 后更新) |
| jvm_gc_collection_seconds_count / sum | Summary | gc:PS ScavengePS MarkSweep | 各垃圾收集器耗时(秒) |
| jvm_info | Gauge | runtime、vendor、version | JVM 版本信息 |
| process_cpu_seconds_total / start_time_seconds / open_fds / max_fds | Counter / Gauge | - | 进程 CPU 时间、启动时间戳、打开/最大文件描述符 |
| jvm_memory_objects_pending_finalization | Gauge | - | 等待 finalizer 队列的对象数 |
| jvm_memory_bytes_used / committed / max / init | Gauge | area:heapnoheap | JVM 内存区使用/已提交/最大/初始字节 |
| jvm_memory_pool_bytes_used / committed / max / init | Gauge | pool,见上 | 内存池字节指标 |
| jvm_memory_pool_allocated_bytes_created | Gauge | pool | 内存池累计分配字节(仅在 GC 后更新) |
| jvm_memory_pool_collection_used / committed / max / init_bytes | Gauge | pool(GC 场景仅PS Eden SpacePS Old GenPS Survivor Space) | 最近一次 GC 后的内存池字节指标 |
| jvm_buffer_pool_used_bytes / capacity_bytes / used_buffers | Gauge | pool:directmapped | 直接/映射缓冲区池使用情况 |
常用查询:
# 各节点堆内存使用率(按集群) sum by (cluster) (jvm_memory_bytes_used{area="heap"}) / sum by (cluster) (jvm_memory_bytes_max{area="heap"}) # 老年代 GC 耗时增长 increase(jvm_gc_collection_seconds_sum{gc="PS MarkSweep"}[5m])四、Prometheus 与 Grafana 集群监控实战
4.1 安装 Prometheus
Prometheus 服务端的安装方式请参考 Prometheus 官方 Installation 文档(Prometheus 官网提供二进制、容器化与包管理器多种安装路径,本文不再赘述)。安装完成后确认服务可访问,再进行抓取配置。
4.2 配置 Prometheus 抓取 SeaTunnel
编辑 Prometheus 配置文件(常见路径为/etc/prometheus/prometheus.yaml),添加 SeaTunnel 实例的指标抓取任务。下面配置同时展示了两种抓取路径(/metrics与/openmetrics,二选一或按需使用):
global: # How frequently to scrape targets from this job. scrape_interval: 15s scrape_configs: # The job name assigned to scraped metrics by default. - job_name: 'seatunnel' scrape_interval: 5s # Metrics export path (Prometheus text format) metrics_path: /hazelcast/rest/instance/metrics # List of labeled statically configured targets for this job. static_configs: # The targets specified by the static config. - targets: [ 'localhost:5801' ] # Labels assigned to all metrics scraped from the targets. # labels: [<labelName>:<labelValue>] # OpenMetrics 格式抓取(如需) - job_name: 'seatunnel-openmetrics' scrape_interval: 5s metrics_path: /hazelcast/rest/instance/openmetrics static_configs: - targets: [ 'localhost:5801' ]要点说明:
metrics_path必须为/hazelcast/rest/instance/metrics(或/openmetrics),与 RestConstant 中的端点路径一致;targets中的端口为 Hazelcast 成员端口5801(config/hazelcast.yaml 中hazelcast.network.port.port),多节点集群可在此列出全部节点;- 可通过
labels给不同集群打标,配合cluster标签进一步区分多套环境。
修改配置后重启 Prometheus 或触发配置热加载,然后在 Prometheus Targets 页面确认seatunneljob 处于 UP 状态。
4.3 安装 Grafana 并接入数据源
Grafana 服务端的安装方式同样请参考 Grafana 官方 Installation 文档。安装并启动后:
- 登录 Grafana,进入Configuration → Data Sources → Add data source;
- 选择Prometheus,在 URL 中填入 Prometheus 服务地址(如
http://localhost:9090); - 点击Save & Test确认连通。
4.4 监控仪表盘
完成数据源接入后:
- 在 Grafana 中选择Import,导入
Seatunnel Cluster监控仪表盘 JSON(官方文档提供的 dashboard JSON 可从 SeaTunnel 发行包或社区渠道获取)。
仪表盘的效果图如下:
4.5 常见告警规则示例
结合上文指标,可以配置如下告警规则(写入 Prometheus rules 文件):
groups: - name: seatunnel-alerts rules: # 节点失联 - alert: SeaTunnelNodeDown expr: node_state{cluster="seatunnel"} == 0 for: 2m labels: severity: critical annotations: summary: "SeaTunnel node {{ $labels.address }} is down" # 作业失败 - alert: SeaTunnelJobFailed expr: job_count{type="failed"} > 0 for: 1m labels: severity: warning annotations: summary: "SeaTunnel cluster has failed jobs" # checkpoint 积压 - alert: CheckpointBacklogHigh expr: engine_state_store_checkpoint_monitor_in_progress_checkpoints{backend="hazelcast"} > 50 for: 5m labels: severity: warning annotations: summary: "Checkpoint backlog is high on cluster {{ $labels.cluster }}" # 线程池拒绝任务 - alert: JobThreadPoolRejection expr: increase(job_thread_pool_rejection_total[5m]) > 0 for: 5m labels: severity: warning五、指标抓取与监控链路小结
至此,一个完整的 SeaTunnel Zeta 引擎监控链路已经打通:
- 在
seatunnel.yaml中开启seatunnel.engine.telemetry.metric.enabled: true(配置解析见 YamlSeaTunnelDomConfigProcessor); - 节点启动时由 ExportsInstanceInitializer 注册全部 Collector,MetricsServlet 在
5801端口提供/hazelcast/rest/instance/metrics与/hazelcast/rest/instance/openmetrics两个抓取端点; - Prometheus 按
scrape_interval周期性抓取,指标统一携带cluster标签(值来自hazelcast.cluster-name,见 AbstractCollector); - Grafana 通过 Prometheus 数据源可视化,并基于集群节点、作业状态、checkpoint 积压、线程池拒绝数等关键指标配置告警。
使用建议:区分"每节点上报"与"仅 master 导出"两类指标——node_state、engine_state_store_local_*、jvm_*等在任意节点抓取均有意义;而job_count、job_thread_pool_*、request_slot_operation_*、engine_state_store_*(逻辑指标)只在活跃 master 端返回,多节点环境请确保 Prometheus 抓取到 master 节点,或在查询时注意聚合口径,避免因抓取目标不同导致数据缺失误判。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考