ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Vector 的 NATS Source 实战指南:从 Subject 与 JetStream 订阅可观测性数据

Vector 的 NATS Source 实战指南:从 Subject 与 JetStream 订阅可观测性数据 Vector 的 NATS Source 实战指南从 Subject 与 JetStream 订阅可观测性数据【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector导读natssource 是 Vector 内置的日志采集组件用于从 NATS、生成的配置定义为主线结合其源码实现配置结构、运行逻辑、认证与 TLS 辅助展开。读完本文你将掌握natssource 的完整配置语法、Core 订阅与 JetStream 拉取两种工作模式、四种认证方式、输出事件结构以及它作为聚合层数据入口的典型部署形态。组件定位与能力概览从 CUE 元数据可以看出该组件的定位是aggregator聚合器部署角色交付方式为best_effort开发状态为stable组件元数据。它通过 TCP 协议连接 NATS 服务端默认端口4222属于incoming方向的协议入口。其核心特性包括组件元数据特性值说明支持 TLS是可校验证书、可校验主机名默认不强制开启且可按 scheme 启用Codecs是默认 framing 为bytes可使用 framing decoding 灵活解析载荷多行聚合否不支持 multiline 聚合Checkpoint否不维护消费位点JetStream 模式下由服务端 durable consumer 跟踪确认机制是仅 JetStream 模式支持消息确认acknowledgement组件说明中特别强调natssource 底层使用 Rust 的nats.rs。配置项详解完整参数表natssource 的全部配置项定义于 生成的配置定义对应的 Rust 结构体为NatsSourceConfigconfig.rs。下面逐项说明。必填参数参数类型说明官方示例值urlstringNATS 连接地址形如nats://server:port端口省略时默认4222支持逗号分隔的多个地址以实现故障转移nats://demo.nats.io、nats://127.0.0.1:4242、nats://localhost:4222,nats://localhost:5222,nats://localhost:6222subjectstring要订阅的 NATS subject支持通配符见下文foo、time.us.east、time.*.east、time.、connection_namestring分配给 NATS 连接的名称别名name便于在服务端识别连接来源vector关于url的多地址支持源码中通过parse_server_addresses将逗号分隔的字符串逐个解析为ServerAddr再交由async_nats客户端连接config.rs。因此一个 source 可以同时指向 NATS 集群的多个节点。关于subject通配符NATS 使用点分命名空间*匹配单层匹配一个或多个尾部层级。例如time.*.east匹配time.us.east但不匹配time.us.westtime.匹配time下的所有层级。使用即可订阅全部消息。可选参数参数类型默认值说明queuestring无要加入的 NATS queue group用于在多个消费者间负载均衡subject_key_fieldstringsubject消息 subject 写入事件的目标字段名subscriber_capacityuint65536底层 NATS 订阅者的缓冲容量决定内部缓冲多少条消息后才丢弃framingobjectbytes帧解析配置决定如何在字节流中切分事件decodingobject取决于 codec反序列化配置决定如何把原始字节解码为事件部分解码器还能决定输出类型log/metric/tracetlsobject无TLS 连接选项见下文authobject无认证策略见下文jetstreamobject无启用 NATS JetStream 模式见下文log_namespacebool全局设置日志命名空间覆盖项文档中隐藏subscriber_capacity的默认值65536定义于源码常量default_subscription_capacityconfig.rs。该值会通过ConnectOptions::subscription_capacity传递给底层客户端config.rs影响背压行为——当管道下游处理不过来时消息在订阅缓冲区内排队。一个最小可用配置以下是源码中GenerateConfig生成的标准示例config.rssources: nats: type: nats connection_name: vector subject: from.vector url: nats://127.0.0.1:4222使用 NATS JetStream从 Stream 拉取消息当配置了jetstream字段时source 进入 JetStream 拉取模式否则走 Core NATS 订阅模式由mode()方法判定config.rs。jetstream配置项结构如下sources: nats: type: nats url: nats://127.0.0.1:4222 connection_name: vector subject: from.vector jetstream: stream: my-stream # 必填要绑定的 Stream 名称 consumer: my-consumer # 必填要拉取的 durable consumer 名称 batch_config: batch: 200 # 可选默认 200单次拉取的最大消息条数 max_bytes: 0 # 可选默认 0单次拉取的字节上限0 表示不限stream与consumer均为必填项batch_config有两个子参数生成的配置定义batch默认200单次批量拉取的最大消息数。max_bytes默认0批量拉取的字节上限满足batch或max_bytes任一条件即返回。源码create_consumer_stream展示了 JetStream 模式的完整初始化链路先通过jetstream::new创建 JetStream 上下文再get_stream获取指定 Stream、get_consumer获取指定 consumer最后用max_messages_per_batch与max_bytes_per_batch构建拉取流source.rs。JetStream 模式下的消息确认与恢复在 JetStream 模式下can_acknowledge()返回trueconfig.rs。每条消息处理流程为source.rs成功事件发送下游成功后调用msg.ack()向服务端确认消息不会重投。解码失败不发送 ack消息将被 NATS 服务端重新投递避免数据丢失。拉取流中断source 会记录告警并进入指数退避重连循环ExponentialBackoff最大延迟 30 秒重新创建 consumer 流。由于使用的是 durable consumer服务端会保存投递状态恢复后从上次位置继续拉取。这意味着 JetStream 模式天然适合对可靠性要求较高的场景而 Core 模式更适合简单的实时订阅。认证与 TLS 配置认证策略auth支持四种方式由strategy标签区分src/nats.rs源码中四种方式的解析均有对应单元测试验证src/nats.rs。用户名 / 密码sources: nats: type: nats url: nats://127.0.0.1:4222 connection_name: vector subject: foo auth: strategy: user_password user_password: user: username password: passwordTokenauth: strategy: token token: value: my-token凭证文件JWT 体系auth: strategy: credentials_file credentials_file: path: /etc/nats/nats.credsNKeyauth: strategy: nkey nkey: nkey: UC4... # 相当于公钥 seed: SUAA... # 相当于私钥种子TLStls配置支持标准 TLS 字段enabled、ca_file、crt_file、key_file、verify_certificate、verify_hostname等。从源码可确认的实现细节src/nats.rs未启用 TLS 时直接返回明文连接启用后可通过ca_file添加根证书通过crt_filekey_file成对提供客户端证书若只提供证书或只提供密钥validate_tls_cert_key_pair会分别报出missing key/missing cert错误src/nats.rs。sources: nats: type: nats url: tls://127.0.0.1:4222 connection_name: vector subject: foo tls: enabled: true ca_file: /etc/ssl/certs/ca.pem crt_file: /etc/ssl/client.pem key_file: /etc/ssl/client.key输出事件结构natssource 输出log 事件每个事件对应一条 NATS 记录组件元数据。在旧版Legacy命名空间下事件包含以下字段字段是否必填类型说明message是stringNATS 消息的原始载荷文本source_type是string来源类型名称固定为natssubject是string消息来源的 NATS subjecttimestamp是timestamp事件时间戳其中subject字段名可通过subject_key_field参数自定义默认subject。源码在process_message中负责注入这些元数据source.rsinsert_standard_vector_source_metadata写入source_type与时间戳insert_source_metadata把msg.subject写入 subject 字段使用InsertIfEmpty语义即字段已存在则不覆盖。在 Vector 命名空间模式下元数据则写入vector.source_type、vector.ingest_timestamp与nats.subject事件主体即message。config.rs 中的两个 schema 测试output_schema_definition_vector_namespace与output_schema_definition_legacy_namespace分别断言了这两种命名空间下的事件结构config.rs。底层工作原理Core 与 JetStream 两条运行路径build方法根据配置选择运行路径config.rsCore NATS 模式create_subscription建立连接并创建订阅——未配置queue时使用subscribe配置后使用queue_subscribe加入队列组source.rs。随后run_nats_core进入循环持续读取订阅流中的消息解码并发送下游收到关闭信号时调用subscriber.drain()优雅排空订阅source.rs。JetStream 模式如前述通过 durable consumer 拉取消息支持 ack 与断线恢复source.rs。两条路径共用process_message完成解码与元数据注入source.rs它使用DecoderFramedRead按 framing 规则切分、按 decoding 规则解码单条 NATS 消息载荷可能解码出多个事件解码错误会记录日志并按可恢复性决定是否中断。每条消息同时会发出bytes_received与events_received内部指标便于观测吞吐量source.rs。从源码结构看queue队列组机制可在多个 Vector 实例订阅同一 subject 时实现水平扩展与负载均衡而subscriber_capacity则为每个订阅提供了高达 65536 条消息的背压缓冲。部署形态与实操建议由于该组件的部署角色为aggregator典型的拓扑是把分散的 NATS 消息汇聚到 Vector 聚合节点经remap、filter等转换后再路由到下游 sink。以下是一个完整示例——从 NATS 读取日志并通过consolesink 输出便于快速验证sources: nats: type: nats connection_name: vector subject: logs. url: nats://127.0.0.1:4222 sinks: stdout: type: console inputs: [nats] encoding: codec: json验证配置与运行# 校验配置在仓库根目录执行 vector validate --config-yaml (cat EOF sources: nats: type: nats connection_name: vector subject: logs. url: nats://127.0.0.1:4222 sinks: stdout: type: console inputs: [nats] encoding: codec: json EOF ) # 运行 vector --config your-config.yaml仓库还提供了 NATS 集成测试测试用例覆盖了 Core 订阅、JetStream 拉取、认证与 TLS 等场景是理解组件行为的权威参考。若需做更深层的源码研究建议按如下顺序阅读组件元数据——组件定位、特性、输出定义生成的配置定义——全部配置项的 schema 与默认值配置结构——配置解析、默认值常量与模式判定运行逻辑——Core/JetStream 两条运行路径与消息处理认证与 TLS 辅助——四种认证方式与 TLS 连接的底层封装。总结natssource 是 Vector 生态中连接 NATS 消息系统的官方入口其核心价值在于以简洁的声明式配置对接 NATS subject 与 JetStream同时保留对认证用户名密码、Token、凭证文件、NKey、TLS、队列组、背压缓冲、消息确认等生产级细节的完整控制。理解 Core 与 JetStream 两种模式的区别实时订阅 vs 可靠拉取 ack并根据实际可靠性要求选择配置是把它用好、用对的关键。【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表