ARTICLE DETAIL

资讯详情

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

Vector 原生互联:深入解析 vector source 组件的 gRPC 事件接收、配置与实现原理

Vector 原生互联:深入解析 vector source 组件的 gRPC 事件接收、配置与实现原理 Vector 原生互联深入解析 vector source 组件的 gRPC 事件接收、配置与实现原理【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vectorvector source是 Vector 项目中用于接收另一台上游Vector 实例通过vectorsink 推送的观测数据的接入组件。它基于 gRPC/HTTP 协议工作以透传方式原样转发日志、指标与 trace是构建 Vector 级联cascading与聚合aggregation拓扑的核心构件。读完本文你将掌握该组件的完整配置项、输出 Schema、确认acknowledgements机制、TLS/保活/压缩处理以及其底层 gRPC 服务实现与健康检查细节可以直接上手搭建多实例数据管道。组件定位为聚合器而生的内部互联端口vectorsource 的唯一职责是接收来自另一台 Vector 实例的数据。在组件元数据 website/cue/reference/components/sources/vector.cue 中它的定位描述为Receives data from another upstream Vector instance using the Vector sink.即通过 vector sink 接收来自上游 Vector 实例的数据。它与此前可能接触过的http、socket等通用接收组件不同vectorsource 与vectorsink 配对使用走的是 Vector 自有的 protobuf 协议而不是通用的 JSON/文本编码因此接收端无需再做解码与字段映射事件到达后几乎零成本进入下游管道。从 CUE 元数据可以提取出该组件的完整体检表维度值说明交付语义deliveryat_least_once配合确认机制实现至少一次投递部署角色deployment_rolesaggregator面向聚合/中转场景设计开发状态developmentstable已进入稳定阶段出口方式egress_methodstream事件以流式方式进入管道是否有状态statefulfalse无本地持久化状态确认机制acknowledgements支持可向上游 sink 返回端到端确认多行聚合multiline不支持无多行合并能力协议protocolshttpgRPC监听传入连接SSL可选可按需开启 TLS源码中的#[configurable_component(source(vector, ...))]注解见 src/sources/vector/mod.rs将其注册为名为vector的 source 组件并声明can_acknowledge() - true表明该组件具备端到端确认能力。工作原理一条 gRPC 服务背后的完整数据链路协议定义PushEvents 与 HealthCheck 两个 RPCvectorsource 的线上协议定义在 proto/vector/vector.protosyntax proto3; package vector; import event.proto; message PushEventsRequest { repeated event.EventWrapper events 1; } message PushEventsResponse {} enum ServingStatus { SERVING 0; NOT_SERVING 1; } message HealthCheckRequest {} message HealthCheckResponse { ServingStatus status 1; } service Vector { rpc PushEvents(PushEventsRequest) returns (PushEventsResponse) {} rpc HealthCheck(HealthCheckRequest) returns (HealthCheckResponse); }PushEvents批量推送事件EventWrapper列表一次调用可携带多条日志/指标/traceHealthCheck自定义健康检查配合标准 gRPC health 协议使用详见下文健康检查一节。服务端实现接收、打标、入管道在 src/sources/vector/mod.rs 中Service::push_eventsL54-L95是核心处理函数链路如下反序列化把PushEventsRequest中的EventWrapper逐个Event::from转为内部事件对象注入标准元数据调用self.log_namespace.insert_standard_vector_source_metadata(log, VectorConfig::NAME, now)根据当前日志命名空间为每条日志附加source_type与时间戳元数据统计注册并发射EventsReceived内部事件CountByteSize(count, byte_size)字节数按estimated_json_encoded_size_of()估算确认若启用了确认机制通过BatchNotifier::maybe_apply_to挂上批量确认接收器投递self.pipeline.clone().send_batch(events)把整批事件送入下游管道管道关闭时报StreamClosedError并返回 gRPCunavailable回执根据handle_batch_statusL110-L121把下游投递结果翻译为 gRPC 状态——Delivered返回OkErrored返回Status::internal(Delivery error)Rejected返回Status::data_loss(Delivery failed)。关键设计点事件本体是透传的。元数据 website/cue/reference/components/sources/vector.cue 的输出定义中明确写道Vector transparently forwards data from another upstream Vector instance. Thevectorsource will not modify or add fields.即上游事件中的业务字段不会被改动source 只负责附加来源相关的标准元数据命名空间差异见下文输出 Schema。服务装配Vector 服务 gRPC 健康服务共用同一端口buildL183-L230中通过RoutesBuilder把两个服务注册到同一监听地址自定义的vector.Vector服务tonic_health的标准健康检查服务。同时用max_decoding_message_size(max_decompressed_size_bytes())限制了单条消息解码上限避免超大消息在无认证的监听端口上造成无界内存分配。快速上手最小配置与进阶配置仓库生成的示例配置可以直接套用。最小配置见 website/generated/example-configs/sources/vector/minimal.yamlsources: my_source_id: type: vector address: 0.0.0.0:6000进阶配置website/generated/example-configs/sources/vector/advanced.yaml额外声明了协议版本sources: my_source_id: type: vector address: 0.0.0.0:6000 version: 2配置保存在config/vector.yaml的sources段下即可通过vector validate --config-toml config/vector.yaml校验随后vector --config config/vector.yaml启动。配置项详解VectorConfig的结构体定义在 src/sources/vector/mod.rsL127-L149全部字段如下配置项类型默认值说明addressSocketAddr0.0.0.0:6000监听地址必须包含端口无认证的 gRPC 监听入口version枚举2未设置协议版本标记当前仅2一个合法值tlsTlsEnableableConfig未启用可选 TLS 配置见下文acknowledgementsbool或结构体跟随全局配置端到端确认开关支持true或{enabled, when_full}两种写法keepaliveGrpcKeepaliveConfig均未设置gRPC 连接保活/老化控制log_namespacebool未设置日志命名空间覆盖全局设置文档中隐藏字段address必须显式包含端口address是核心必填项源码中类型为std::net::SocketAddr注释明确要求Itmustinclude a port。默认值由Default实现提供L161-L172为0.0.0.0:6000同时resources()L245-L247返回Resource::tcp(self.address)用于拓扑层面的端口占用声明。配置时写成0.0.0.0:6000、127.0.0.1:6000或[::]:6000均可只要带上端口。version当前仅支持 2VectorConfigVersion枚举只定义了V2一个成员serde(rename 2)表示配置中写作version: 2。它是配置兼容性的版本标记日常使用通常不需要显式填写。tls可选的加密传输CUE 元数据显示该组件tls.enabled: true、enabled_default: false默认关闭支持证书校验can_verify_certificate: true。build中通过MaybeTlsSettings::from_config(self.tls.as_ref(), true)将配置转换为实际监听层。启用示例sources: vector_in: type: vector address: 0.0.0.0:6000 tls: enabled: true crt_file: /etc/vector/tls/server.crt key_file: /etc/vector/tls/server.key ca_file: /etc/vector/tls/ca.crt verify_certificate: false对应地上游vectorsink 侧配置tls.verify_certificate/tls.ca_file等即可完成双向校验sink 侧 TLS 能力见 website/cue/reference/components/sinks/vector.cue。acknowledgements端到端至少一次该字段使用serde(default, deserialize_with bool_or_struct)反序列化因此既可以直接写布尔值sources: vector_in: type: vector address: 0.0.0.0:6000 acknowledgements: true也可以写成带when_full策略的结构体形式。构建时cx.do_acknowledgements(self.acknowledgements)会把组件级配置与全局确认设置合并。启用后上游 sink 的批次只有在整批事件被下游管道成功消费后才会收到成功回执配合at_least_once交付语义实现端到端可靠性。keepalive连接生命周期控制keepalive对应GrpcKeepaliveConfig定义在 src/sources/util/grpc/mod.rsL73-L91字段单位默认值说明max_connection_age_secs秒未设置不限龄连接最长存活时间到期后服务端关闭连接max_connection_age_grace_secs秒未设置在max_connection_age_secs基础上的宽限期仅在设置了前者时生效示例sources: vector_in: type: vector address: 0.0.0.0:6000 keepalive: max_connection_age_secs: 300 max_connection_age_grace_secs: 30实现上由MaxConnectionAgeIo包装 TCP 流src/sources/util/grpc/mod.rs L107-L127到达 deadline 后读操作返回 EOF写操作则需等到该连接上没有活动请求active_requests 0才返回BrokenPipe确保正在处理中的 RPC 不被粗暴打断。源码中的config_keepalive、max_connection_age_closes_idle_connection、max_connection_age_allows_client_reconnect等测试src/sources/vector/mod.rs验证了上述行为。log_namespace隐藏的命名空间覆盖项log_namespace是标记了docs::hidden的隐藏字段类型为Optionbool用于覆盖全局日志命名空间设置影响下文输出 Schema中元数据字段的存放位置。输出 Schema日志/指标/trace 全透传outputs()L232-L243通过NativeDeserializerConfig.schema_definition(...)生成输出定义并声明DataType::all_bits()——即该组件同时输出日志、指标与 trace三种数据类型。元数据 website/cue/reference/components/sources/vector.cue 中的输出说明日志事件字段透传另附加source_type必填示例值vector指标counter / distribution / gauge / histogram / set 全部透传统一附加source_type标签trace透传上游 trace 事件。source_type的取值与命名空间有关。源码中的两个 Schema 测试给出了精确答案Vector 命名空间output_schema_definition_vector_namespace元数据字段为vector.source_type字节类型与vector.ingest_timestamp时间戳类型Legacy 命名空间output_schema_definition_legacy_namespace事件字段为source_type字节类型与timestamp时间戳类型。也就是说不改动业务字段指的是事件负载本身而来源标记和时间戳这类标准元数据始终会被附加只是存放位置随命名空间不同而不同。压缩与解压gzip / zstd / identityvectorsource 的压缩协商由DecompressionAndMetricsLayer统一处理见 src/sources/util/grpc/decompression.rs。该层通过grpc-accept-encoding头向客户端声明支持gzip、zstd、identity三种编码并在收到请求后按grpc-encoding头解压载荷。几个实现细节值得注意gRPC 每条消息都有 5 字节头1 字节压缩标志 4 字节长度前缀见 src/sources/util/grpc/decompression.rs 的GRPC_MESSAGE_HEADER_LEN解压有全局字节数上限约束max_decompressed_size_bytes压缩帧还会经过线级预过滤防止单一超大消息打爆内存由于解压与字节统计被集中到该层source 代码中刻意不再调用accept_compressed避免职责重复。这一能力与vectorsink 的压缩选项一一对应。sink 侧生成的进阶配置website/generated/example-configs/sinks/vector/advanced.yaml为sinks: my_sink_id: type: vector inputs: - my-source-or-transform-id address: http://127.0.0.1:6000 compression: nonecompression可取none、gzip、zstd。源码中以source sink 互操作方式验证了三种模式的端到端数据一致性见 src/sources/vector/mod.rs 的receive_message、receive_gzip_compressed_message、receive_zstd_compressed_message测试并断言 100 条随机事件经 sink 发送、source 接收后逐条相等assert_event_data_eq。健康检查自定义 RPC 标准 gRPC health 双通道监听端口上同时提供两套健康检查自定义HealthCheckRPCService::health_checkL98-L107恒定返回ServingStatus::Serving源码注释保留 TODO未来接入真实健康状态标准 gRPC health 协议tonic_health的health_reporter注册了vector.Vector服务的 serving 状态空服务名的聚合健康检查同样可用。源码中的custom_health_check_works与standard_grpc_health_check_works测试分别验证了两条路径。这意味着负载均衡器、K8s probe 或运维脚本既可以用grpc_health_probe这类标准工具探测也可以直接调用 Vector 自定义的HealthCheckRPC。与 vector sink 配对搭建级联/聚合拓扑vectorsource 的典型使用场景是聚合器拓扑多个边缘 Vector 用vectorsink 指向中心聚合实例的vectorsource。配对关系在 sink 侧元数据 website/cue/reference/components/sinks/vector.cue 中描述为 Sends data to another downstream Vector instance via the Vector source。一个完整的级联示例# 中心聚合器接收端 sources: vector_in: type: vector address: 0.0.0.0:6000 acknowledgements: true sinks: long_term_storage: type: clickhouse inputs: [vector_in] ... # 边缘节点发送端 sources: app_logs: type: file include: [/var/log/app/*.log] sinks: to_aggregator: type: vector inputs: [app_logs] address: http://aggregator-host:6000 compression: gzip配合说明端口约定接收端 source 的address必须与发送端 sink 的address端口一致如上面的6000压缩协同sink 按compression压缩source 自动按grpc-encoding解压两端无需显式对齐编码清单可靠性source 声明at_least_once并支持确认sink 侧通过批次回执获得端到端投递状态适合对丢数敏感的场景。可观测性source 自身的 gRPC 指标vectorsource 内置的遥测指标由元数据声明见 website/cue/reference/components/sources/vector.cue 的telemetry段经由build_grpc_trace_layersrc/sources/util/grpc/mod.rs L440-L477在请求/响应路径上发射指标名含义grpc_server_handler_duration_secondsgRPC 处理器耗时直方图grpc_server_messages_received_total已接收 gRPC 消息计数grpc_server_messages_sent_total已发送 gRPC 消息计数配合EventsReceived每条消息的条数与字节数可以在运维侧同时观测连接层面的 RPC 数量与事件层面的数据量两个维度的健康状况。测试验证仓库给出的行为契约组件自带丰富测试可作为行为契约参考src/sources/vector/mod.rsgenerate_config配置可被正确生成/反序列化config_keepalive验证 keepalive 配置解析max_connection_age_closes_idle_connection连接超过max_connection_age_secs后被关闭读返回 EOFoutput_schema_definition_vector_namespace/output_schema_definition_legacy_namespace两种命名空间下的 Schema 字段契约receive_message/receive_gzip_compressed_message/receive_zstd_compressed_message与vectorsink 的端到端互操作含压缩模式custom_health_check_works/standard_grpc_health_check_works两套健康检查路径max_connection_age_allows_client_reconnect连接老化后客户端能自动重连。注意事项与适用边界监听端口无内置认证vectorsource 的 gRPC 端口本身不带鉴权生产环境建议通过防火墙、独立网络、TLS 或 K8s NetworkPolicy 限制访问范围仅面向 Vector 实例它是为vectorsink 配套设计的私有协议入口不适合作为通用 HTTP/JSON 摄取端点那应选用http_server等组件默认值口径源码Default实现的监听地址为0.0.0.0:6000生成的示例配置亦采用该端口文档元数据中另有 9000 的参考端口标注实际部署以address显式配置为准多行聚合不可用组件不支持multiline配置多行合并请在发送端如filesource 或对应 transform完成。综上vectorsource 是 Vector 集群化部署中实例到实例传输的关键一环它用自有的 protobuf/gRPC 协议做到事件零损耗透传通过确认机制支撑端到端可靠性并以简洁的配置项覆盖地址、版本、TLS、保活与命名空间等核心诉求。结合本文给出的源码路径与配置示例你可以在自己的拓扑中直接复制、验证并扩展这套级联方案。【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表