SeaTunnel Zeta 引擎 Telemetry 监控接入指南:Prometheus 指标导出、配置与 Grafana 可视化
2026/9/17 11:31:11 网站建设 项目流程

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,注册NodeMetricExportsClusterMetricExportsJobMetricExportsJobThreadPoolStatusExportsReportMetricsOperationExportsRequestSlotOperationExportsEngineStateStoreMetricExportsEngineStateStoreLogicalMetricExports等 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.metricTelemetryMetricConfig类型,承载上述开关;
  • 同级的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 提供,反映集群整体与单节点运行状态:

MetricNameTypeLabelsDESCRIPTION
cluster_infoGaugehazelcastVersion,Hazelcast 版本;master,SeaTunnel master 地址集群信息
cluster_timeGaugehazelcastVersion,Hazelcast 版本集群启动时间
node_countGauge-集群节点总数
node_stateGaugeaddress,服务实例地址,例如127.0.0.1:5801SeaTunnel 节点是否在线
hazelcast_executor_executedCountGaugetype,执行器类型:asyncclientclientBlockingclientQueryiooffloadablescheduledsystem集群节点上 Hazelcast 执行器的已执行任务数
hazelcast_executor_isShutdownGaugetype,同上执行器是否已关闭
hazelcast_executor_isTerminatedGaugetype,同上执行器是否已终止
hazelcast_executor_maxPoolSizeGaugetype,同上执行器最大线程池大小
hazelcast_executor_poolSizeGaugetype,同上执行器当前线程池大小
hazelcast_executor_queueRemainingCapacityGaugetype,同上执行器队列剩余容量
hazelcast_executor_queueSizeGaugetype,同上执行器队列大小
hazelcast_partition_partitionCountGauge-集群节点的分区总数
hazelcast_partition_activePartitionGauge-集群节点的活跃分区数
hazelcast_partition_isClusterSafeGauge-分区集群是否安全
hazelcast_partition_isLocalMemberSafeGauge-分区本地成员是否安全

