ARTICLE DETAIL

资讯详情

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

Quickwit 队列源(Queue Source)设计解析:基于消息队列的 Exactly-Once 索引机制

Quickwit 队列源(Queue Source)设计解析:基于消息队列的 Exactly-Once 索引机制 Quickwit 队列源Queue Source设计解析基于消息队列的 Exactly-Once 索引机制【免费下载链接】quickwitCloud-native OSS search engine for observability项目地址: https://gitcode.com/GitHub_Trending/qu/quickwitQuickwit 是一款面向可观测性的云原生开源搜索引擎。在其索引体系中队列源queue source是一套将消息队列如 AWS SQS与对象存储通知打通、并借助 Metastore Shard API 实现每条消息内容恰好被索引一次exactly once的关键机制。本文以 design.md 为骨架结合 queue_sources 模块的源码实现系统讲解 Quickwit 如何从队列消费消息、跟踪其生命周期、处理可见性超时、回收去重状态并最终保证对象文件恰好被索引一次。读完本文你将掌握队列源在 Quickwit 中的定位与适用边界、Shard 表去重与分片所有权仲裁的具体流程、可见性扩展任务的工作原理以及这套设计与 SQS 文件源、ingest v2 等周边模块的衔接方式。背景为什么需要队列源与 Exactly-Once索引过程中除了常见的失败进程崩溃、网络抖动、提交失败队列系统本身还天然会产生重复投递SQS、Pub/Sub 这类至少一次at least once投递语义的消息队列在没有消费者确认前会把同一消息重新暴露给消费者。如果 Quickwit 索引器把同一批对象文件重复读入索引管线就会产生重复文档。因此队列源的核心命题是在队列只保证 at-least-once 投递的前提下让每个对象文件object file恰好被索引一次。design.md 给出的答案是借助 Metastore 中的shard 表来跟踪每个文件对象的处理进度每个文件对象被建模为一个shardshard ID 就是该文件的 URI某个 shard 的索引进度在与 split 发布split publishing相同的数据库事务中提交到 shard 表——也就是说进度提交和索引结果发布要么同时成功、要么同时失败经过一段时间称为去重窗口deduplication window之后shard 被垃圾回收以控制 shard 表的规模。这个设计把去重状态从内存搬到了具备事务能力的 Metastore 中从而让多个索引管线并发消费同一队列成为可能谁先在 shard 表上登记了某个 shard 的归属权谁就负责索引这份内容其他管线看到归属权冲突时便让出。总体架构QueueCoordinator 与三大组件design.md 明确指出QueueCoordinator是从接收消息到索引完成后确认消息acknowledge这一整套机制的具体实现。在源码中它位于 coordinator.rs其 API 刻意设计得与Sourcetrait 非常接近注释原文Its API closely resembles thecrate::source::Sourcetrait从而让新的队列源实现非常直接。QueueCoordinator与三个主要组件交互组件源码位置职责Queue抽象队列接口mod.rs屏蔽具体队列实现SQS、Pub/Sub…只暴露三个核心操作QueueLocalState本地状态local_state.rs内存中跟踪本管线已知消息的四态转换QueueSharedState共享状态shared_state.rsShard API 客户端负责分片归属权仲裁与后台清理从 coordinator.rs 可以看到QueueCoordinator同时持有queue、shared_state、local_state以及publish_token、visibility_settings等字段其中publish_token由PublishToken::resolve(node_id, )生成是当前管线在该索引源上的所有权凭据。抽象队列接口Queuedesign.md 强调Queue抽象只要求底层队列保证至少一次投递并把真实队列的 API 面缩减为 3 个函数。对应 mod.rs 中的 trait 定义receive(max_messages, suggested_deadline)拉取可处理的消息。实现方自行决定无消息时的等待策略通常用长轮询有消息时应尽快返回max_messages由实现方钳制例如 SQS 单次最多 10 条。modify_deadlines(ack_id, suggested_deadline)延长消息的可见性超时visibility timeout即推迟该消息对其他消费者重新可见的时间点返回新的截止时间。acknowledge(ack_ids)在索引成功后将消息从队列中彻底删除。trait 注释还强调了一个工程细节从消息被 receive 的那一刻起调用方就必须及时维护消息可见性否则消息会重新暴露给其他索引管线导致它们去竞争提交锁commit lock产生额外的重复处理压力。在 Quickwit 当前仓库中Queue有两个实现SqsQueuesqs_queue.rs基于aws-sdk-sqs。它使用 20 秒长轮询通过ApproximateReceiveCount系统属性获得投递次数并以不同的重试策略区分三类调用receive用标准重试、acknowledge不重试消息若被再次投递反而会自然重试、modify_deadlines用激进重试以尽力保住消息所有权源码注释Retry aggressively to avoid loosing the ownership of the message。MemoryQueueForTestsmemory_queue.rs一个带可见性过期回队逻辑的内存队列用于单元测试例如test_process_multiple_coordinator中模拟两个 coordinator 竞争同一消息的场景。此外helpers.rs 中的QueueReceiver是一个有状态包装器把缓慢的receive()调用切成短迭代通过tokio::select!与固定迭代时长竞争避免 actor 系统长时间阻塞在长轮询上而无法及时处理 mailbox 消息或退出。消息生命周期QueueLocalState 的四态机design.md 描述QueueLocalState是一个内存数据结构跟踪当前 source 对最近收到消息的认知管理消息在 4 个状态之间的迁移ready for read待读read in progress读取中awaiting commit等待提交completed已完成local_state.rs 的实现把这四个状态映射为四个字段pub struct QueueLocalState { ready_for_read: VecDequeReadyMessage, // 已收到、可开始读取的消息队列 read_in_progress: OptionInProgressMessage, // 正在读取并送入 DocProcessor 的消息 awaiting_commit: BTreeMapPartitionId, String, // 已读完、仍在索引中记录其 ack_id completed: BTreeSetPartitionId, // 已完全索引并提交 }结合 coordinator.rs 中poll_messages与emit_batches的驱动逻辑状态迁移的完整路径是消息被receive后先经过预处理器pre-process见下文随后通过checkpoint_messages与共享状态核对归属权拿到可处理位置的消息进入ready_for_read队列set_ready_for_read此时会为该消息生成可见性扩展任务spawn_visibility_taskemit_batches从ready_for_read取出一条消息get_ready_for_read会跳过可见性扩展已失败的消息调用start_processing创建InProgressMessage并置为read_in_progressInProgressMessage内部是一个ObjectUriBatchReader逐批读出文档送入DocProcessor读到 EOF 后drop_currently_read把该消息移入awaiting_commit同时请求最后一次可见性扩展为提交留出时间当索引提交完成、checkpoint 推进后suggest_truncate回调中调用mark_completed把分区移入completed并返回其 ack_id 供Queue.acknowledge删除消息。这一套状态机的行为在 coordinator.rs 的测试中有非常直白的验证例如test_process_local_duplicate_message证明同一分区在本地已跟踪时不会重复处理test_process_shared_complete_message证明当共享状态显示该分区已到 EOF 时消息被直接确认丢弃acknowledge而不产生任何 batchtest_process_multiple_coordinator证明第二个 coordinator 会从共享状态学到该消息很可能仍在处理中而跳过它。消息预处理从通知到分区 ID在进入四态机之前原始消息需要先被最小化转换以发现其分区 IDpartition ID。这一步在 message.rs 的RawMessage::pre_process中完成支持两种MessageTypeS3Notification消息体是 S3 事件通知 JSON。uri_from_s3_notification会丢弃s3:TestEvent测试事件只接受ObjectCreated:*事件从中提取 bucket 与 object key 拼出s3://bucket/keyURI并做 URL 解码测试test_uri_from_s3_notification_url_decode验证了hello%3A%3Aworld%3A%3Alogs.json解码为hello::world::logs.json。RawUri消息体本身就是一个对象 URI 字符串如s3://bucket/key。该函数返回两类错误Discardable可以确认并丢弃如 S3 测试事件或ObjectRemoved事件与UnexpectedFormat无法解析仅计数并限速记日志design 注释建议配合死信队列使用。pre_process之后消息被unique_by(partition_id)去重防止同批内重复再交给共享状态仲裁。从源码注释可以推断这种先发现分区 ID、再决定是否深度处理的顺序是有意为之如果分区已被处理过就完全不必做昂贵的对象读取工作。去重与所有权仲裁QueueSharedState 与 Shard APIdesign.md 用较大篇幅描述了共享状态的职责它是Shard API 的客户端——一个主要为 ingest v2 设计的 Metastore API相比旧的 checkpoint API以 blob 形式存放在 index model 的某个字段中是明显的改进。OpenShards所有权判定队列源以能唯一标识消息内容的 ID作为 shard ID 打开 shard文件源即文件 URI。每个 source 拥有唯一的 publish token通过OpenShards请求提交。OpenShards的响应会返回第一个调用该 API 的调用者 token据此有三种分支shared_state.rs 的acquire_partitions与 design.md 完全对应返回的 token 与当前管线 token 一致→ 本管线拥有该消息内容的归属权可以继续索引。token 不一致另一个管线拥有归属权→ 检查 shard 内容若已完全处理EOF→ 直接确认acknowledge并丢弃消息若最后更新时间戳较旧例如超过两倍 commit timeout→ 认为处理已过期例如原属管线已崩溃执行AcquireShards把 shard 的 token 更新为本管线的 token表明后续处理由本管线接管若最后更新时间戳较新 → 认为另一管线仍在处理中直接丢弃消息且不做 acknowledge等其可见性超时后自然重新投递。design.md 特别指出AcquireShards存在竞态两个管线可能并发获取同一 shard双方都会以为自己拥有归属权最终其中一个会在提交commit阶段失败——这是设计上接受的结果因为提交阶段的原子性保证不会产生重复数据。源码中的实现细节还揭示了几个补充规则shared_state.rs位置position为 EOF或者本管线拥有且位置为 Beginning都可直接进入处理队列不是本管线拥有且未过期→ 跳过对应 design 的直接丢弃、不确认位置已过期但需重新获取 → 批量调用acquire_shards一个防御性检查若 shard 由本管线拥有但位置不在 Beginning则直接bail!报错提示这绝不应发生请上报 issue。时间戳与宽限期QueueSharedState持有一个reacquire_grace_period重新获取宽限期字段在 coordinator.rs 中创建共享状态时被设为2 * commit_timeout_secs与 design.md 中旧时间戳例如两倍 commit timeout的表述一致。acquire_partitions用update_timestamp与当前时间差来判定 shard 是否过期is_stale。测试test_re_acquire_shards_within_grace_period与test_re_acquire_shards_after_grace_period分别验证了宽限期内不可抢占、宽限期后可AcquireShards两种行为。后台清理PruneShardsQueueSharedState::new会tokio::spawn一个后台清理任务run_cleanup_taskshared_state.rs以pruning_interval为周期调用 Metastore 的prune_shards参数包括max_age_secsshard 最大存活时长与max_countshard 最大数量上限。若两者都未配置则跳过。该任务的生命周期通过Weak()与主 state 绑定主 state 被回收时任务自动退出。design.md 还强调了一个与并发扩展相关的重要设计垃圾回收由队列源自身负责——每个带队列源的管线都会 spawn 一个 GC 任务为避免管线数量增多时对 Metastore 造成过大压力GC 调用由控制平面control plane做去抖debounced。可见性扩展任务Visibility Extension Taskdesign.md 指出为了让消息在处理期间对其他管线保持不可见每条收到的消息都会派生一个可见性扩展任务负责在可见性截止时间临近时持续延长可见性超时。当消息最后一批被读取并送入索引管线时请求一次最后的可见性扩展为索引完成留出时间典型值为两倍 commit timeout随后停止该扩展任务。源码实现在 visibility.rs 中是一套基于 Quickwit actor 模型的完整机制spawn_visibility_task为每条消息 spawn 一个名为QueueVisibilityTask的 actor返回VisibilityTaskHandle。handle 内部持有Arc()强引用actor 通过Weak()感知外部引用是否释放当 handle 被 dropactor 退出可见性不再被维护测试test_visibility_task_stop_on_drop验证了这一点。actor 的Loop处理器周期性调用extend_visibility即queue.modify_deadlines(ack_id, deadline_for_default_extension)默认扩展 1 分钟再根据剩余时间调度下一次Loop。next_extension会预留request_timeout3 秒与request_margin1 秒即在距当前截止时间减去这两个缓冲后才发起扩展请求避免扩展请求本身超时而错过截止时间源码注释We prefer applying ample margins ... to avoid missing deadlines while also keeping the number of extension requests (and associated cost) small。最后一次扩展由RequestLastExtension消息触发drop_currently_read在消息读完后调用handle.request_last_extension()将截止时间再延长deadline_for_last_extension2 * commit_timeout然后置last_extension_requested true扩展循环随之停止测试test_visibility_task_request_last_extension验证了扩展行为与最终期限。VisibilitySettings::from_commit_timeout统一从 commit timeout 推导四个时长参数接收截止时间为2 分钟 commit_timeout默认扩展为 1 分钟最终扩展为2 * commit_timeout请求超时 3 秒请求余量 1 秒。值得一提的是get_ready_for_read会跳过extension_failed()的消息即扩展任务已失败因为这类消息反正会重新出现在队列中。通往配置SQS 文件源的落地design.md 是设计文档但它描述的能力在 Quickwit 中已落地为可配置的SQS 文件源。在 source_config/mod.rs 中FileSourceSqs定义了完整配置项配置项默认值说明queue_url无必填SQS 队列地址如https://sqs.us-east-1.amazonaws.com/123456789012/queue-namemessage_type无必填s3_notification或raw_uri对应源码FileSourceMessageType::S3Notification/RawUrideduplication_window_duration_secs3600去重窗口时长秒即 shard 在 metastore 中的最大存活时长deduplication_window_max_messages100000去重窗口内最多跟踪的 shard 数量deduplication_cleanup_interval_secs60shard 清理任务执行周期秒在 coordinator.rs 的try_from_sqs_config中可以看到这三组去重参数如何被映射到QueueSharedStatededuplication_window_duration_secs→shard_max_agededuplication_window_max_messages→shard_max_countdeduplication_cleanup_interval_secs→shard_pruning_interval。也就是说去重窗口越大可容忍的重复投递时间跨度越长但 shard 表占用的 Metastore 空间也越大清理间隔越短表回收越及时但对 Metastore 的调用频率也越高。一个典型的 SQS 文件源配置片段如下参考 source_config/mod.rs 中的文档示例source: source_id: my-sqs-source source_type: file params: notification: type: sqs queue_url: https://sqs.us-east-1.amazonaws.com/123456789012/queue-name message_type: s3_notification deduplication_window_duration_secs: 3600 deduplication_window_max_messages: 100000 deduplication_cleanup_interval_secs: 60适用边界与性能约束design.md 非常坦诚地给出了这套设计的适用边界这是读者在使用队列源前必须理解的内容由于每一条消息都会被 Metastore 跟踪这套设计在高消息速率下表现不佳。例如对于每条消息只含一个事件的数据流它并不高效。作为经验法则为了保护 Metastore不建议用本设计处理超过每秒 50 条消息。这意味着高吞吐只能靠放大每条消息的内容来实现——例如使用带队列通知的文件源file source with queue notifications让每条消息携带一个较大的文件对象索引引擎再去对象存储中批量读取。这一约束直接决定了架构选型它天然适合对象存储事件通知 批量文件索引的场景而非高频小事件流场景后者应优先考虑 ingest v2 等其他路径。通用性与扩展方向design.md 最后强调该模块被设计得足够通用以便接入其他队列实现例如GCP Pub/Sub从对象存储之外的数据源取数例如直接从消息体取数据。从源码看这些方向已经留下了明确的扩展点Queuetrait 是抽象接口注释即提到based on the AWS SQS and Google Pubsub APIsmessage.rs 中的MessageType枚举注释里留有GcsNotification与RawData的占位当前为注释状态PreProcessedPayload同样有被注释掉的RawData(OwnedBytes)变体。可以推断这些是尚未启用的未来扩展点——新的队列源只需实现Queuetrait 并在MessageType中新增一种解析方式即可接入现有机制。总结一条消息的完整旅程把 design.md 与源码拼合起来一条消息在队列源中的完整旅程如下Queue.receive长轮询收到消息附带初始可见性截止时间pre_process解析出分区 ID文件 URIQueueSharedState.acquire_partitions通过OpenShards判定归属权本管线拥有 → 继续他方已处理完EOF→ 确认丢弃他方处理过期 →AcquireShards接管他方仍在处理 → 不确认地丢弃等待重投消息进入QueueLocalState的ready_for_read同时 spawn 可见性扩展 actor 持续续期emit_batches逐批读取对象内容送入索引管线读完后请求最后一次扩展2 × commit_timeout消息转入awaiting_commitsplit 发布与 shard 进度在同一事务中提交到 Metastoresuggest_truncate将分区标记为completed并acknowledge删除队列消息后台PruneShards任务周期性清理超过去重窗口的 shardGC 调用由控制平面去抖以保护 Metastore。这套设计用队列可见性超时 Metastore 分片归属权 事务化进度提交三件套在至少一次投递的队列之上构建了 exactly-once 的索引语义是理解 Quickwit 对象存储索引路径与 ingest v2 演进关系时不可跳过的一环。若想深入验证本文结论建议直接阅读 design.md 及其同目录下的 coordinator.rs、shared_state.rs、local_state.rs、visibility.rs 四个文件其中内嵌的单元测试几乎为每种状态分支提供了可运行的证据。/output_article【免费下载链接】quickwitCloud-native OSS search engine for observability项目地址: https://gitcode.com/GitHub_Trending/qu/quickwit创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表