AutoMQ automq-metrics 模块实践指南:基于 OpenTelemetry 的 Kafka 可观测性数据采集与多端导出
【免费下载链接】automqDiskless Kafka® on S3. 10x Cost-Effective. No Cross-AZ Traffic Cost. Autoscale in seconds. Single-digit ms latency. Multi-AZ Availability.项目地址: https://gitcode.com/GitHub_Trending/au/automq
AutoMQ 的automq-metrics模块(源码位于 automq-metrics/src/main/java/com/automq/opentelemetry)是构建在 OpenTelemetry SDK 之上的统一遥测组件,专为 AutoMQ Kafka 设计,负责 JVM 指标、JMX 指标与 Kafka 自带 Yammer 指标的采集,并通过 Prometheus、OTLP 或 S3 兼容存储三种渠道导出。读完本文,你将掌握该模块的完整模块结构、AutoMQTelemetryManager的初始化与生命周期管理、三类 Exporter 的 URI 配置与底层实现原理,并能独立完成生产环境的接入与排障。
一、模块定位与整体结构
automq-metrics是 AutoMQ 仓库中的独立 Gradle/Java 模块,其 README(automq-metrics/README.md)将包结构划分为四个部分,与源码目录逐一对应:
com.automq.opentelemetry/ ├── AutoMQTelemetryManager.java # 主管理类:初始化与生命周期 ├── TelemetryConstants.java # 常量定义(资源属性、基数限制等) ├── common/ │ ├── OTLPCompressionType.java # OTLP 压缩类型(none/gzip) │ └── OTLPProtocol.java # OTLP 协议类型(grpc/http) ├── exporter/ │ ├── MetricsExporter.java # Exporter 接口(暴露 asMetricReader()) │ ├── MetricsExportConfig.java # 导出配置接口 │ ├── MetricsExporterProvider.java # Exporter 工厂提供者(SPI) │ ├── MetricsExporterType.java # Exporter 类型枚举(otlp/prometheus/ops) │ ├── MetricsExporterURI.java # Exporter URI 解析器 │ ├── OTLPMetricsExporter.java # OTLP Exporter 实现 │ ├── PrometheusMetricsExporter.java # Prometheus Exporter 实现 │ └── s3/ # S3 指标 Exporter 实现 │ ├── CompressionUtils.java # 数据压缩工具 │ ├── PrometheusUtils.java # Prometheus 命名兼容工具 │ ├── S3MetricsExporter.java # S3 指标 Exporter 实现 │ └── S3MetricsExporterAdapter.java # 适配 S3 指标导出 └── yammer/ ├── DeltaHistogram.java # 增量直方图实现 ├── OTelMetricUtils.java # OpenTelemetry 指标工具 ├── YammerMetricsProcessor.java # Yammer 指标处理器 └── YammerMetricsReporter.java # Yammer 指标上报器模块提供了统一、可插拔的遥测数据管理能力:既能把 JVM 运行时、JMX Bean、Kafka Yammer 指标汇聚到一套 OpenTelemetry 管线中,也能根据 URI 灵活地把指标投递到 Prometheus、OTLP 兼容后端或 S3 兼容对象存储。
二、核心功能一览
1. 指标采集
- JVM 指标:自动采集 CPU、内存池、垃圾回收、线程状态等运行时指标。对应实现位于 AutoMQTelemetryManager.java,通过
MemoryPools.registerObservers、Cpu.registerObservers、GarbageCollector.registerObservers、Threads.registerObservers注册异步观测器,返回的AutoCloseable被登记在autoCloseableList中,随shutdown()统一释放。 - JMX 指标:通过 YAML 配置文件定义并采集 JMX Bean 指标,底层复用 OpenTelemetry 的
JmxMetricInsight、MetricConfiguration与RuleParser(见 AutoMQTelemetryManager.java)。 - Yammer 指标:以监听器方式桥接 Kafka 既有的 Yammer 指标体系到 OpenTelemetry。
2. 多种 Exporter 支持
- Prometheus:通过内嵌 HTTP Server 以 Prometheus 文本格式暴露指标。
- OTLP:支持 gRPC 与 HTTP/Protobuf 两种协议,投递到 OTLP 后端。
- S3:把指标序列化为 JSON 行并批量上传到 S3 兼容对象存储。
3. 灵活配置
- 支持通过 Properties 配置文件设置参数(模块内以
MetricsExportConfig接口 + 构造函数参数承载)。 - 可配置导出间隔、压缩方式、超时时间等。
- 支持指标基数(cardinality)限制以控制内存占用:默认
20000,常量定义见 TelemetryConstants.java,可通过AutoMQTelemetryManager#setMetricCardinalityLimit(int)调整。
三、快速开始:五分钟接入遥测
1. 基本用法
首先实现MetricsExportConfig接口,向模块提供集群与导出参数(接口定义见 MetricsExportConfig.java):
import com.automq.opentelemetry.AutoMQTelemetryManager; import com.automq.opentelemetry.exporter.MetricsExportConfig; // 实现 MetricsExportConfig public class MyMetricsExportConfig implements MetricsExportConfig { @Override public String clusterId() { return "my-cluster"; } @Override public boolean isLeader() { return true; } @Override public int nodeId() { return 1; } @Override public ObjectStorage objectStorage() { // Return your object storage instance for S3 exports return myObjectStorage; } @Override public List<Pair<String, String>> baseLabels() { return Arrays.asList( Pair.of("environment", "production"), Pair.of("region", "us-east-1") ); } @Override public int intervalMs() { return 60000; } // 60 seconds }随后初始化单例并启动:
// 创建导出配置 MetricsExportConfig config = new MyMetricsExportConfig(); // 初始化遥测管理器单例 AutoMQTelemetryManager manager = AutoMQTelemetryManager.initializeInstance( "prometheus://localhost:9090", // exporter URI "automq-kafka", // service name "broker-1", // instance ID config // export config ); // 启动 Yammer 指标上报(可选) MetricsRegistry yammerRegistry = // Get Kafka's Yammer registry manager.startYammerMetricsReporter(yammerRegistry); // Application running... // 关闭遥测系统 AutoMQTelemetryManager.shutdownInstance();从源码层面看,initializeInstance采用双重检查锁的单例模式(AutoMQTelemetryManager.java),同一 JVM 内只能存在一个实例;构造函数还会调用SLF4JBridgeHandler.install()把 OpenTelemetry SDK 自身的 JUL 日志重定向到 SLF4J,保证日志体系统一。初始化流程init()依次完成:构建SdkMeterProvider(含 Resource 属性与各 Exporter 的MetricReader注册)→ 构建并注册全局OpenTelemetrySdk→ 注册 JVM/JMX 指标(AutoMQTelemetryManager.java)。
2. 获取 Meter 实例并创建自定义指标
// 获取单例实例 AutoMQTelemetryManager manager = AutoMQTelemetryManager.getInstance(); // 获取 Meter 用于自定义指标 Meter meter = manager.getMeter(); // 创建自定义指标 LongCounter requestCounter = meter .counterBuilder("http_requests_total") .setDescription("Total number of HTTP requests") .build(); requestCounter.add(1, Attributes.of(AttributeKey.stringKey("method"), "GET"));getMeter()返回以TelemetryConstants.TELEMETRY_SCOPE_NAME(automq_for_kafka)为作用域名称的 Meter(AutoMQTelemetryManager.java),尚未初始化时调用会抛出IllegalStateException,因此务必先调用initializeInstance。
四、配置详解
1. 基础配置
配置通过MetricsExportConfig接口与构造函数参数提供:
| 参数 | 说明 | 示例 |
|---|---|---|
exporterUri | 指标导出 URI | prometheus://localhost:9090 |
serviceName | 遥测数据中的服务名 | automq-kafka |
instanceId | 服务实例唯一标识 | broker-1 |
config | MetricsExportConfig 实现 | 见上文示例 |
其中exporterUri的解析逻辑集中在 MetricsExporterURI.java:URI 的 scheme 决定 Exporter 类型(对应枚举 MetricsExporterType.java),支持otlp、prometheus、ops(即 S3)三种内建类型。有两个值得注意的源码细节:
- 支持多 Exporter 并联:
parse()方法将 URI 字符串按逗号,拆分,逐段解析并聚合多个 Exporter(MetricsExporterURI.java)。例如"prometheus://localhost:9090,otlp://localhost:4317"可同时向 Prometheus 与 OTLP 后端导出。 - 支持 SPI 扩展:模块通过
ServiceLoader.load(MetricsExporterProvider.class)加载外部 Exporter Provider,遇到未知 scheme 时会依次询问各 Provider 是否支持(MetricsExporterURI.java),方便第三方扩展新的导出目标。
2. Prometheus Exporter
// 使用 prometheus:// URI scheme AutoMQTelemetryManager manager = AutoMQTelemetryManager.initializeInstance( "prometheus://localhost:9090", "automq-kafka", "broker-1", config );底层实现是 OpenTelemetry 的PrometheusHttpServer(PrometheusMetricsExporter.java):默认监听localhost:9090,支持通过 URI 查询参数host/port覆盖(解析逻辑见 MetricsExporterURI.java)。它通过setAllowedResourceAttributesFilter精确控制哪些 Resource 属性会转成 Prometheus 标签:仅保留job(serviceName)、instance(instanceId)、host.name以及baseLabels()中的自定义标签,避免高基数 Resource 属性污染指标序列。初始化时还会把job/instance属性写入 Resource(见 AutoMQTelemetryManager.java),保证与原生 Prometheus 抓取语义对齐。
3. OTLP Exporter
// 使用 otlp:// URI scheme 并携带可选查询参数 AutoMQTelemetryManager manager = AutoMQTelemetryManager.initializeInstance( "otlp://localhost:4317?protocol=grpc&compression=gzip&timeout=30000", "automq-kafka", "broker-1", config );OTLPMetricsExporter(OTLPMetricsExporter.java)支持的查询参数及其默认值如下:
| 查询参数 | 说明 | 默认值 | 可选值 |
|---|---|---|---|
protocol | OTLP 传输协议 | grpc | grpc/http(见 OTLPProtocol.java) |
compression | 压缩方式 | none | none/gzip(见 OTLPCompressionType.java) |
endpoint | OTLP 后端地址 | 由 URI 的scheme://authority拼出 | 任意合法 endpoint |
timeout | 导出超时 | 固定30000ms | 见下方说明 |
按协议分支,asMetricReader()分别构建OtlpGrpcMetricExporter或OtlpHttpMetricExporter,并统一包装为PeriodicMetricReader,导出间隔取自config.intervalMs()(OTLPMetricsExporter.java)。需要说明的是:模块源码中DEFAULT_EXPORTER_TIMEOUT_MS目前固定为 30000ms,README 示例中的timeout参数暂未参与超时控制,若你的后端响应较慢,建议结合自身网络情况关注该默认值。
4. S3 指标 Exporter
// 使用 s3:// URI scheme AutoMQTelemetryManager manager = AutoMQTelemetryManager.initializeInstance( "s3://access-key:secret-key@my-bucket.s3.amazonaws.com", "automq-kafka", "broker-1", config // config.clusterId(), nodeId(), isLeader() 用于 S3 导出 );S3 Exporter 是 AutoMQ 最有特色的导出通道,适用于把指标沉淀到对象存储长期归档。实际创建路径为:URI 解析 →S3MetricsExporterAdapter(S3MetricsExporterAdapter.java)→S3MetricsExporter(S3MetricsExporter.java)。若config.objectStorage()返回null,该 Exporter 会被跳过(MetricsExporterURI.java)。
S3 导出链路的关键机制如下:
- 对象路径:
automq/metrics/{clusterId}/{nodeId}/{yyyyMMddHH}/{uuid},按小时分桶、UUID 随机命名,避免并发节点写入冲突(S3MetricsExporter.java)。 - 数据格式:每行一条 JSON,字段含
kind(absolute)、timestamp(秒)、name(Prometheus 兼容名)、counter/gauge/histogram载荷以及tags标签集(S3MetricsExporter.java)。直方图以count/sum/buckets(含upper_limit)结构保存。 - 默认标签:自动注入
host_name、job(=clusterId())、instance(=nodeId()),再合并baseLabels()(S3MetricsExporter.java)。 - 批量缓冲:指标先写入 16MB 的堆外
ByteBuf(DEFAULT_BUFFER_SIZE = 16 * 1024 * 1024),由独立 daemon 上传线程按UPLOAD_INTERVAL + 随机抖动(0~60s)周期 flush,缓冲不足时立即上传(S3MetricsExporter.java)。 - 过期清理:只有
config.isLeader()为 true 的节点才会执行清理任务,删除超过CLEANUP_INTERVAL的旧对象;由于部分 S3 实现单次请求最多 1000 个 key,删除操作会按 1000 个/批分片并allOf等待全部完成(S3MetricsExporter.java)。这意味着集群中应由一个节点承担清理职责,避免多节点重复删除。 - 环境变量:上传间隔与清理间隔可通过环境变量覆盖(S3MetricsExporter.java):
| 环境变量 | 默认值 |
|---|---|
AUTOMQ_OBSERVABILITY_UPLOAD_INTERVAL | 60000ms |
AUTOMQ_OBSERVABILITY_CLEANUP_INTERVAL | 120000ms |
S3 模式下的完整配置实现:
// S3 导出配置实现 public class S3MetricsExportConfig implements MetricsExportConfig { private final ObjectStorage objectStorage; public S3MetricsExportConfig(ObjectStorage objectStorage) { this.objectStorage = objectStorage; } @Override public String clusterId() { return "my-kafka-cluster"; } @Override public boolean isLeader() { // 集群中仅一个节点返回 true return isCurrentNodeLeader(); } @Override public int nodeId() { return 1; } @Override public ObjectStorage objectStorage() { return objectStorage; } @Override public List<Pair<String, String>> baseLabels() { return Arrays.asList(Pair.of("environment", "production")); } @Override public int intervalMs() { return 60000; } } // 使用 S3 导出初始化遥测管理器 ObjectStorage objectStorage = // 创建对象存储实例(复用 AutoMQ s3stream 的 ObjectStorage) MetricsExportConfig config = new S3MetricsExportConfig(objectStorage); AutoMQTelemetryManager manager = AutoMQTelemetryManager.initializeInstance( "s3://access-key:secret-key@my-bucket.s3.amazonaws.com", "automq-kafka", "broker-1", config ); // Application running... // 关闭遥测系统 AutoMQTelemetryManager.shutdownInstance();5. JMX 指标配置
JMX 采集规则通过 YAML 文件定义,初始化后调用setJmxConfigPaths设置路径:
AutoMQTelemetryManager manager = AutoMQTelemetryManager.initializeInstance( exporterUri, serviceName, instanceId, config ); // 初始化后设置 JMX 配置路径 manager.setJmxConfigPaths("/jmx-config.yaml,/kafka-jmx.yaml");配置文件要求
- 目录要求:配置文件必须放在项目 classpath 中(如
src/main/resources目录),支持子目录结构,如/config/jmx-metrics.yaml。 - 路径格式:路径必须以
/开头,表示从 classpath 根目录定位;多个配置文件用逗号分隔(getJmxConfigPaths()会先按逗号拆分、trim去空白并过滤空串,见 AutoMQTelemetryManager.java)。 - 文件格式:使用 YAML 格式(
.yaml或.yml后缀),文件名可自定义,建议使用有意义的命名。
推荐目录结构
src/main/resources/ ├── jmx-kafka-broker.yaml # Kafka Broker 指标配置 ├── jmx-kafka-consumer.yaml # Kafka Consumer 指标配置 ├── jmx-kafka-producer.yaml # Kafka Producer 指标配置 └── config/ ├── custom-jmx.yaml # 自定义 JMX 指标配置 └── third-party-jmx.yaml # 第三方组件 JMX 配置配置示例
JMX 配置文件示例(jmx-config.yaml):
rules: - bean: kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec metricAttribute: name: kafka_server_broker_topic_messages_in_per_sec description: Messages in per second unit: "1/s" attributes: - name: topic value: topic加载时模块通过getResourceAsStream(path)从 classpath 读取文件,文件缺失或解析失败会记录 ERROR 日志但不会中断其他配置文件的加载(AutoMQTelemetryManager.java)。JMX 采集基于 OpenTelemetry 的JmxMetricInsight,会按config.intervalMs()周期抓取 JMX 数据。
五、支持的指标类型
1. JVM 指标
- 内存使用(堆内存、非堆内存、内存池)
- CPU 使用率
- 垃圾回收统计
- 线程状态
2. Kafka 指标(经 Yammer 桥接)
Yammer 桥接器并不全量转发所有指标,而是由 OTelMetricUtils.java 白名单过滤,聚焦三类对诊断最关键的指标组:
| Yammer 分组 | 指标类型 | 指标名 | OTel 指标名前缀 |
|---|---|---|---|
kafka.network | RequestMetrics | RequestBytes | kafka.request.size |
kafka.network | RequestMetrics | TotalTimeMs | kafka.request.time |
kafka.network | RequestMetrics | RequestQueueTimeMs | kafka.request.queue.time |
kafka.network | RequestMetrics | ResponseQueueTimeMs | kafka.response.queue.time |
kafka.log | LogFlushStats | LogFlushRateAndTimeMs | kafka.logs.flush.time |
kafka.controller | ControllerEventManager | EventQueueTimeMs | kafka.event.queue.time |
kafka.controller | ControllerEventManager | EventQueueProcessingTimeMs | kafka.event.queue.processing.time |
桥接的核心逻辑在 YammerMetricsProcessor.java:只处理Histogram与Timer两类指标(Counter/Gauge/Metered会抛出UnsupportedOperationException),通过 DeltaHistogram.java 计算相邻采样之间的增量均值(delta mean),再以buildWithCallback注册 OTel gauge 回调。Yammer 指标 scope(如request=Fetch.type=consumer)会被解析成request/type等标签,其中request键自动改名为type,与 Kafka 原生指标语义对齐(YammerMetricsProcessor.java)。YammerMetricsReporter本身实现了MetricsRegistryListener,注册到MetricsRegistry后能动态感知新增/移除的指标(YammerMetricsReporter.java)。
3. 自定义指标
通过 OpenTelemetry API 创建,支持:
- Counter
- Gauge
- Histogram
- UpDownCounter
六、最佳实践
1. 生产环境配置
public class ProductionMetricsConfig implements MetricsExportConfig { @Override public String clusterId() { return "production-cluster"; } @Override public boolean isLeader() { // 实现你的主节点选举逻辑 return isCurrentNodeController(); } @Override public int nodeId() { return getCurrentNodeId(); } @Override public ObjectStorage objectStorage() { return productionObjectStorage; } @Override public List<Pair<String, String>> baseLabels() { return Arrays.asList( Pair.of("environment", "production"), Pair.of("region", System.getenv("AWS_REGION")), Pair.of("version", getApplicationVersion()) ); } @Override public int intervalMs() { return 60000; } // 1 分钟 } // 生产环境初始化 AutoMQTelemetryManager manager = AutoMQTelemetryManager.initializeInstance( "prometheus://0.0.0.0:9090", // 或用 S3 URI 做对象存储归档 "automq-kafka", System.getenv("HOSTNAME"), new ProductionMetricsConfig() );生产环境建议:把isLeader()与节点真实主备/控制节点状态绑定,确保 S3 清理任务只由单一节点执行;baseLabels()中的标签要控制基数,避免把节点 IP、请求 ID 等易变值放入标签。
2. 开发环境配置
public class DevelopmentMetricsConfig implements MetricsExportConfig { @Override public String clusterId() { return "dev-cluster"; } @Override public boolean isLeader() { return true; } // 开发环境单节点 @Override public int nodeId() { return 1; } @Override public ObjectStorage objectStorage() { return null; } // OTLP 场景不需要 @Override public List<Pair<String, String>> baseLabels() { return Arrays.asList(Pair.of("environment", "development")); } @Override public int intervalMs() { return 10000; } // 10 秒,便于快速反馈 } // 开发环境初始化 AutoMQTelemetryManager manager = AutoMQTelemetryManager.initializeInstance( "otlp://localhost:4317", "automq-kafka-dev", "local-dev", new DevelopmentMetricsConfig() );开发环境建议:导出间隔缩短到 10 秒以加速调试反馈;objectStorage()返回null即可,不会触发 S3 导出。
3. 资源管理
- 设置合理的指标基数限制以避免内存泄漏:默认值为
20000,基数选择器在 AutoMQTelemetryManager.java 中对每个 Instrument 生效,可通过setMetricCardinalityLimit(int)按需调低。 - 应用关闭时调用
shutdownInstance()(或实例的shutdown())释放资源:它会依次关闭 JVM 指标观测器、forceFlush并关闭所有MetricReader、最后关闭OpenTelemetrySdk(AutoMQTelemetryManager.java),确保缓冲中的指标在退出前完成导出。 - 关注 Exporter 健康状态:S3 上传失败、OTLP 端点不可达都会在
com.automq.opentelemetry日志中输出 ERROR 记录。
七、Prometheus 命名兼容细节(源码级)
S3 导出采用 Prometheus 兼容命名,核心转换逻辑在 PrometheusUtils.java:
- 单位映射:OTel 单位(如
s、ms、By、1/s、By/s)映射为 Prometheus 单位后缀(seconds、milliseconds、bytes、per_second、bytes_per_second)。 - 计数器后缀:单调递增的 counter 自动追加
_total后缀。 - 保留后缀防护:
sanitizeMetricName会剥离_total、_created、_bucket、_info等保留后缀,避免重名冲突。 - 标签名清洗:
mapLabelName将.替换为_(mapLabelName("host.name")→host_name),与 S3 默认标签的注入逻辑一致。
这些规则有对应单元测试验证,见 PrometheusUtilsTest.java,例如mapMetricsName("foo", "1", true, false)期望得到foo_ratio_total,mapMetricsName("kafka_tabletopic_fields", "1/s", false, true)期望得到kafka_tabletopic_fields_per_second。
八、故障排查
常见问题
- 指标未导出
- 检查传给
initializeInstance()的 exporter URI 是否正确; - 确认目标端点可达(Prometheus 端口、OTLP 后端地址、S3 bucket);
- 查看日志中的错误信息;
- 确保
MetricsExportConfig.intervalMs()返回合理值(如 60000)。
- 检查传给
- JMX 指标缺失
- 确认通过
setJmxConfigPaths()设置的路径正确且以/开头; - 检查 YAML 配置格式与 bean 名称;
- 验证 JMX Bean 真实存在;
- 确保文件位于 classpath 中(模块通过
getResourceAsStream定位)。
- 确认通过
- 内存占用偏高
- 在
MetricsExportConfig中实现基数限制(或调低setMetricCardinalityLimit); - 检查
baseLabels()是否存在高基数标签; - 考虑通过
intervalMs()增大导出间隔。
- 在
日志配置
模块日志使用 SLF4J,可通过日志框架(logback.xml、log4j2.xml)开启调试:
<!-- 对于 Logback --> <logger name="com.automq.opentelemetry" level="DEBUG" /> <logger name="io.opentelemetry" level="INFO" />DEBUG级别会输出 Exporter 初始化参数(端点、协议、压缩方式、间隔)、指标注册与导出错误等关键信息,是排查上述问题的主要手段。
九、依赖与许可
- Java 8+
- OpenTelemetry SDK 1.30+(含
opentelemetry-instrumentation-jmx、opentelemetry-exporter-prometheus、opentelemetry-exporter-otlp、opentelemetry-exporter-otlp-http、runtime-metrics 系列) - Apache Commons Lang3(
Pair、StringUtils) - SLF4J 日志框架(含
jul-to-slf4j桥接) - Jackson(JSON 序列化)、Netty(堆外缓冲)与 AutoMQ s3stream 的
ObjectStorage抽象
模块采用 Apache License 2.0 开源。整体接入路径为:实现MetricsExportConfig→ 构造 exporter URI →initializeInstance初始化 → 按需setJmxConfigPaths/startYammerMetricsReporter→ 应用退出时shutdownInstance,通过上述配置即可让 AutoMQ Kafka 集群的遥测数据稳定流向 Prometheus、OTLP 或 S3。
【免费下载链接】automqDiskless Kafka® on S3. 10x Cost-Effective. No Cross-AZ Traffic Cost. Autoscale in seconds. Single-digit ms latency. Multi-AZ Availability.项目地址: https://gitcode.com/GitHub_Trending/au/automq
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考