ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Hazelcast Scheduled 与 Durable Executor 服务指标统计的设计与实现

Hazelcast Scheduled 与 Durable Executor 服务指标统计的设计与实现 缓存KV存储消息队列流处理后端【免费下载链接】hazelcastHazelcast 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点击查看免费下载本设计文档对应仓库路径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 ServiceIExecutorService通用分布式任务提交任务绑定分区执行有已有完整统计Scheduled Executor ServiceIScheduledExecutorService支持单次/固定频率调度、任务可持久化与迁移无Durable Executor ServiceIDurableExecutorService任务写入 ring buffer、具备持久性与至少一次执行语义无在 4.1 之前三种实现中只有普通 Executor Service 维护运行统计Scheduled 与 Durable 两种均未采集任何指标。本设计的目标就是把 Executor Service 已有的同一套统计模型扩展到另外两种执行器并让它们统一通过 Management Center 可见、可监控。设计统一复用LocalExecutorStats统计模型统计模型接口整个统计能力建立在LocalExecutorStats接口之上其定义位于 hazelcast/src/main/java/com/hazelcast/executor/LocalExecutorStats.javapublic 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所有已完成任务的实际执行时长累计毫秒。其实现类为LocalExecutorStatsImplhazelcast/src/main/java/com/hazelcast/internal/monitor/impl/LocalExecutorStatsImpl.java。从源码可以看到两个值得注意的实现细节无锁并发计数所有计数器字段pending、started、completed、cancelled、totalStartLatency、totalExecutionTime都是volatile long通过AtomicLongFieldUpdater做原子自增/累加避免在热路径上引入锁竞争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对象内部用ConcurrentHashMapString, 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构造时startPendingrun()开始时startExecution(name, start - creationTime)执行结束后未被取消的前提下finishExecution(name, Clock.currentTimeMillis() - start)在 ring buffer 提交或分区迁移后重放任务时若RejectedExecutionException发生则调用rejectExecution(name)回退 pending 计数。指标导出DynamicMetricsProvider与 Metrics 子系统统计数据采集完成后通过 Hazelcast 的 metrics 子系统对外发布供 Management Center 拉取。核心机制是两个服务类实现DynamicMetricsProvider接口并在init()阶段向 metrics 注册表注册自己DistributedScheduledExecutorServiceSERVICE_NAME hz:impl:scheduledExecutorServiceDistributedDurableExecutorServiceSERVICE_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()返回的MapString, LocalExecutorStatsImpl逐条写入指标上下文。两个服务各自的getStats()都直接返回executorStats.getStatsMap()。Metrics 前缀约定文档明确了两种服务使用的指标前缀对应常量定义于 hazelcast/src/main/java/com/hazelcast/internal/metrics/MetricDescriptorConstants.java服务常量前缀值Scheduled ExecutorSCHEDULED_EXECUTOR_PREFIXscheduledExecutorDurable ExecutorDURABLE_EXECUTOR_PREFIXdurableExecutor该文件头部有一段重要注释这些常量必须在 minor 版本之间保持稳定修改它们可能破坏 Management Center 或其他指标消费方。因此scheduledExecutor.*与durableExecutor.*系列指标名称是事实上的对外契约。每个命名执行器的指标都会附带nameexecutorName判别维度discriminator以区分同一个 JVM 上多个不同名字的执行器。示例输出文档给出了启用统计后典型的指标输出以scheduledExecutor前缀为例unit为毫秒或计数[nameexecutorName,unitms,metricscheduledExecutor.creationTime]1598016899537 [nameexecutorName,unitcount,metricscheduledExecutor.pending]0 [nameexecutorName,unitcount,metricscheduledExecutor.started]1 [nameexecutorName,unitcount,metricscheduledExecutor.completed]0 [nameexecutorName,unitcount,metricscheduledExecutor.cancelled]0 [nameexecutorName,unitms,metricscheduledExecutor.totalStartLatency]2 [nameexecutorName,unitms,metricscheduledExecutor.totalExecutionTime]0durableExecutor前缀的指标结构与之一一对应只是前缀不同。这七项指标名称分别来自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且默认值都是trueScheduledExecutorConfigprivate boolean statisticsEnabled true;配套setStatisticsEnabled(boolean)DurableExecutorConfigprivate 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.javaDurable 侧可参考 hazelcast/src/test/java/com/hazelcast/durableexecutor/DurableSingleNodeTest.java集群层面的指标开关hazelcast.metrics.datastructuresMETRICS_DATASTRUCTURES控制所有数据结构类动态指标含本设计的两个执行器是否注册到 metrics 注册表。在实际运维中除了在 Management Center 的数据结构监控页查看这些指标也可以通过 Hazelcast 的指标导出端点默认 REST/诊断端点会按上述前缀输出同名指标直接观测scheduledExecutor.*与durableExecutor.*系列数值用于容量规划、调度延迟分析与执行器健康度监控。小结从 01-executor-service-stats.md 这份设计出发可以看到 Hazelcast 在 4.1 中完成了一次低成本、高一致性的观测性补齐模型复用Scheduled 与 Durable Executor 直接复用普通 Executor 的LocalExecutorStats六指标模型无需另起炉灶采集统一ExecutorStats容器按执行器名称维护每实例计数TaskRunner/DurableExecutorContainer在任务提交、开始、完成、拒绝、取消的生命周期节点精确记账导出统一两个分布式服务类实现DynamicMetricsProvider在METRICS_DATASTRUCTURES开启时向 metrics 子系统注册分别以scheduledExecutor、durableExecutor前缀对外发布开关可控statisticsEnabled默认开启、可逐执行器关闭既保证开箱即用的可观测性也为高吞吐场景留出去除采集开销的余地边界清晰Durable Executor 因不提供取消能力cancelled指标恒为零这是由产品语义决定而非实现缺陷。对于需要精确掌握调度任务排队深度、启动延迟与执行耗时的集群这套指标是定位瓶颈和评估执行器健康状况的直接依据。赞分享缓存KV存储消息队列流处理后端【免费下载链接】hazelcastHazelcast 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点击查看免费下载相关推荐Hazelcast 可卸载 Entry Processor 执行器统计Offloaded Entry Processor Executor Stats设计解析Hazelcast 可卸载 Entry Processor 执行器统计Offloaded Entry Processor Executor Stats设计解缓存KV存储消息队列流处理后端FastDFS存储服务器统计指标API设计与文档FastDFS存储服务器统计指标API设计与文档 引言为什么需要统计指标API 在大规模分布式文件系统Distributed File System,分布式文件系统存储后端深度拆解5大架构难题redux-persist高级开发者的优化指南深度拆解5大架构难题redux persist高级开发者的优化指南 redux persist作为Redux生态中状态持久化的核心技术在实际应用中常面临架构前端上一篇SillyTavern提示词优化终极指南从新手到专家的实战秘籍下一篇Ventoy革命性教程5分钟打造万能U盘启动盘告别重复制作烦恼创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进