Vector 端到端确认机制(End-to-end Acknowledgements)架构详解:从 RFC 6517 到生产实现
Vector 端到端确认机制End-to-end Acknowledgements架构详解从 RFC 6517 到生产实现【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector端到端确认End-to-end Acknowledgement是 Vector 保证数据管线成功投递全部事件的核心机制事件只有在被所有目标 sink 成功投递或经历不可恢复的失败之后源头才向产生数据的客户端返回确认。本文以仓库内的设计文档 rfcs/2021-03-26-6517-end-to-end-acknowledgement.md 为主体结合 lib/vector-common/src/finalization.rs、lib/vector-core/src/event/metadata.rs、lib/vector-core/src/source_sender/sender.rs 等落地实现讲清该机制的功能需求、数据模型、边界场景fan-out、merge、丢事件、磁盘缓冲以及acknowledgements/authoritative/buffer.acknowledgements三个配置项的正确用法。为什么需要端到端确认用户的诉求很朴素希望得到保证交给数据管线的每一条事件都确实被投递成功。为此 Vector 需要为每条事件维护与其源头source的关联信息使源头能够在正确的时机向协议对端确认那些进入管线的事件。文档把事件完成整段管线旅程的时刻称为finalization终结化——它可能由三种原因触发成功投递、永久性失败、或在处理过程中被丢弃。RFC 的 Scope 明确了边界本文只讨论功能需求与数据模型变更及其潜在影响不涉及单个 source / transform / sink 为了支持该特性所需的细节改动。六类边界场景为什么朴素实现行不通Vector 的拓扑中天然存在 fan-out、merge、drop 等行为直接事件送达 sink 即确认的朴素方案会出错。RFC 逐一枚举了这些场景这也是后续数据模型设计的出发点。单 source 线性经过 transforms 到单 sink最简单的场景。事件在管线中仅被搬运 修改每个事件需要携带一个指向来源 source 的引用与一个唯一的消息标识符。不同 source 对消息标识符的要求各不相同例如 Kafka 用 offsetHTTP 用请求批次因此标识符的生成必须交给 source 自己。多 source 汇聚到单 sink多个 source 的输出可能被 sink 打包进同一批次。该场景不需要额外的数据追踪但状态的回传必须以事件为粒度、逐条流向各自的 source而不能按批次整体回传——这直接决定了最终设计中批量通知器的结构。单 source 扇出到多 sink一条事件被多个 sink 消费时确认状态必须等所有 sink 都终结该事件后才发出。这就要求终结追踪在所有相关 sink 之间共享而不是随事件克隆并用计数器之类的机制判定全部投递完成。拓扑层参与该计数器的初始化以避免把事件交给最后一个 sink与更早的 sink 认为投递已全部完成之间的竞态条件。合并类 transform如mergemerge会把多条事件合并成一条新事件因此新事件的追踪信息也必须合并携带所有参与合并的原始事件的 source 引用且该 transform 可能从多个 source 取数导致合并结果天然属于多个 source。拆分类 transform如routeroute的效果与单事件发往多 sink相同区别在于它可能把事件送向一条、多条甚至零条输出管线全部被路由丢弃的情况。丢事件类 transform如dedupe、filter、sample这类 transform 会根据条件丢弃事件。丢弃时唯一现实的处理是将事件标记为已投递。RFC 建议维护独立的 delivered已投递与 dropped已丢弃状态位不过这对 source 的确认行为没有实质差别。用户态 transformLua 或 WASM用户脚本可能复现上述任意行为甚至凭空创建与源事件毫无关系的新事件。Vector 无法可靠地为脚本合并/丢弃事件自动打标因此把这一行为交给脚本作者Vector 只提供便捷库函数协助脚本处理常见的终结化任务。非确认型 source 与磁盘缓冲 sink非确认型 sourceapache_metrics、prometheus_parser、stdin等 source 在协议层面无法提供确认用户也可能只想对部分 source 选择性开启该特性。为最小化开销不能或不应确认的 source不必参与终结化处理否则只会白白丢弃状态。磁盘缓冲 sink确认可能走两条路径——若事件落盘即视为安全则事件入缓冲后即可终结缓冲层负责终结化序列化前剥离终结数据并改造为只在事件被投递后才从缓冲中删除reload 时会为事件挂上新的终结追踪。若必须等真实投递则终结状态需在事件从缓冲读回后确认。此时 source 引用必须序列化成字符串/整数形式但配置 reload 可能改名、Vector 重启会诞生同名而不同的 source所以持久化的 source 标识符必须在 reload 之间保持稳定、又在每次 Vector 运行之间彼此不同——这是磁盘缓冲场景最难啃的点。内部提案事件终结化Event Finalization提案由两部分组成事件终结化与源之间的通信机制以及承载事件状态追踪的事件元数据。通道机制与三步终结流程拓扑实例化时会给每个 source 分配一个唯一标识符 token。source 用该标识符创建一条或多条通知通道各自带唯一标识从而同时获得谁需要收到终结状态的标识与如何送达事件标识符和状态的机制。终结化是三步流程sink 完成事件投递后把投递状态记录到所有事件克隆共享的 finalizer中状态从初始的Dropped提升为Delivered、Errored或Failed。Recorded状态一旦写入就不会再变。若某个 sink 被配置为authoritative权威它会立即把所有源批次的状态更新为Recorded防止后续多余的状态更新否则由事件最后一个副本在共享 finalizer 被 drop 时完成这次更新。当批次中最后一条事件终结后批次的最终状态只发送一次给 source经由 one-shot channelsource 据此确认整个批次。事件元数据加到事件元数据上的结构共三层每个事件有可选的 finalizer。事件被扇出到多个 transform/sink 时事件被克隆到每个目的地但各目的地共享同一个 finalizer。不做确认的 source 不会初始化 finalizer若事件后来与需要确认的事件合并finalizer 可以被补设。每个 finalizer 含事件状态标记与一个或多个共享的 source 批次通知器BatchNotifier引用。事件合并时这些 source 批次在此聚合批次中最后一条事件终结时整个批次统一确认。逐条接收事件的 source 需要为每条事件单独创建一个批次。批次通知器含一条到来源 source 的 one-shot 通道、批次当前状态与唯一标识符。标识符用于在事件序列化之后重建该通道。数据结构RFC 原始设计struct EventMetadata { // … existing fields … finalizers: Box[ArcEventFinalizer], } struct EventFinalizer { status: EventStatus, source: ArcBatchNotifier, identifier: Uuid, } struct BatchNotifier { status: MutexBatchStatus, notifier: tokio::sync::oneshot::SenderBatchStatus, identifier: Uuid, } enum BatchStatus { Delivered, Errored, Failed, } enum EventStatus { Dropped, // default status Delivered, Errored, Failed, Recorded, }从 RFC 到实现finalization.rs 中的落地形态RFC 中的设计在仓库里演化为 lib/vector-common/src/finalization.rs 中的一组类型命名与语义保持了高度一致EventFinalizersVecArcEventFinalizer的封装提供add、merge、update_status、update_sources等方法并实现Finalizable/MergeFinalizabletrait事件克隆时 finalizer 集合随Arc共享事件合并时通过merge聚合多个来源finalization.rs。EventFinalizer内部用AtomicCellEventStatus保存状态默认Dropped持有BatchNotifier。update_status通过fetch_update合并状态update_batch在把事件状态写入批次后把自身状态切换为Recorded杜绝后续再更新其Drop实现会自动调用update_batch——这正是引用计数保证事件不会不给出状态就逃逸的机制finalization.rs。BatchNotifier内部是ArcOwnedBatchNotifier含AtomicCellBatchStatus默认Delivered与Optiononeshot::SenderBatchStatus。OwnedBatchNotifier的Drop会通过send_status把最终状态送回 sourceBatchStatusReceiver是对 one-shot receiver 的Future封装receiver 在发送前被 drop 时返回Erroredfinalization.rs。状态合并规则EventStatus::update中Recorded优先级最高且不可再变Rejected覆盖Errored/DeliveredErrored覆盖Delivered向Dropped更新是非法的debug 构建会断言——实现里还把 RFC 的Failed演化成了语义更精确的Rejectedfinalization.rs。该文件自带完整单元测试finalization.rs例如sends_notification验证 drop finalizer 后批次以Delivered送达clone_events验证克隆的两个 finalizer 只有全部drop 后才发送通知multi_event_batch验证同一批次的多条事件全部终结后才确认event_status_updates逐对验证了全部状态转移矩阵。事件侧EventMetadata 与 SourceSenderfinalizers 挂在 EventMetadata 上在 lib/vector-core/src/event/metadata.rs 中Inner结构体通过finalizers: EventFinalizers字段把终结器集合挂到每条事件的元数据上并提供with_finalizer、with_batch_notifier、update_status、update_sources、take_finalizers、merge_finalizers等操作入口merge方法会在合并事件时把两边的 finalizer 一并合并metadata.rs。注意EventMetadata本身是ArcInner的包装克隆元数据是廉价指针拷贝这也契合克隆共享、写时复制的整体设计。SourceSender源头发送的统一出口lib/vector-core/src/source_sender/sender.rs 是 source 向管线发事件的通道。它实现AddBatchNotifier把共享的BatchNotifier挂到整个EventArray上与Finalizable提供take_finalizers/take_finalizer_groups供缓冲层取走并重组终结器。测试用的new_test_finalize(status)与new_test_errors(error_at)展示了在没有真实 sink 的测试管线中如何手动调用metadata.update_status(...)metadata.update_sources()模拟终结可作为理解终结流程的最小示例sender.rs。拓扑与配置acknowledgements、authoritative、buffer.acknowledgementsRFC 的 Doc-level Proposal 在仓库中已成为三处可配置项下面给出完整配置语义。全局与 source 级acknowledgements在配置的全局层与每个 source上新增acknowledgements布尔选项控制该 source 是否参与端到端确认默认true。在实现中它对应 lib/vector-core/src/config/mod.rs 的SourceAcknowledgementsConfigOptionboolNone时回退到全局值并经由GlobalOptions.acknowledgements合并。一个重要的实现细节是source 级配置最终由 sink 侧反向传播决定——src/config/mod.rs 的propagate_acknowledgements会遍历所有启用了确认的 sink把sink_acknowledgements沿拓扑反向往上游 source 传播从而让只有真实接了确认型 sink 的 source 才产生额外开销。sink 级authoritativeRFC 提出在 sink 层新增authoritative布尔设置默认false标记权威 sink事件送达权威 sink 时立即向该事件的所有 source 发送终结状态并将事件状态置为Recorded等同于合并进其他事件后的 no-op 行为防止后续再投递一个状态。若没有任何 sink 声明为权威则必须等所有sink 都终结事件后才允许发送确认。在 src/config/mod.rs 的acknowledgements_tests模块中可以看到针对propagate_acknowledgements的完整测试含source 被启用确认的 sink 引用后即使 source 自身不支持确认也会告警的行为断言。缓冲级buffer.acknowledgements缓冲配置新增acknowledgements布尔选项控制事件持久化到缓冲时是否即被确认默认false即默认走真实投递后再确认路径。这正好对应 RFC 讨论部分磁盘缓冲 sink一节的两条路径当前仓库中disk_v2缓冲变体在 lib/vector-buffers/src/variants/disk_v2/mod.rs 实现了读者在确认后才推进读取的机制而EventFinalizerGroupsfinalization.rs正是为一条缓冲记录里含多条独立终结事件、在重试恢复时按事件顺序重新挂回 finalizer而设计。完整配置示例TOML# 全局开关默认 true可整体关闭端到端确认 acknowledgements true [sources.my_source_id] # 单个 source 开关默认跟随全局值 acknowledgements true [sinks.my_sink_id] # 权威 sink投递到它就立即终结该事件的所有源批次默认 false authoritative true # 事件落入磁盘缓冲即确认默认 false buffer.acknowledgements true对应地SourceContextsrc/config/source.rs中已包含acknowledgements: bool字段SourceConfig::build拿到上下文后即可决定是否创建BatchNotifierBatchNotifier::maybe_new_with_receiver正是为此提供的条件化创建入口见 finalization.rs。设计权衡Rationale、Drawbacks 与 Alternatives为什么这样设计Rationale这是支撑完整功能所需最少的元数据增量——单个共享引用且Option被优化进Arc中。用Arc引用计数做终结化能保证事件不可能逃逸而不给出状态指示Drop 即上报。不需要或未配置确认的 source 不贡献 finalizer 列表事件创建时零额外分配。向 source 回传终结状态无需查表或遍历拓扑直接走通道。拓扑重配置导致 source 被移除时除了发送时检查通道已关闭之外无需额外处理。代价Drawbacks唯一被点名的缺点每条事件都增加固定的体积开销即使配置根本不支持或不需要端到端确认也不例外。被否定的替代方案Alternatives用常规Vec存储 source 集合合并多 source 时代码更少但合并事件本就有额外开销且数组需再增加一个 word 的数据量。仅存 source 唯一标识符字符串所有终结状态上报都必须走字典查找而不是直接发通道增加运行时开销。落地路线图Plan of AttackRFC 末尾给出实现顺序如今这些项均已在仓库中落地引入struct SourceContext见 src/config/source.rs设置唯一 source 标识符建立 source 元数据与事件丢弃处理修改 sink 以提供终结通知修改 source 以提供终结处理小结端到端确认是 Vector 保证数据不丢的基石特性通过共享 finalizer 按事件/批次聚合状态 one-shot 通道回传这套数据模型同时覆盖了 fan-out、merge、route、丢事件、用户态脚本与磁盘缓冲等全部边界场景。理解acknowledgements全局/source 级、authoritativesink 级与buffer.acknowledgements缓冲级三者的分工与默认值是正确配置端到端确认的关键而 lib/vector-common/src/finalization.rs 中的类型与测试则是把 RFC 概念转化为可运行语义的最佳阅读入口。【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考