
手写一个 Vector Event Stream Sink从配置结构体到内部事件的全流程实现指南【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector导读本文以 Vector 官方教程 docs/tutorials/sinks/1_basic_sink.md 为骨架完整讲解如何从零编写一个event stream风格即流式风格的 Sinkbasic它把收到的每个事件以 Debug 格式打印到标准输出。你将逐步掌握配置结构体与configurable_component宏、SinkConfig/StreamSink两个核心 trait 的实现、Cargo.toml特性开关的接入、端到端确认acknowledgements以及BytesSent/EventsSent内部可观测性事件的埋点并了解如何用vdev或裸cargo run立刻运行起来。文中所有关键步骤都会对照仓库真实源码给出佐证让你照着写的同时理解 Vector 底层如何驱动你的 Sink。背景两种 Sink 风格为什么要写 Event StreamVector 目前存在两种风格的 Sinkevent风格与event streams风格。旧式 event 风格基于futures::Sink已经被标记为废弃deprecated但当前仍有相当一部分 Sink 使用该风格开发社区正在逐步迁移跟踪进度见 event streams 迁移 issue对应仓库 issue 9261。这一点在 VectorSink 的源码中有直接体现——lib/vector-core/src/sink.rs 中from_event_sink转换函数被标注了#[deprecated]注释明确写着 Deprecated in favor ofVectorSink::from_event_streamsink. See [vector/9261]#[deprecated] pub fn from_event_sink(sink: impl SinkEvent, Error () Send Unpin static) - Self { VectorSink::Sink(Box::new(EventSink::new(sink))) }而新的 event streams 风格基于StreamSinkEventtrait本教程只讲解这一种风格。本教程最终产出的basicsink 会打印每个事件的 Debug 表示属于最纯粹的教学样例如果你希望看一个生产环境中打印到控制台的同类实现可以参考仓库里现成的 console sink 实现下文会对照讲解。第一步创建模块与模块级文档注释在src/sinks/下新建一个 Rust 模块文件basic.rs并写上模块级注释说明这个 Sink 是做什么的。Vector 要求每个组件都有清晰的目的说明这份注释也会进入后续自动生成的组件文档//! Basic sink. //! A sink that will send its output to standard out for pedagogical purposes.第二步导入所需符号教程使用的导入分两类一类是 Vector 为所有流式 Sink 准备的 prelude 集合另一类是内部事件相关的符号。prelude 定义在 src/sinks/prelude.rs它会统一 re-exportasync_trait、futures的StreamExt/BoxStream、tower的Service/ServiceBuilder以及vector_lib中的StreamSink、VectorSink、EventFinalizers、EventStatus、Finalizable、CountByteSize等大量常用符号省去逐个手动导入的繁琐use crate::sinks::prelude::*; use vector_lib::internal_event::{ ByteSize, BytesSent, EventsSent, InternalEventHandle, Output, Protocol, };注意internal_event中ByteSize、BytesSent、EventsSent等类型从vector_lib导出仓库中 console sink 的导入写法与教程完全一致它们将用于后面发射内部事件小节。第三步配置结构体与configurable_component宏Sink 开发的第一步是创建一个代表该 Sink 配置的结构体。Vector 启动时读入的 YAML/TOML 配置文件会被反序列化到这个结构体的字段上从而让用户自定义 Sink 行为#[configurable_component(sink(basic))] #[derive(Clone, Debug)] /// A basic sink that dumps its output to stdout. pub struct BasicConfig { #[serde( default, deserialize_with crate::serde::bool_or_struct, skip_serializing_if crate::serde::is_default )] pub acknowledgements: AcknowledgementsConfig, }几个要点#[configurable_component(sink(basic))]这是 Vector 的 config 宏用来让配置结构体具备可配置组件configurable component的能力。宏展开时会同时挂上serde反序列化、typetag动态注册下文 SinkConfig 会用到以及文档生成所需的元数据。该宏的实现位于 lib/vector-config-macros/src/configurable_component.rs从注释可以看出#[configurable_component(source)]/sink(...)这样的写法会被转换为更具体的组件宏并携带组件类型信息供文档工具链使用。必须在结构体上方提供 doc 注释否则 Vector 编译会失败——因为文档是依赖这些注释自动生成的。acknowledgements字段类型为AcknowledgementsConfig用于配置该 Sink 的端到端确认end-to-end acknowledgements能力——即 Sink 成功投递事件后通知上游 source 的能力。通过deserialize_with crate::serde::bool_or_struct可以让用户在配置里直接写acknowledgements: true布尔简写或完整的对象结构bool-or-struct 反序列化default与skip_serializing_if则保证缺省与序列化行为正确。接着为配置结构体实现GenerateConfigtrait。这个 trait 被vector generate命令用来生成该 Sink 的默认配置骨架impl GenerateConfig for BasicConfig { fn generate_config() - serde_json::Value { serde_yaml::from_str({}).unwrap() } }因为basic没有任何必填字段所以生成的默认配置就是一个空对象{}。真实组件同样如此例如 console config 的GenerateConfig实现会构造一个带默认值的完整结构体并序列化为 JSON每个GenerateConfig实现还配有一个test_generate_config测试辅助函数验证其可用性见同一文件的测试模块。第四步实现SinkConfigtraitSinkConfig是 Vector 用来从配置生成可运行 Sink 的入口 trait通过typetag实现多态注册#[async_trait::async_trait] #[typetag::serde(name basic)] impl SinkConfig for BasicConfig { async fn build(self, _cx: SinkContext) - crate::Result(VectorSink, Healthcheck) { let healthcheck Box::pin(async move { Ok(()) }); let sink VectorSink::from_event_streamsink(BasicSink); Ok((sink, healthcheck)) } fn input(self) - Input { Input::log() } fn acknowledgements(self) - AcknowledgementsConfig { self.acknowledgements } }务必注意#[typetag::serde(name basic)]中给出的字符串必须与上面configurable_component(sink(basic))的组件名一致否则配置解析时无法通过名称找到该 Sink。trait 的三个方法职责如下build异步构建 Sink 的两个组件——健康检查healthcheck与真正的数据通路。这里返回的Healthcheck是一个 boxed futureVector 会在启动时轮询它以确认下游目标可用。本例只是输出到控制台假设必然可用因此直接返回Ok(())。input声明该 Sink 接受的输入类型Input::log()表示只接受日志事件log event。Vector 会据此在拓扑校验阶段检查上游组件是否与其类型匹配。acknowledgements返回该 Sink 的确认配置引用Vector 用它决定端到端确认在拓扑中的传播行为。build中VectorSink::from_event_streamsink(BasicSink)这行是整个转换的关键。在 lib/vector-core/src/sink.rs 中可以看到VectorSink是一个两变体枚举pub enum VectorSink { Sink(Boxdyn SinkEventArray, Error () Send Unpin), Stream(Boxdyn StreamSinkEventArray Send), }from_event_streamsinklib/vector-core/src/sink.rs接受一个实现了StreamSinkEvent的对象包装成内部EventStream后放进VectorSink::Stream变体而VectorSink::runlib/vector-core/src/sink.rs会调用StreamSink的run方法驱动数据流动pub fn from_event_streamsink(sink: impl StreamSinkEvent Send static) - Self { let sink Box::new(sink); VectorSink::Stream(Box::new(EventStream { sink })) }第五步实现BasicSink与StreamSinktrait现在实现真正干活的部分。BasicSink非常简单不需要任何字段struct BasicSink;为了让 Vector 驱动它需要实现StreamSinkEventtrait。该 trait 定义在 lib/vector-core/src/sink.rs只有一个异步方法run#[async_trait::async_trait] pub trait StreamSinkT { async fn run(self: BoxSelf, input: stream::BoxStream_, T) - Result(), (); }run的核心参数input是一个BoxStream_, Event即正在流向该 Sink 的事件流Sink 要做的就是从这条流中不断拉取事件并投递到目的地。由于async_trait宏处理生命周期时存在限制教程采用了一个惯用做法run只做转发真正的逻辑放在直接实现在BasicSink上的run_inner方法里#[async_trait::async_trait] impl StreamSinkEvent for BasicSink { async fn run( self: BoxSelf, input: futures_util::stream::BoxStream_, Event, ) - Result(), () { self.run_inner(input).await } }run_inner的逻辑极简逐条拉取事件并打印 Debug 表示直到流结束impl BasicSink { async fn run_inner(self: BoxSelf, mut input: BoxStream_, Event) - Result(), () { while let Some(event) input.next().await { println!({:?}, event); } Ok(()) } }这里input.next().await来自 prelude 中 re-export 的futures::StreamExt。此时你的 Sink 已经可以工作了——只是还没有接入 Vector 的注册机制也还没有处理确认与可观测性。第六步把 Sink 接入 Vector特性开关与注册每个 Sink 都被放在一个独立的 feature flag 后面这样用户可以只编译自己需要的组件避免不必要的编译负担。需要在根 Cargo.toml 中添加sinks-basic特性把它排进字母序sinks-azure_blob [dep:azure_core, dep:azure_identity, dep:azure_storage, dep:azure_storage_blobs] sinks-azure_logs_ingestion [dep:azure_core, dep:azure_identity, dep:azure_storage_blob] sinks-basic [] sinks-blackhole [] sinks-chronicle []然后在sinks-logs聚合特性中加入它。sinks-logs是编译日志类 Sink 的总开关Cargo.toml 中可以看到它聚合了 amqp、aws_cloudwatch_logs、kafka、loki 等全部日志 Sinkbasic接受 log 输入因此属于这一组sinks-logs [ sinks-amqp, sinks-apex, sinks-aws_cloudwatch_logs, sinks-aws_kinesis_firehose, sinks-aws_kinesis_streams, sinks-aws_s3, sinks-aws_sqs, sinks-axiom, sinks-azure_blob, sinks-azure_logs_ingestion, sinks-basic, sinks-blackhole, sinks-chronicle,从仓库现状看sinks-blackhole []Cargo.toml与sinks-console []Cargo.toml这类无第三方依赖的 Sink 特性定义与教程新增的sinks-basic []完全同构可以照此对比校验自己的写法。特性注册完成后还需要在src/sinks/的模块树中按#[cfg(feature sinks-basic)]声明basic模块同时把BasicConfig注册进组件枚举这一步是让配置文件里type: basic能被解析到BasicConfig的关键。第七步端到端确认Acknowledgements当 Sink 处理完一个事件后必须确认它这样确认信息才能沿拓扑回传给上游 source。Vector 的事件携带 finalizer终结器通过更新其状态来表达投递结果。需要修改run_innerasync fn run_inner(self: BoxSelf, mut input: BoxStream_, Event) - Result(), () { - while let Some(event) input.next().await { while let Some(mut event) input.next().await { println!({:#?}, event); let finalizers event.take_finalizers(); finalizers.update_status(EventStatus::Delivered); } Ok(()) }改动包含两层含义mut event把event声明为可变才能取出并修改其内部状态。event.take_finalizers()从事件中取出EventFinalizers集合Finalizabletrait 提供见 src/sinks/prelude.rs 的 re-export然后调用update_status(EventStatus::Delivered)标记投递成功。Vector 会据此向上游确认该事件已送达。EventStatus一共有三种状态选择依据是错误性质EventStatus::Delivered事件成功送达。本例打印成功即视为送达。EventStatus::Errored投递过程中发生错误但并非永久性错误。Vector 会尝试重新投递该事件。EventStatus::Rejected发生了无论重试多少次都必然失败的永久性错误事件被最终拒绝。真实的 console sink 展示了一个更完整的处理范例它先取出 finalizer编码失败时置Errored写入 stdout/stderr 失败时记录错误日志、置Errored并停止 Sink只有写入成功才置Delivered并继续发射指标。第八步发射内部事件可观测性Vector 强调自身可观测每个组件都应发射描述自身运行状态的内部事件方便用户判断整体健康状况。你的 Sink 在投递事件后必须更新已投递事件数等指标。这里需要发射两种事件。BytesSent统计向下游发送的字节数BytesSent记录 Sink 向下游发送了多少字节。先算出要发送的字节数再发射事件async fn run_inner(self: BoxSelf, mut input: BoxStream_, Event) - Result(), () { let bytes_sent register!(BytesSent::from(Protocol(console.into(),))); while let Some(mut event) input.next().await { let bytes format!({:#?}, event); println!({}, bytes); - println!({:#?}, event); bytes_sent.emit(ByteSize(bytes.len())); let finalizers event.take_finalizers(); finalizers.update_status(EventStatus::Delivered); } Ok(()) }要点register!宏把BytesSent指标注册到全局指标注册表Protocol(console)标记传输协议标签ByteSize(bytes.len())携带本次发送的字节数。这里特意把println!({:#?})改为先format!成字符串再println!({})为的是能精确拿到bytes.len()。EventsSent统计发送的事件数与编码后大小EventsSent由 Vector 的每个组件发射用于统计向下游发送的事件数量并附带每个事件的编码后估算大小async fn run_inner(self: BoxSelf, mut input: BoxStream_, Event) - Result(), () { let bytes_sent register!(BytesSent::from(Protocol(console.into(),))); let events_sent register!(EventsSent::from(Output(None))); while let Some(mut event) input.next().await { let bytes format!({:#?}, event); println!({}, bytes); bytes_sent.emit(ByteSize(bytes.len())); let event_byte_size event.estimated_json_encoded_size_of(); events_sent.emit(CountByteSize(1, event_byte_size)); let finalizers event.take_finalizers(); finalizers.update_status(EventStatus::Delivered); } Ok(()) }Output(None)标记该指标对应的输出CountByteSize(1, event_byte_size)同时携带事件个数1与estimated_json_encoded_size_of()估算出的 JSON 编码后大小——该 trait 方法来自vector_lib::EstimatedJsonEncodedSizeOf同样由 prelude re-export。真实的 console sink 里两处发射逻辑与教程完全同构bytes_sent.emit(ByteSize(bytes.len()))与events_sent.emit(CountByteSize(1, event_byte_size))可以互相印证。关于 Vector 内部可观测性的完整规范见仓库的 instrumentation 规范文档。第九步运行你的 Sink写一个 Vector 配置文件./basic.yml把标准输入 source 直接接到basicsinksources: stdin: type: stdin sinks: basic: type: basic inputs: - stdin这里type: basic正是typetag注册名与组件名inputs: [stdin]声明数据来源。方式一使用 vdev推荐vdev是 Vector 提供的构建辅助工具代码在 vdev/使用说明见 vdev/README.md它能自动解析配置文件、识别所需特性并构建运行cargo vdev run ./basic.ymlvdev会读取配置文件中的组件自动打开sources-stdin、sinks-basic等对应特性进行编译省去手动拼特性的麻烦。方式二裸 cargo run不使用vdev时需要显式指定特性sources-stdin特性在 Cargo.toml 中定义cargo run --no-default-features --features sources-stdin, sinks-basic -- -c ./basic.yml启动后在终端输入一行文本并回车Vector 就会打印出对应 log 事件的 Debug 信息——恭喜你的第一个 Sink 跑通了对照真实实现从教程到生产的三个建议教程的basic是教学骨架仓库中功能最接近的成熟组件是consolesink对比它可以在三个维度上把玩具推向可用编码encoding真实 sink 通常不直接println!而是用 Encoder 配置如 JSON serializer newline delimiter统一处理编码SinkType::StreamBased指明流式类型。健康检查console 的build返回future::ok(()).boxed()src/sinks/console/config.rs与教程的Box::pin(async move { Ok(()) })等价但更简洁。合规测试console 提供了组件规范合规测试component_spec_compliancesrc/sinks/console/sink.rs通过test_util::components::run_and_assert_sink_compliance把VectorSink::from_event_streamsink(...)包起来喂入单事件流验证 Sink 行为符合 Vector 组件规范。为自己的新 Sink 添加类似的测试是保证其长期健康的最佳实践。如果你要写一个真正面向网络的 Sink例如 HTTP建议直接学习教程的下一篇 2_http_sink.md它会引入tower服务、批处理与重试等生产级机制。【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考