☰
Hazelcast Scheduled 与 Durable Executor 服务指标统计的设计与实现
2026/10/8 7:55:53 网站建设 项目流程

本设计文档对应仓库路径:docs/design/executors/01-executor-service-stats.md,自 Hazelcast 4.1 起生效。该文档描述了一次对 Hazelcast 分布式执行框架的观测性补全:在传统 Executor Service 之外,为 Scheduled Executor Service 与 Durable Executor Service 补齐了同规格的任务级指标统计,并通过 metrics 子系统将数据暴露给 Management Center。读完本文,你将掌握三种 Executor 在统计能力上的差异、六类指标的定义与生命周期、metric 名称与前缀约定、如何通过ScheduledExecutorConfig/DurableExecutorConfig开关统计,以及底层源码中统计计数与指标导出的完整调用链。

背景:三种 Executor Service,只有一种有统计

Hazelcast 通过公共 API 暴露了三种分布式执行器实现(详见 executor 包、scheduledexecutor 包、durableexecutor 包):

Executor 类型核心特征统计能力(本设计之前)
Executor Service(IExecutorService)通用分布式任务提交,任务绑定分区执行有,已有完整统计
Scheduled Executor Service(IScheduledExecutorService)支持单次/固定频率调度、任务可持久化与迁移无
Durable Executor Service(IDurableExecutorService)任务写入 ring buffer、具备持久性与至少一次执行语义无

在 4.1 之前,三种实现中只有普通 Executor Service 维护运行统计;Scheduled 与 Durable 两种均未采集任何指标。本设计的目标,就是把 Executor Service 已有的同一套统计模型扩展到另外两种执行器,并让它们统一通过 Management Center 可见、可监控。

设计:统一复用LocalExecutorStats统计模型

统计模型接口

整个统计能力建立在LocalExecutorStats接口之上,其定义位于 hazelcast/src/main/java/com/hazelcast/executor/LocalExecutorStats.java:

public interface LocalExecutorStats extends LocalInstanceStats { /** Returns the number of pending operations on the executor service. */ long getPendingTaskCount(); /** Returns the number of started operations on the executor service. */ long getStartedTaskCount(); /** Returns the number of completed operations on the executor service. */ long getCompletedTaskCount(); /** Returns the number of cancelled operations on the executor service. */ long getCancelledTaskCount(); /** Returns the total start latency of operations started. */ long getTotalStartLatency(); /** Returns the total execution time of operations finished. */ long getTotalExecutionLatency(); }

接口暴露六个核心指标(外加从LocalInstanceStats继承的getCreationTime()):

  • pendingTaskCount:当前排队中、尚未开始执行的任务数;
  • startedTaskCount:已经开始执行的任务数;
  • completedTaskCount:已完成的任务数;
  • cancelledTaskCount:被取消的任务数;
  • totalStartLatency:所有任务从提交到开始执行的累计等待时间(毫秒);
  • totalExecutionLatency:所有已完成任务的实际执行时长累计(毫秒)。

其实现类为LocalExecutorStatsImpl(hazelcast/src/main/java/com/hazelcast/internal/monitor/impl/LocalExecutorStatsImpl.java)。从源码可以看到两个值得注意的实现细节:

  1. 无锁并发计数:所有计数器字段(pending、started、completed、cancelled、totalStartLatency、totalExecutionTime)都是volatile long,通过AtomicLongFieldUpdater做原子自增/累加,避免在热路径上引入锁竞争;
  2. Probe 注解驱动的指标导出:每个字段上都标注了@Probe(name = EXECUTOR_METRIC_..., unit = MS),其中creationTime、totalStartLatency、totalExecutionTime的单位是毫秒(MS),四个计数指标不带单位,由 metrics 子系统按count处理。

计数器的状态迁移逻辑也全部集中在该实现类中:

public void startPending() { PENDING.incrementAndGet(this); } public void startExecution(long elapsed) { TOTAL_START_LATENCY.addAndGet(this, elapsed); STARTED.incrementAndGet(this); PENDING.decrementAndGet(this); } public void finishExecution(long elapsed) { TOTAL_EXECUTION_TIME.addAndGet(this, elapsed); COMPLETED.incrementAndGet(this); } public void rejectExecution() { PENDING.decrementAndGet(this); } public void cancelExecution() { CANCELLED.incrementAndGet(this); }

也就是说:任务提交时pending +1;任务真正开跑时pending -1、started +1,并把「提交到开始」的耗时累加进totalStartLatency;任务结束时completed +1,把执行耗时累加进totalExecutionTime;被拒绝执行的任务从 pending 中回退;被取消的任务单独计入cancelled。

共享的统计容器ExecutorStats

三种 Executor 服务并不是各自实现一套统计,而是共用同一个ExecutorStats容器类(hazelcast/src/main/java/com/hazelcast/map/impl/ExecutorStats.java)。该类源码注释明确写道:“Scheduled、durable 和 executor 服务的实现都使用这个类来收集统计”,其职责包括:

  • 每个 Executor Service 实例持有一个ExecutorStats对象;
  • 内部用ConcurrentHashMap<String, LocalExecutorStatsImpl>按执行器名称(executorName)维护每个命名执行器的统计;
  • 提供startExecution / finishExecution / startPending / rejectExecution / cancelExecution转发方法,以及getStatsMap()、clear()、removeStats(executorName)等管理操作。

