ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Envoy Kafka Stats Sink 完全指南:将指标流直写 Kafka 的统计输出扩展

Envoy Kafka Stats Sink 完全指南:将指标流直写 Kafka 的统计输出扩展 Envoy Kafka Stats Sink 完全指南将指标流直写 Kafka 的统计输出扩展【免费下载链接】envoyCloud-native high-performance edge/middle/service proxy项目地址: https://gitcode.com/GitHub_Trending/en/envoy导读Kafka Stats Sink 是 Envoy 的统计输出stats sink扩展它绕过传统的 gRPC 指标采集链路通过 librdkafka 将指标序列化后直接写入 Kafka topic特别适合大规模部署场景下降低指标采集带来的内存压力。本文以仓库中的官方 API 文档与配置说明为主线结合KafkaStatsSinkConfig的 proto 定义与contrib/kafka/stat_sinks源码实现系统讲解其配置参数、两种序列化格式、认证方式与底层工作原理读完你即可独立完成 Kafka 指标直写管线的搭建与调优。本文主题对应仓库中的三份核心资料API 导航入口 docs/root/api-v3/config/contrib/kafka_stats_sink/kafka_stats_sink.rst、配置详解文档 kafka_stat_sink.rst以及消息定义 kafka_stats_sink.proto。一、Kafka Stats Sink 是什么根据 kafka_stat_sink.rst 的说明KafkaStatsSinkConfig配置了一个直接通过 librdkafka 将指标生产produce到 Apache Kafka topic 的统计输出扩展。它的典型价值体现在高规模部署场景当指标通过 gRPCmetrics_service采集器转发时中间采集层会引入不可忽视的内存压力而 Kafka Stats Sink 允许 Envoy 直连 Kafka消除中间采集基础设施让指标以消息形式落地由下游消费者如时序数据库、监控平台自行消费。使用前需注意两点限制文档中明确以 attention 标注该扩展仅包含在 contrib 镜像中普通发行版默认不携带该扩展目前处于实验阶段experimental仍在积极开发中功能会持续扩充配置结构未来可能发生变化。二、扩展类型与注册信息Kafka Stats Sink 的扩展类型为envoy.stat_sinks.kafka对应的配置类型 URL 为type.googleapis.com/envoy.extensions.stat_sinks.kafka.v3.KafkaStatsSinkConfig在源码层面扩展工厂定义于 contrib/kafka/stat_sinks/source/config.h其中KafkaStatsSinkName常量即envoy.stat_sinks.kafka并通过REGISTER_FACTORY(KafkaStatsSinkFactory, Server::Configuration::StatsSinkFactory)注册进 Envoy 的扩展注册表见 config.cc。这意味着你只需在 bootstrap 配置的stats_sinks列表中按名字引用它即可。三、消息定义与配置参数详解完整的配置字段定义在 kafka_stats_sink.proto共 8 个字段。下表汇总了各字段的语义、类型与默认行为字段类型必填说明broker_liststring是Kafka broker 地址列表host:port格式、逗号分隔至少指定一个 brokertopicstring是指标生产的目标 Kafka topicbatch_sizeuint32否单条 Kafka 消息内聚合的指标条数为 0 或未设置时一次 flush 的所有指标放在一条消息里指标很多时用于控制单条消息体积report_counters_as_deltasbool否为 true 时 counter 上报为自上次 flush 以来的增量而非累计绝对值默认 falseemit_tags_as_labelsbool否为 true 时使用去除 tag 后的指标名并将 tag 作为独立标签/JSON 字段输出为 false 时使用含 tag 值的完整指标名。默认 trueproducer_configmapstring, string否额外 librdkafka 生产者配置键值对直接透传给 librdkafka可配置压缩compression.type、认证security.protocol、sasl.mechanism等、批处理batch.num.messages等buffer_flush_timeout_msuint32否缓冲消息后强制 produce 的最大等待时间毫秒对应 librdkafka 的linger.ms未设置时默认 500msformatenum否指标消息的序列化格式默认 JSON3.1 序列化格式枚举SerializationFormatproto 定义了两种取值见 kafka_stats_sink.protoJSON 0默认指标被编码为人类可读的 JSON 对象便于用标准 Kafka 工具直接查看与消费PROTOBUF 1每条 Kafka 消息的 value 是二进制序列化的envoy.service.metrics.v3.StreamMetricsMessage内部包含io.prometheus.client.MetricFamily条目——这与 gRPCenvoy.stat_sinks.metrics_service输出端使用的线上格式完全一致消费者可以复用已有的 Protobuf 反序列化器无需为 Kafka 单独开发解码逻辑。四、完整配置示例以下 YAML 来自 kafka_stat_sink.rst 的官方示例覆盖了全部常用参数可直接作为 bootstrap 配置模板stats_flush_interval: 10s stats_sinks: - name: envoy.stat_sinks.kafka typed_config: type: type.googleapis.com/envoy.extensions.stat_sinks.kafka.v3.KafkaStatsSinkConfig broker_list: kafka1:9092,kafka2:9092 topic: envoy-metrics batch_size: 100 format: PROTOBUF emit_tags_as_labels: true report_counters_as_deltas: true buffer_flush_timeout_ms: 500解读几个关键搭配stats_flush_interval: 10s决定 Envoy 周期性抓取指标快照并调用 sink 的频率实际生产节奏由该值驱动broker_list支持多 broker 逗号分隔满足高可用需求batch_size: 100表示每 100 条指标封装为一条 Kafka 消息能显著减少消息数量与网络开销format: PROTOBUF选择与 metrics_service 一致的紧凑二进制格式buffer_flush_timeout_ms: 500控制延迟与吞吐的平衡——消息在缓冲区内等待聚合最多 500ms 强制发送。五、认证与加密配置AuthenticationKafka Stats Sink 的认证与加密完全通过producer_config映射表实现——该 map 的键值对会被直接透传给 librdkafka。官方给出的 SASL/SCRAM TLS 示例stats_sinks: - name: envoy.stat_sinks.kafka typed_config: type: type.googleapis.com/envoy.extensions.stat_sinks.kafka.v3.KafkaStatsSinkConfig broker_list: kafka:9093 topic: envoy-metrics format: PROTOBUF producer_config: security.protocol: SASL_SSL sasl.mechanism: SCRAM-SHA-256 sasl.username: envoy sasl.password: secret ssl.ca.location: /etc/ssl/certs/ca.pem这里security.protocol: SASL_SSL同时启用 SASL 认证与 TLS 加密sasl.mechanism指定机制如 SCRAM-SHA-256ssl.ca.location指定 CA 证书路径。由于producer_config与 librdkafka 全量配置属性对齐你还可以通过它配置compression.type消息压缩、batch.num.messages批处理条数、request.required.acks生产确认级别等 librdkafka 的完整能力完整属性清单可查阅 librdkafka 的 CONFIGURATION.md 配置参考文档中注明的外部资料此处不展开。六、源码级实现原理配置解析与 sink 装配在 contrib/kafka/stat_sinks/source/config.cc 中完成核心流程如下配置校验createStatsSink通过MessageUtil::downcastAndValidate将传入消息转换为KafkaStatsSinkConfig并进行静态校验broker_list与topic均有min_len: 1的 validate 规则构建 librdkafka 全局配置创建RdKafka::Conf::CONF_GLOBAL配置对象先写入bootstrap.servers broker_list再将buffer_flush_timeout_ms未设置时经PROTOBUF_GET_WRAPPED_OR_DEFAULT取默认值500写入linger.ms透传用户配置遍历producer_config的每个键值对逐一调用setConfProperty任一属性设置失败都会返回带错误信息的InvalidArgumentError创建生产者通过LibRdKafkaUtilsImpl的默认实例创建RdKafka::Producer失败则返回InternalError装配 flush 器按report_counters_as_deltas默认 false、emit_tags_as_labels默认 true、format默认 Json构造KafkaMetricsFlusher最终生成KafkaStatsSink实例。sink实例本身kafka_stats_sink_impl.h持有 producer、topic 与 batch_size实现Stats::Sink接口flush()时调用KafkaMetricsFlusher完成序列化再由produce()逐条发送到 Kafka。构建依赖方面BUILD 中kafka_stats_sink_impl_lib直接依赖//bazel/deps:librdkafka且该扩展标注了skip_on_windows即当前不支持 Windows 构建。6.1 JSON 序列化细节flushJsonkafka_stats_sink_impl.cc将一次快照组织为{metrics: [...]}结构每个指标条目包含typecounter/gauge/histogram之一name与可选tags当emit_tags_as_labels为 true 时name取tagExtractedName()tag 以{env:prod,region:us-east}形式的独立对象输出为 false 时直接用含 tag 的完整metric.name()timestamp_ms快照时间戳毫秒数值字段counter默认输出累计valuereport_counters_as_deltas为 true 时输出value增量并附加delta: true标记gauge输出当前valuehistogram输出sample_count、sample_sum、buckets数组每项含upper_bound与cumulative_count以及quantiles数组每项含quantile与value。值得注意的实现细节序列化采用Json::StringStreamer流式写入并预分配 4096 字节缓冲区以减少分配只输出used()为 true 的指标未使用的指标会被过滤当batch_size为 0 时所有指标进入同一条消息即使快照为空也会产出一条{metrics:[]}消息。6.2 Protobuf 序列化细节flushProtobufkafka_stats_sink_impl.cc构造envoy.service.metrics.v3.StreamMetricsMessage将每个指标封装为io.prometheus.client.MetricFamilyfamily的name同样遵循emit_tags_as_labels语义tag 提取名 独立 label或完整指标名每个metric条目写入timestamp_mscounter 写入counter.valuedelta 模式下取快照 delta否则取累计值gauge 写入gauge.valuehistogram 写入histogram的桶与分位数数据按batch_size将多条MetricFamily聚合进同一个StreamMetricsMessage达到阈值即SerializeAsString()产出二进制消息。由于复用 Prometheus 的MetricFamily结构与 metrics_service 同款流格式任何已经为 gRPC metrics 采集编写过解码器的团队都可以零成本复用。七、测试验证行为即契约仓库为 Kafka Stats Sink 提供了完整测试kafka_stats_sink_impl_test.cc 覆盖了核心行为FlushCountersAndGauges验证 counter 以 delta 模式输出value:5与delta:true、gauge 输出当前值、时间戳字段正确快照时间 1000ms 对应timestamp_ms:1000FlushWithBatchingbatch_size1时 counter 与 gauge 被拆分为两条独立消息验证批处理逻辑FlushEmptySnapshot空快照也会产生一条{metrics:[]}消息AbsoluteCounterValuesreport_counters_as_deltasfalse时输出累计值value:10且不出现delta字段。测试夹具通过 mock 计数器、gauge、直方图与快照直接驱动KafkaMetricsFlusher::flush()这些断言正是上节序列化规则的可执行契约可作为二次开发或格式调试的参考基线。此外 config_test.cc 覆盖工厂侧的配置解析与错误处理路径。八、使用注意事项总结镜像选择仅在 contrib 镜像中提供需要按 contrib 构建方式启用对应扩展对应contrib目录下的 Bazel target稳定性预期处于实验阶段升级 Envoy 版本时应关注KafkaStatsSinkConfig的字段变更平台限制构建目标标注skip_on_windowsLinux/macOS 之外的平台不可用必填项broker_list与topic至少各填一项否则配置校验失败生产参数调优batch_size与buffer_flush_timeout_ms共同决定单条消息体积与实时性大规模指标场景建议两者配合使用如 100 条 / 500ms避免单条消息过大或延迟过高。【免费下载链接】envoyCloud-native high-performance edge/middle/service proxy项目地址: https://gitcode.com/GitHub_Trending/en/envoy创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表