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 结构体为NatsSourceConfig(config.rs)。下面逐项说明。
必填参数
| 参数 | 类型 | 说明 | 官方示例值 |
|---|---|---|---|
url | string | NATS 连接地址,形如nats://server:port,端口省略时默认4222;支持逗号分隔的多个地址以实现故障转移 | nats://demo.nats.io、nats://127.0.0.1:4242、nats://localhost:4222,nats://localhost:5222,nats://localhost:6222 |
subject | string | 要订阅的 NATS subject,支持通配符(见下文) | foo、time.us.east、time.*.east、time.>、> |
connection_name | string | 分配给 NATS 连接的名称,别名name,便于在服务端识别连接来源 | vector |
关于url的多地址支持,源码中通过parse_server_addresses将逗号分隔的字符串逐个解析为ServerAddr,再交由async_nats客户端连接(config.rs)。因此一个 source 可以同时指向 NATS 集群的多个节点。
关于subject通配符:NATS 使用点分命名空间,*匹配单层,>匹配一个或多个尾部层级。例如time.*.east匹配time.us.east但不匹配time.us.west,time.>匹配time下的所有层级。使用>即可订阅全部消息。
可选参数
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
queue | string | 无 | 要加入的 NATS queue group,用于在多个消费者间负载均衡 |
subject_key_field | string | subject | 消息 subject 写入事件的目标字段名 |
subscriber_capacity | uint | 65536 | 底层 NATS 订阅者的缓冲容量,决定内部缓冲多少条消息后才丢弃 |
framing | object | bytes | 帧解析配置,决定如何在字节流中切分事件 |
decoding | object | 取决于 codec | 反序列化配置,决定如何把原始字节解码为事件(部分解码器还能决定输出类型:log/metric/trace) |
tls | object | 无 | TLS 连接选项(见下文) |
auth | object | 无 | 认证策略(见下文) |
jetstream | object | 无 | 启用 NATS JetStream 模式(见下文) |
log_namespace | bool | 全局设置 | 日志命名空间覆盖项(文档中隐藏) |
subscriber_capacity的默认值65536定义于源码常量default_subscription_capacity(config.rs)。该值会通过ConnectOptions::subscription_capacity传递给底层客户端(config.rs),影响背压行为——当管道下游处理不过来时,消息在订阅缓冲区内排队。
一个最小可用配置
以下是源码中GenerateConfig生成的标准示例(config.rs):
sources: 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()返回true(config.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: passwordToken
auth: strategy: token token: value: my-token凭证文件(JWT 体系)
auth: strategy: credentials_file credentials_file: path: /etc/nats/nats.credsNKey
auth: strategy: nkey nkey: nkey: UC4... # 相当于公钥 seed: SUAA... # 相当于私钥种子TLS
tls配置支持标准 TLS 字段(enabled、ca_file、crt_file、key_file、verify_certificate、verify_hostname等)。从源码可确认的实现细节(src/nats.rs):
- 未启用 TLS 时直接返回明文连接;
- 启用后可通过
ca_file添加根证书,通过crt_file+key_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 | 是 | string | NATS 消息的原始载荷文本 |
source_type | 是 | string | 来源类型名称,固定为nats |
subject | 是 | string | 消息来源的 NATS subject |
timestamp | 是 | timestamp | 事件时间戳 |
其中subject字段名可通过subject_key_field参数自定义(默认subject)。源码在process_message中负责注入这些元数据(source.rs):insert_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.rs):
- Core 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),仅供参考