从源码看,cluster_info在导出时会快照一次 master 地址以避免 master 选举期间的竞态,若 master 地址暂不可用则跳过该条指标;node_count直接取集群成员列表大小(ClusterMetricExports.java#L62-L95)。hazelcast_executor_*系列则逐一读取asyncclientclientBlockingclientQueryiooffloadablescheduledsystem八类执行器的 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做聚合。

MetricNameTypeLabelsDESCRIPTION
engine_state_store_local_owned_entriesGaugeaddress,实例地址;store,状态存储名;backend,状态存储后端该节点上状态存储的本地自有条目数
engine_state_store_local_backup_entriesGaugeaddressstorebackend该节点上状态存储的本地备份条目数
engine_state_store_local_heap_cost_bytesGaugeaddressstorebackend后端暴露时的本地堆占用字节数
# 按状态存储聚合总条目数 sum by (cluster, store, backend) (engine_state_store_local_owned_entries)

3.3 引擎状态存储逻辑指标(Engine State Store Logical Metrics)

逻辑指标暴露的是面向业务语义的计数,仅由活跃 master 导出(worker 节点抓不到)。指标名保持后端无关(backend-neutral),但当前实现仍统一带backend="hazelcast"标签。

MetricNameTypeLabelsDESCRIPTION
engine_state_store_running_job_metrics_task_contextsGaugebackendengine_runningJobMetrics中当前存储的任务指标上下文总数
engine_state_store_running_job_metrics_active_partition_keysGaugebackendengine_runningJobMetrics中当前活跃的顶层分区桶数
engine_state_store_checkpoint_monitor_jobsGaugebackendengine_checkpoint_monitor当前追踪的作业数
engine_state_store_checkpoint_monitor_in_progress_checkpointsGaugebackendengine_checkpoint_monitor中进行中的 checkpoint 数
engine_state_store_checkpoint_monitor_retained_history_entriesGaugebackendengine_checkpoint_monitor中保留的 checkpoint 历史条目数
engine_state_store_finished_job_recordsGaugestore,finished job 存储名;backendfinished job 存储中的当前记录数
engine_state_store_finished_job_cleanup_totalCounterstorebackend从 finished job 存储过期清理事件累计数
engine_state_store_connector_jar_tracked_jarsGaugebackendengine_connectorJarRefCounters中当前追踪的 connector jar 数
engine_state_store_connector_jar_total_referencesGaugebackendengine_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 节点端点不会返回这些指标。

MetricNameTypeLabelsDESCRIPTION
job_thread_pool_activeCountGaugeaddresscoordinator 作业执行器缓存线程池活跃线程数
job_thread_pool_corePoolSizeGaugeaddress核心线程池大小
job_thread_pool_maximumPoolSizeGaugeaddress最大线程池大小
job_thread_pool_poolSizeGaugeaddress当前线程池大小
job_thread_pool_queueTaskCountGaugeaddress队列中的任务数
job_thread_pool_completedTask_totalCounteraddress已完成任务数
job_thread_pool_task_totalCounteraddress已提交任务总数
job_thread_pool_rejection_totalCounteraddress被拒绝任务数

实战中,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 状态均不受影响。

MetricNameTypeLabelsDESCRIPTION
report_metrics_operation_totalCounteraddress,worker 实例地址;result,取值successfailureinterruptedworker 发送的ReportMetricsOperation调用总数
report_metrics_operation_last_payload_task_countGaugeaddress最近一次上报负载中包含的任务指标数
report_metrics_operation_last_invocation_latency_msGaugeaddress最近一次上报延迟(毫秒),含本地指标采集与 worker 到 master 的调用耗时
report_metrics_operation_max_invocation_latency_msGaugeaddress自 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 标签。

MetricNameTypeLabelsDESCRIPTION
request_slot_operation_totalCounteraddress,master 实例地址;result,取值successno_slotfailuremaster 向 worker 发送的RequestSlotOperation调用总数
request_slot_operation_last_invocation_latency_msGaugeaddress最近一次 master 侧调用延迟(毫秒)
request_slot_operation_max_invocation_latency_msGaugeaddress自 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"}

MetricNameTypeLabelsDESCRIPTION
job_countGaugetype,作业类型:canceledcancellingcreatedfailedfailingfinishedrunningscheduledSeaTunnel 集群全部作业计数
# 集群运行中作业数 job_count{type="running"} # 集群失败作业数(可配置告警阈值) job_count{type="failed"}

3.8 JVM 指标

JVM 指标由 Prometheus Java 客户端的DefaultExports(标准 hotspot collectors)提供,在 ExportsInstanceInitializer 中通过DefaultExports.initialize()一次性初始化,覆盖线程、类加载、内存池、GC、进程等维度:

MetricNameTypeLabelsDESCRIPTION
jvm_threads_current / daemon / peak / started_total / deadlocked / deadlocked_monitorGauge / Counter- /state当前、守护、峰值、累计启动线程数;死锁检测(含仅等待对象监视器的死锁周期)
jvm_threads_stateGaugestateNEWTERMINATEDRUNNABLEBLOCKEDWAITINGTIMED_WAITINGUNKNOWN按线程状态计数
jvm_classes_currently_loaded / loaded_total / unloaded_totalGauge / Counter-当前已加载类数、累计加载/卸载类数
jvm_memory_pool_allocated_bytes_totalCounterpoolCode CachePS Eden SpacePS Old GenPS Survivor SpaceCompressed Class SpaceMetaspace内存池累计分配字节(仅在 GC 后更新)
jvm_gc_collection_seconds_count / sumSummarygcPS ScavengePS MarkSweep各垃圾收集器耗时(秒)
jvm_infoGaugeruntimevendorversionJVM 版本信息
process_cpu_seconds_total / start_time_seconds / open_fds / max_fdsCounter / Gauge-进程 CPU 时间、启动时间戳、打开/最大文件描述符
jvm_memory_objects_pending_finalizationGauge-等待 finalizer 队列的对象数
jvm_memory_bytes_used / committed / max / initGaugeareaheapnoheapJVM 内存区使用/已提交/最大/初始字节
jvm_memory_pool_bytes_used / committed / max / initGaugepool,见上内存池字节指标
jvm_memory_pool_allocated_bytes_createdGaugepool内存池累计分配字节(仅在 GC 后更新)
jvm_memory_pool_collection_used / committed / max / init_bytesGaugepool(GC 场景仅PS Eden SpacePS Old GenPS Survivor Space最近一次 GC 后的内存池字节指标
jvm_buffer_pool_used_bytes / capacity_bytes / used_buffersGaugepooldirectmapped直接/映射缓冲区池使用情况

常用查询:

# 各节点堆内存使用率(按集群) 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 文档。安装并启动后:

  1. 登录 Grafana,进入Configuration → Data Sources → Add data source
  2. 选择Prometheus,在 URL 中填入 Prometheus 服务地址(如http://localhost:9090);
  3. 点击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 引擎监控链路已经打通:

  1. seatunnel.yaml中开启seatunnel.engine.telemetry.metric.enabled: true(配置解析见 YamlSeaTunnelDomConfigProcessor);
  2. 节点启动时由 ExportsInstanceInitializer 注册全部 Collector,MetricsServlet 在5801端口提供/hazelcast/rest/instance/metrics/hazelcast/rest/instance/openmetrics两个抓取端点;
  3. Prometheus 按scrape_interval周期性抓取,指标统一携带cluster标签(值来自hazelcast.cluster-name,见 AbstractCollector);
  4. Grafana 通过 Prometheus 数据源可视化,并基于集群节点、作业状态、checkpoint 积压、线程池拒绝数等关键指标配置告警。

使用建议:区分"每节点上报"与"仅 master 导出"两类指标——node_stateengine_state_store_local_*jvm_*等在任意节点抓取均有意义;而job_countjob_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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询