这意味着「按执行器名称隔离统计」对三种 Executor 是一致的行为:同一个集群里不同名字的 Executor 各有各的计数器,销毁某个执行器时通过removeStats清理对应条目。

指标采集的生命周期:任务提交即开始

文档强调“On submit of each task, we start to collect these statistics”(每个任务提交时就开始采集统计)。以 Scheduled Executor 为例,其执行入口 hazelcast/src/main/java/com/hazelcast/scheduledexecutor/impl/TaskRunner.java 完整演示了这一生命周期:

  • 构造TaskRunner时(任务入队):executorStats.startPending(name);
  • call()开跑前:startExecution(name, start - creationTime)—— 记录排队等待耗时并切换状态;
  • finally中:finishExecution(name, Clock.currentTimeMillis() - start);
  • 对于AT_FIXED_RATE固定频率任务:每轮执行完后会重新调用startPending(name)并把creationTime重置,使周期任务的“下一轮排队”也被计入 pending,从而让指标在调度循环中保持连续;
  • 取消路径(ScheduledExecutorContainer的cancelExecution)与拒绝路径(rejectExecution)分别维护对应计数。

Durable Executor 侧的逻辑与之对称,见 hazelcast/src/main/java/com/hazelcast/durableexecutor/impl/DurableExecutorContainer.java:

  • 内部类TaskProcessor构造时startPending;
  • run()开始时startExecution(name, start - creationTime);
  • 执行结束后(未被取消的前提下)finishExecution(name, Clock.currentTimeMillis() - start);
  • 在 ring buffer 提交或分区迁移后重放任务时,若RejectedExecutionException发生,则调用rejectExecution(name)回退 pending 计数。

指标导出:DynamicMetricsProvider与 Metrics 子系统

统计数据采集完成后,通过 Hazelcast 的 metrics 子系统对外发布,供 Management Center 拉取。核心机制是两个服务类实现DynamicMetricsProvider接口,并在init()阶段向 metrics 注册表注册自己:

  • DistributedScheduledExecutorService(SERVICE_NAME = "hz:impl:scheduledExecutorService")
  • DistributedDurableExecutorService(SERVICE_NAME = "hz:impl:durableExecutorService")

两者的init()逻辑几乎一致,且受一个集群属性控制:

boolean dsMetricsEnabled = nodeEngine.getProperties().getBoolean(ClusterProperty.METRICS_DATASTRUCTURES); if (dsMetricsEnabled) { nodeEngine.getMetricsRegistry().registerDynamicMetricsProvider(this); }

即:数据结构的 metrics 开关(METRICS_DATASTRUCTURES)默认开启,只有当该属性被显式关闭时,这两个执行器服务才不会注册动态指标提供者。

导出动作本身由provideDynamicMetrics完成,它委托给ProviderHelper.provide(...)(hazelcast/src/main/java/com/hazelcast/internal/metrics/impl/ProviderHelper.java)将getStats()返回的Map<String, LocalExecutorStatsImpl>逐条写入指标上下文。两个服务各自的getStats()都直接返回executorStats.getStatsMap()。

Metrics 前缀约定

文档明确了两种服务使用的指标前缀,对应常量定义于 hazelcast/src/main/java/com/hazelcast/internal/metrics/MetricDescriptorConstants.java:

服务常量前缀值
Scheduled ExecutorSCHEDULED_EXECUTOR_PREFIXscheduledExecutor
Durable ExecutorDURABLE_EXECUTOR_PREFIXdurableExecutor

该文件头部有一段重要注释:这些常量必须在 minor 版本之间保持稳定,修改它们可能破坏 Management Center 或其他指标消费方。因此scheduledExecutor.*与durableExecutor.*系列指标名称是事实上的对外契约。

每个命名执行器的指标都会附带name=<executorName>判别维度(discriminator),以区分同一个 JVM 上多个不同名字的执行器。

示例输出

文档给出了启用统计后典型的指标输出(以scheduledExecutor前缀为例,unit为毫秒或计数):

[name=executorName,unit=ms,metric=scheduledExecutor.creationTime]=1598016899537 [name=executorName,unit=count,metric=scheduledExecutor.pending]=0 [name=executorName,unit=count,metric=scheduledExecutor.started]=1 [name=executorName,unit=count,metric=scheduledExecutor.completed]=0 [name=executorName,unit=count,metric=scheduledExecutor.cancelled]=0 [name=executorName,unit=ms,metric=scheduledExecutor.totalStartLatency]=2 [name=executorName,unit=ms,metric=scheduledExecutor.totalExecutionTime]=0

durableExecutor前缀的指标结构与之一一对应,只是前缀不同。这七项指标名称分别来自MetricDescriptorConstants中的EXECUTOR_METRIC_CREATION_TIME、EXECUTOR_METRIC_PENDING、EXECUTOR_METRIC_STARTED、EXECUTOR_METRIC_COMPLETED、EXECUTOR_METRIC_CANCELLED、EXECUTOR_METRIC_TOTAL_START_LATENCY、EXECUTOR_METRIC_TOTAL_EXECUTION_TIME,可以看到三种 Executor 复用了同一套指标命名体系,只是前缀不同。

统计开关与一个重要的能力差异

默认开启,可配置关闭

统计数据默认开启,可以通过配置关闭。对应配置项在两个配置类中均为statisticsEnabled,且默认值都是true:

  • ScheduledExecutorConfig:private boolean statisticsEnabled = true;,配套setStatisticsEnabled(boolean);
  • DurableExecutorConfig:private boolean statisticsEnabled = true;,配套setStatisticsEnabled(boolean)。

程序化配置示例(YAML/XML 中同样存在statistics-enabled对应字段):

Config config = new Config(); // 关闭某个命名 Scheduled Executor 的统计 config.getScheduledExecutorConfig("my-scheduled") .setStatisticsEnabled(false); // 关闭某个命名 Durable Executor 的统计 config.getDurableExecutorConfig("my-durable") .setStatisticsEnabled(false);

从源码实现看,statisticsEnabled并不只影响指标导出:在 TaskRunner 与 DurableExecutorContainer 中,它同样作为采集端开关存在——关闭后,任务生命周期内不会再调用startPending/startExecution/finishExecution等计数方法,从而完全避免统计采集带来的开销。

注意:Durable Executor 没有取消统计

文档末尾有一条明确说明:

NOTE: Durable executor service has no cancellation capability hence no stats available for it.

Durable Executor Service 的语义是任务一经提交便持续执行直至完成(任务写入 ring buffer、具备持久性与至少一次执行保障),不提供取消能力,因此durableExecutor.cancelled指标虽然占用位,但不会产生非零值;cancelled统计仅对普通 Executor 与 Scheduled Executor 有意义。这解释了为什么示例输出中cancelled一栏恒为 0,也提醒读者在解读 Durable 执行器指标时不要期待取消相关的观测信息。

验证与可观测入口

  • 指标实现本身的单元测试见 hazelcast/src/test/java/com/hazelcast/internal/monitor/impl/LocalExecutorStatsImplTest.java,覆盖LocalExecutorStatsImpl的计数迁移逻辑;
  • 普通 Executor 的统计行为测试见 hazelcast/src/test/java/com/hazelcast/executor/SingleNodeTest.java 与 ClientExecutorServiceTest.java,Durable 侧可参考 hazelcast/src/test/java/com/hazelcast/durableexecutor/DurableSingleNodeTest.java;
  • 集群层面的指标开关hazelcast.metrics.datastructures(METRICS_DATASTRUCTURES)控制所有数据结构类动态指标(含本设计的两个执行器)是否注册到 metrics 注册表。

在实际运维中,除了在 Management Center 的数据结构监控页查看这些指标,也可以通过 Hazelcast 的指标导出端点(默认 REST/诊断端点会按上述前缀输出同名指标)直接观测scheduledExecutor.*与durableExecutor.*系列数值,用于容量规划、调度延迟分析与执行器健康度监控。

小结

从 01-executor-service-stats.md 这份设计出发,可以看到 Hazelcast 在 4.1 中完成了一次低成本、高一致性的观测性补齐:

  1. 模型复用:Scheduled 与 Durable Executor 直接复用普通 Executor 的LocalExecutorStats六指标模型,无需另起炉灶;
  2. 采集统一:ExecutorStats容器按执行器名称维护每实例计数,TaskRunner/DurableExecutorContainer在任务提交、开始、完成、拒绝、取消的生命周期节点精确记账;
  3. 导出统一:两个分布式服务类实现DynamicMetricsProvider,在METRICS_DATASTRUCTURES开启时向 metrics 子系统注册,分别以scheduledExecutor、durableExecutor前缀对外发布;
  4. 开关可控:statisticsEnabled默认开启、可逐执行器关闭,既保证开箱即用的可观测性,也为高吞吐场景留出去除采集开销的余地;
  5. 边界清晰:Durable Executor 因不提供取消能力,cancelled指标恒为零,这是由产品语义决定而非实现缺陷。

对于需要精确掌握调度任务排队深度、启动延迟与执行耗时的集群,这套指标是定位瓶颈和评估执行器健康状况的直接依据。

  • 缓存
  • KV存储
  • 消息队列
  • 流处理
  • 后端

【免费下载链接】hazelcast

Hazelcast is a unified real-time data platform combining stream processing with a fast data store, allowing customers to act instantly on>项目地址:https://gitcode.com/gh_mirrors/ha/hazelcast

点击查看免费下载
上一篇:SillyTavern提示词优化终极指南:从新手到专家的实战秘籍
下一篇:Ventoy革命性教程:5分钟打造万能U盘启动盘,告别重复制作烦恼

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询