
Rerun 流式数据集查询架构深度解析re_datafusion 的 safe horizon、PipelineBudget 与死锁规避设计【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerunre_datafusion是 Rerun 仓库中为多模态机器人数据查询提供 DataFusion 接口的核心 crate负责把 gRPC/对象存储中的 chunk 数据流式导入内存ChunkStore并以RecordBatch增量输出。本文以 ARCHITECTURE.md 为主体结合 pipeline_budget.rs、cpu_worker.rs、io_loop.rs 与 segment_chunk_manifest.rs 源码深入讲解其safe horizon 驱动的增量发射、共享 CPU runtime 线程模型、字节预算 段计数双重门控以及针对八类预算死锁危险的规避策略并给出全部可调环境变量与编译期默认值。读完本文你将理解 Rerun 如何在数据量远超内存时保持单分区工作集有界以及#1736死锁事故背后的诊断思路。为什么数据集查询必须流式化一次数据集查询返回的数据可能远超内存容量一个分区可能包含数百个 segment每个 segment 又有成千上万个 chunk。而下游消费者DataFusion 算子、查询引擎通常只需要以RecordBatch形式增量消费结果。如果先把完整结果集物化再产出内存必然爆炸。SegmentStreamExec的流式查询路径给出的答案是在源 chunk 变得可安全释放时就发射批次从而把每个分区的常驻工作集working set限制在有界范围内而不是等全部结果物化后再输出。这套机制的内存背压原理在 PIPELINE_BUDGET.md 中有专门说明——为什么一个字节有界的 channel 不足以解决 IO → 内存 store 管线的问题。核心概念latest-at 语义与 safe horizon增量发射不是想发就能发它受制于 Rerun 的latest-at 语义实体/a在时间T的值是/a上满足time_min T的最近一个 chunk。因此时间T处的一行数据只有在所有可能影响它的 chunk 都已到达时才能发射。若提前发射下游消费者无法区分完整行与携带前向值carry-forward的不完整行行会静默错误而非可见过期见下文非目标 #2。由此引出全篇最关键的定义safe horizon是最大的T使得对每一个实体所有满足time_min T的 chunk 都已知被收到。CPU worker 通过跟踪每个 segment 的 manifest服务端宣告的 chunk 起始时间表计算 horizon每当 horizon 前进就把时间落在(last_emitted, horizon]区间内的行作为一个或多个RecordBatch发射并把任何未来T horizon的行都不再可达的 chunk 做 GC。horizon 的维护逻辑在 segment_chunk_manifest.rs 中实现其核心不变量为safe_horizon min(所有实体中最早的未收到 time_min) - 1该文件用一个entity_path → (time_min → outstanding_count)的映射跟踪未到 chunk并用一个反向索引entity_heads: BTreeMapTimeInt, usize让safe_horizon查询在O(log n)内完成——这一点很重要因为flush_incremental在每个 chunk 到达时都会调用safe_horizon总开销随开放 segment 数 × 每个 manifest 的实体数增长常量时间查询能控制住聚合成本。顶层数据流IO loop 与 CPU worker 的分工SegmentStreamExec::execute是 DataFusion 的入口。它对每个输出分区各生成两个任务中间由一条有界的tokio::mpscchannel 连接消息类型为CpuWorkerMsg定义见 cpu_worker.rs┌──────────────────────┐ ┌──────────────────────┐ ┌─────────────────────┐ │ server / network │ fetch │ IO loop │ tx │ CPU worker │ RecordBatch │ (gRPC direct URL) │ ────────►│ (per-partition task)│ ──────► │ (per-partition tsk)│ ───────────► └──────────────────────┘ └──────────────────────┘ └─────────────────────┘ ▲ ▲ │ │ chunk_info │ pipeline budget │ └──────────────────────────────────┴──────────────────────────────────┘ reserve / release / publish_segment_finalized完整流程四步CPU workerchunk_store_cpu_worker_thread被派发到进程级共享 CPU runtimeIO loopchunk_stream_io_loop被派发到首次轮询输出流的那个运行时ambient runtime。两者通过单条有界 channel 通信。IO loop 迭代分区内已按 segment 分组的chunk_info行用create_request_batches见 fetch_plan.rs把抓取任务分组为目标大小的批次再把这些批次分配到最多MAX_CONCURRENT_SEGMENTS个 segment 的segment 波wave中。在一个波可以分配 CPU worker 状态或进入抓取缓冲之前IO loop 先通过PipelineBudget的段计数门以零字节预留的方式准入该波。波内每个抓取任务在实际抓取开始之前从共享预算中预留字节解码完成后通过ReservationGuard::commit提交解码后增量并把解码 chunk 推向下游。CPU worker 持有一个HashMapSegmentId, CurrentStores把到达的 chunk 路由到正确的 per-segment 内存ChunkStore。每次 chunk 插入后运行flush_incremental发射 safe horizon 允许的行GC 掉 horizon 已越过的 chunk。CpuWorkerMsg的三种消息也值得注意cpu_worker.rsSegmentChunkCount { segment_id, count }服务端宣告的 chunk 总数驱动segment 完成检测SegmentManifest { segment_id, manifest }由chunk_info行的{filtered_timeline}:start列构建的时序 manifest驱动增量发射与 GC它被Box化以控制枚举体最大变体尺寸——宽 segment 上 manifest 可能有上千条目Chunks(SortedChunksWithSegment)一个 segment 的一批已解码 chunk每个 segment 可有多条。线程模型进程级共享 CPU runtime两个任务运行在不同 Tokio 运行时上这种拆分是刻意为之IO loop运行在 ambient runtime——首次轮询输出流时通过Handle::current()捕获的那个运行时CPU worker运行在一个进程级专用多线程运行时CpuRuntime。CpuRuntime::try_get首次使用时惰性创建并存入static从不 drop因此其工作线程命名为datafusion_cpu_worker存活于整个进程生命周期。所有SegmentStreamExec——无论并发查询、会话还是 scan 叶子——都把cpu_worker任务派发到这个唯一线程池。这是 DataFusionthread_pools示例源自 InfluxDB IOx中每进程一个专用 CPU executor的经典设计换来两个性质CPU 负载下 IO 仍保持响应store 插入、horizon 驱动的发射、GC 都是 CPU 密集型的若跑在 ambient runtime 上会饿死 gRPC 轮询和流消费者每进程 CPU 线程总数有界查询的target_partitions只决定生成多少个cpu_worker任务任务排队进入固定大小的池N 个并发查询共享线程而不是相乘。池的尺寸线程池默认每可用核心一个线程可用RERUN_SDK_NUM_CPUS环境变量覆盖解析后钳制在[1, 可用核心数]首次读取后缓存见 cpu_count.rs。同一个变量还用于设置 Python catalog 会话的datafusion.execution.target_partitions所以一个旋钮同时缩放计划扇出与线程池宽度。网络并发query_dataset扇出流式机制运行之前scan 必须先知道要读哪些 chunk。DataframeQueryTableProvider::scan发出QueryDatasetRequest并收集chunk_info行即上文数据流第 2 步的输入。过滤器下推可以把一个 scan 展开成多个请求——索引范围的OR会为每个分支生成一个请求见apply_filter_expr_to_queries——于是分散读会扇出成数十甚至数百个独立的query_dataset流。这些流并发运行受两层约束每个 scan 内一个宽度为query_dataset_max_concurrency的buffer_unordered定义于 pipeline_budget.rs。响应乱序收集、按 chunk id 去重因此不依赖顺序每进程每个在途请求还持有一个进程级query_dataset_semaphore的 permit。这与共享 CPU runtime、direct-fetch 上限direct_fetch_semaphore是同一推理——受限资源是进程级的客户端是每查询的仅靠每 scan 宽度无法阻止 N 个共存 scan 一起打开 N × width 条流从而触发服务端QueryDataset的流并发限制器它以ResourceExhausted快速失败按请求重试。宽度与 permit 数量是同一个数字DEFAULT_QUERY_DATASET_MAX_CONCURRENCY 16可用RERUN_QUERY_DATASET_MAX_CONCURRENCY覆盖所以孤立 scan 永远不会被自己的 permit 阻塞——信号量只在共租co-tenancy下才起作用。每 segment 的生命周期对每个 segmentCPU worker 跟踪四类状态cpu_worker.rs 中CurrentStores承载expected_chunks来自CpuWorkerMsg::SegmentChunkCount——驱动 segment 完成检测manifest: SegmentChunkManifest来自CpuWorkerMsg::SegmentManifest由chunk_info响应的{filtered_timeline}:start列构建——驱动 safe-horizon 计算last_emitted_time——已发射行的高水位标记用作下一次的filtered_index_range.min - 1使emit_up_to不会重复发射None表示尚未发射任何行last_horizon——最近一次返回的safe_horizon仅供非回归debug_assert!校验。生命周期四个阶段阶段触发条件执行内容Open打开收到第一条SegmentManifest或Chunks消息CurrentStores::new构建内存 store 与可复用的QueryCacheHandleStream流式每条Chunks消息插入 chunk对 manifest 执行record_arrival运行flush_incremental若 horizon 前进则发射 GCFinalize终结received_chunks expected_chunksworker 移除条目调用flush通过emit_up_to(None)排空释放残余字节Drop为 no-opEnd-of-stream流结束channel 关闭且 segment 未完成worker 记录警告并不 flush直接 drop 条目Drop for CurrentStores退回预算预留Drop for CurrentStores的实现cpu_worker.rs是这套设计的安全网无论flush成功、?提前返回、上游错误、消费者中途挂断还是 panicDrop都会恰好执行一次字节退回与段槽位释放保证预留不会泄漏给兄弟分区。非目标四个不要与它们背后的教训文档明确列出早期迭代中被尝试、但因正确性或死锁不变量被破坏而放弃的设计并警告不得未经检查测试用例就重新引入。这是#1736事故的诊断线索。1. 不要重新引入 IO 侧重排缓冲流式化之前chunk_stream_io_loop持有一个BTreeMaptask_idx, VecChunksWithSegment按task_idx顺序排空让 CPU worker 一次只看到一个 segment 的抓取。两个问题预算满时的队头阻塞解码 chunk 因慢的task_idx前驱而滞留缓冲即使 CPU worker 空闲可消费它们的字节预留持续占用预算新的 IO 抓取被满是 CPU 侧尚不可触及的内存的预算阻塞——本可释放预算的 chunk 恰恰卡在缓冲后面突发抓取延迟下的无界放大单个慢前驱把后续所有抓取钉在缓冲里重排缓冲唯一的上限就是管道预算本身而缓冲的全部效果正是让预算持续饱和。重构后的 CPU worker 按SegmentId键控HashMapSegmentId, CurrentStores哪个 segment 的 chunk 先到就先处理谁重排缓冲强制的一次一个 segment不变量不再需要。2. 不要发射越过 safe horizon 的行flush_incremental每轮都用filtered_index_range (last_emitted, horizon]重建QueryHandle。上限是正确性要求而非可调参数在 latest-at 语义下发射越过 horizon 的行等于把实体值当作最终值发布而此时该实体的后续 chunk 尚未到达。下游消费者无法区分完整行与带 carry-forward 值的不完整行行会静默错误。同样的约束也解释了为何流结束清理选择丢弃不完整 segment 而不是 flush 它——carry-forward 值看起来正确非空但反映的是仍在途的数据。3. 不要因为字节预算够了就丢弃段计数门MAX_CONCURRENT_SEGMENTS 3在字节预算之上存在有两个理由CPU worker 的 per-segment HashMap 随在途 segment 数线性增长。没有上限时一个长尾慢 segment 可能让 IO 侧在任一段终结前打开上百个并发 segmentCPU 内存成本每段ChunkStore manifest cache最终会盖过字节预算的上限per-segment manifest 的outstanding_time_mins_per_entity映射是O(N_entities × N_chunks_per_entity)。虽然每次flush_incremental只扫一个 segment 的 manifest但 worker 的聚合每 tick CPU 成本仍随开放 segment 数 × 每个 manifest 大小增长。有界的 segment 数约束住这个聚合。该上限在PipelineBudget::try_admit的active_segments锁下准入发生在 segment 波被允许分配 CPU worker 状态或进入抓取缓冲之前。波内抓取仍先预留字节再解码但不再用抓取并发缓冲去等待尚不可准入的未来 segment。只准入一个代表 segment_id 会让create_request_batches的多 segment 抓取偷偷打开超出上限的 segment。4. 不要恢复 CPU 侧的publish_segment_started路径MAX_CONCURRENT_SEGMENTS上限需要当前哪些 segment 开放的单一事实来源。早期迭代把这个簿记放在 CPU 侧CPU worker 收到某新 segment 的第一个 chunk 时调用publish_segment_started(segment_id)。这看似自然CPU worker 已经按 segment 键控CurrentStores但准入必须与字节预留原子化。具体失败场景IO loop 已有 segmentK-2、K-1、K开放上限 3正要抓取K1的第一批 chunkIO 为K1调用reserve(bytes)。段门还不知道K1——没有 chunk 到过 CPU 侧active_segments仍显示 3——检查通过字节被预留该抓取在途时IO 拉下一批K2的第一个抓取。同样的检查同样的结果等 CPU worker 终于收到K1的第一个 chunk 并调用publish_segment_started时IO 已经为多个越过上限的 segment 超量预留了字节。上限退化为建议性它本要约束的 HashMap manifest 成本会无界爆炸。修复方式在 IO 侧、与准入字节预留的同一把active_segments锁下发出信号。PipelineBudget::try_admit同时接收(segment_ids, bytes)要么两个门都清、槽位与字节都被占用要么都不清、调用者停车。segment 波调度器在发送 CPU 元数据或启动该波抓取之前用同一路径做零字节预留。IO 侧从create_request_batches已经知道批次覆盖哪些 segment准入时它已具备全部信息。publish_segment_finalized是对称的反向信号但它刻意留在 CPU 侧终结释放槽位而非占用槽位所以最终一致的信号是安全的——短暂多计开放段只会略微降低并发绝不会违反上限。而且终结是在 CPU 侧可观测的依赖 horizon 发射排空 segment 与下游消费者速率IO loop 都看不到所以信号天然存在于Drop for CurrentStores。预算门控字节预算 段计数 stall-breakerPipelineBudget::try_admit是 IO 抓取的单一原子准入点必须同时通过三道门或触发 stall 逃生┌───────────────────────────────────────────────┐ reserve ────────►│ force_overcommit set? yes → admit everything │ │ no ↓ │ │ active_segments new_segments MAX? yes → park │ no ↓ │ │ try_acquire(bytes) │ │ - CAS on current against budget │ │ - fail → park; success → admit │ └───────────────────────────────────────────────┘ │ ┌──────────────────────────────┴───┐ │ wait_queue: BinaryHeapReverse │ │ PriorityWaiter { task_time_min,│ │ seq, notify } │ │ │ │ → earliest-time wakes first │ └──────────────────────────────────┘唤醒来源release(bytes)——CPU 侧归还解码字节gc_up_to_horizon的增量归还或flush的最终归还adjust_reservation(estimated, reserved, actual)——IO 侧解码后若actual reservedpublish_segment_finalized(segment_id)——CPU 侧Drop for CurrentStoresnotify_empty_emit跨过STALL_EMPTY_EMIT_THRESHOLD且预算处于STALL_SATURATION_THRESHOLD→ 置force_overcommit并唤醒一个。force_overcommit的复位任何release真实的字节进展notify_row_emittedCPU 侧已发射行。try_acquire用 CAS 循环原子推进current读当前值若current reserved_bytes budget拒绝提交否则仅在无并发修改时原子推进。等待队列是按task_time_min最早时间优先键控的BinaryHeapwake_next弹出最小时间且可准入的 waiter跳过已取消条目与仍被段门阻塞的 waiter避免释放的预算浪费在还无法使用它的高优先级 waiter 上。停车前还有一次入队后重检以关闭丢失唤醒的竞态窗口。第一个 park 与每第十次 re-park 会输出 info 级背压日志无需RUST_LOGdebug即可在生产可见。关于内存背压的总体动机可参阅 PIPELINE_BUDGET.mdchunk 从远端S3 / gRPC抓取、解码为 Arrow、再插入ChunkStore供查询执行若无背压IO 管线会以网络允许的速度抓取在 CPU 线程处理并 GC 之前就可能消耗无界内存。PipelineBudget通过跟踪在途解码字节并在耗尽时阻塞 IO 任务制造一个自然的滑动窗口IO 侧至多领先 CPU 侧budget字节。GC 中的 carry-forward 保护gc_up_to_horizon在基于区间的直觉下看似正确——丢弃time_max horizon的 chunk——但在 Rerun 的 latest-at 语义下这是错误的。示例实体/a只有t10一个 chunk实体/b有t20和t40两个 chunk。/b20到达后safe horizon 是39B 最早的未收到时间是 40减 1。此时若因/a10的time_max 39就丢弃它所有 10的行都会把/a发射为 null而不是把/a10的值 carry-forward。两种保护通过ChunkStore::gc选项应用protected_chunks——store 内所有实体在 horizon 处的latest_at_relevant_chunks_for_all_components的并集包含静态 chunk。这是任何未来T horizon的行在 latest-at 下可能解析到的 chunk 集合protected_time_ranges——filtered timeline 上的(horizon1, inf]。horizon 之后的内容按定义不可读。两者交集才是可 GC 的目标——通常是一个 segment 已发射 chunk store 的主体因为 latest-at 通常只解析到每个实体最近的少量 chunk。Manifest/chunk 分歧静默数据丢失的检测SegmentChunkManifest::record_arrival返回bool#[must_use]。manifest 由服务端chunk_info行构建CPU 侧随后看到的是抓取路径的实际chunk。两者可能分歧服务端chunk_info与 chunk 抓取响应不同步部署进行中、缓存过期chunk 在 filtered timeline 上的time_min与{timeline}:start宣告值不匹配chunk 编码或拆分逻辑 bugCPU 侧的:start提取漏行build_segment_manifests的 bug。分歧若不上报就是静默数据丢失chunk 插入 store但如果safe_horizon已越过该 chunk 的time_min因为 manifest 不知道要以其为门控行范围过滤器(last_emitted, horizon]会完全排除该 chunk 的行它们永不发射。record_arrival的契约返回false表示不符合预期集合worker 触发re_log::debug_panic!re_log::error_once!。chunk 仍会插入——丢弃它更糟——但日志会把问题暴露出来供调查。预算危险类别设计必须保持关闭的八个缺口字节预算在release阻塞在reserve已停车的同一批 chunk 上时死锁若没有 segment 能进展其剩余 chunk 卡在reserve而预算只能通过 segment 完成来退款管线就从网络速率受限转为释放速率受限吞吐归零。重构后的设计通过随 safe horizon 前进逐 chunk 释放而非每 segment 一次避开此不变量使释放速率即使在饱和下也被 chunk 到达速率下限约束。以下 A–H 是设计吸收的死锁形状压力含一个近失是每个机制存在的承重理由A — 小数据集上的尺寸钳制。total_uncompressed × fraction / num_partitions在小数据集上可能低于MIN_BUDGET_PER_PARTITION钳制占主导有效预算变成MIN × N。没有逐 chunk 释放时任何 segment 工作集超过MIN的分区都会死锁。逐 chunk 释放把工作集钉在~MAX_CONCURRENT_SEGMENTS × bytes_per_chunk对典型 chunk 远低于 64 MiB 下限。stall-breaker 是对残余病态情况的最后安全网。B — 算子通过环境变量把上限设得过低。如RERUN_PIPELINE_BUDGET_MAX128MiB低于观测到的峰值工作集。与 A 相同代码路径人为驱动相同解法。C — 多分区 × 宽数据集。每分区工作集是~MAX_CONCURRENT_SEGMENTS × bytes_per_chunk与 schema 宽度无关。宽数据集产生更小更频繁的 chunk而不是撑爆预算。D — 乱序 segment chunk 到达。两个机制协同task_time_min上的优先级唤醒推进 horizon 的 chunk 时间最小预算释放时最先被唤醒抢占争抢同一槽位的晚时间 chunkstall-breaker若连续STALL_EMPTY_EMIT_THRESHOLD次flush_incremental在预算STALL_SATURATION_THRESHOLD饱和下发射零行force_overcommit绕过两道门让停车的 horizon 推进 chunk 准入、落到 CPU 侧、推进 horizon并通过随后的真实释放偿还超支。E — CPU worker 在release前出错。flush(...).await?传播Err。没有 RAII guard 时?会在匹配的release运行前从chunk_store_cpu_worker_thread返回把分区字节钉住整个查询并饿死同查询的兄弟分区panic 越过 release 行同理。Drop for CurrentStores在released标志为false时释放store_bytes()flush内的显式成功路径在 drop 前置位该标志退款恰好一次。任何?、panic 或取消路径都经Drop退款。单元测试test_current_stores_drop_refunds_budget见 pipeline_budget/tests.rs。F — segment 中途流取消。消费者挂断LIMIT、计划取消。CPU worker 的CurrentStores在flush完成前被 drop。Drop for CurrentStores与 E 相同方式覆盖此路径released false退款触发。G — EMA 过度估计。若干大膨胀样本后estimate_multiplier爬向MAX_ESTIMATE_MULTIPLIER。即使真实膨胀已回落到 ~1.0后续预留仍为estimated × multiplier。MAX_ESTIMATE_MULTIPLIER钳制使其有界EMA 平滑因子衰减离群点的影响。H — 单个抓取大于整个预算——不是死锁。reserved budget路径以 warn 级日志绕过门。current因该抓取超支。无死锁为文档化的边缘情况。预算的自适应估计预算总尺寸只是故事的一半——每次预留还需要合理的单次抓取大小。服务端报告的未压缩 chunk 尺寸是线编码估计实际解码SizeBytesArrow heap dictionary index 开销会上下漂移。PipelineBudget维护一个学习型乘数reserved estimated_uncompressed * estimate_multiplier每次完成的抓取通过adjust_reservation回馈(estimated, actual)样本原始比值被钳制在[MIN_ESTIMATE_MULTIPLIER, MAX_ESTIMATE_MULTIPLIER]单个病态 chunk 不会把未来所有预留钉在天花板再以平滑因子ESTIMATE_EMA_ALPHAα0.2折入 EMA——数样本内收敛同时容忍一次性离群点。乘数从INITIAL_ESTIMATE_MULTIPLIER1.5x起步前几次冷启动预留过度记账而非不足记账随后收敛到数据集真实比值。乘数是 per-PipelineBudget的且每个查询都新建预算所以跨查询不持久化——不同 schema/编解码的膨胀比不同用过期的旧乘数比短暂冷启动更糟。所有原子操作使用AcqRel/Acquire序以保证弱序架构ARM上的跨线程可见性。编译期默认值默认值值选择理由BUDGET_FRACTION0.25IO 至多领先 CPU 侧四分之一的查询总解码估计量MIN_BUDGET_PER_PARTITION64 MiB为真实工作负载下的MAX_CONCURRENT_SEGMENTS * bytes_per_chunk留足空间MAX_BUDGET_PER_PARTITION1 GiB为大数据集封顶最坏情况下的每分区 RSSMAX_CONCURRENT_SEGMENTS3约束 CPU worker 的 HashMap 大小与O(N_open_segments)的 horizon 重算工作INITIAL_ESTIMATE_MULTIPLIER1.5冷启动多预留约 50%若膨胀高于典型值前几次抓取不会瞬时 OOMESTIMATE_EMA_ALPHA0.2EMA 数样本内收敛同时容忍一次性离群点STALL_EMPTY_EMIT_THRESHOLD20高到horizon 未移动的停顿不触发低到真实死锁快速打破STALL_SATURATION_THRESHOLD0.95等待慢 IO 且预算大部分空闲的查询不应误触发旁路FLUSH_BATCH_ROWS2048继承自非流式路径#1794/#1822FLUSH_BATCH_BYTES200 MiB同上重要现状说明ARCHITECTURE.md 自我标注为target end-state目标最终状态流式重构是增量落地的并非每个章节都反映main分支当前代码。这一点在源码中有直接印证pipeline_budget.rs 顶部注释明确写道——当前 CPU worker 在释放前会缓冲整个 segment任何低于最大解码 segment 工作集的每分区上限都会死锁在 PR #1736 的rerun-synthetic-structs-10k50 segment 上实测复现自适应尺寸产生 283 MB 总预算钉在 282/28372 个 IO 任务停在 wait #122 分钟零release因此目前实际生效的默认值是FRACTION1.0、MIN4 GiB、MAX1 TiB预算实际上被解除待逐 chunk 流式释放重构落地后再切回上表的 0.25 / 64 MiB / 1 GiB 目标值。阅读本文时请以目标架构理解上表以当前源码默认值理解实际行为。运行时环境变量总览字节类常量可在运行时通过RERUN_PIPELINE_BUDGET_*环境变量覆盖阈值与计数类常量刻意不可覆盖它们为上述架构不变量而选而非按部署调优。尺寸值接受 SI/IEC 后缀64MB、1GiB、512KiB或裸正整数按字节值会去除首尾空白空字符串视为未设置不可解析或越界的值记 error 日志并回退到编译期默认若覆盖后MIN MAX两者都回退默认而非在后续clamp()上 panic。变量类型接受范围默认值RERUN_PIPELINE_BUDGET_MINsize 04GiBRERUN_PIPELINE_BUDGET_MAXsize 01TiBRERUN_PIPELINE_BUDGET_FRACTIONfloat(0.0, 1.0]1.0RERUN_DIRECT_FETCH_MAX_CONCURRENCYint 1128RERUN_SEGMENT_ADMISSION_CAPint3..1024unsetRERUN_ADAPTIVE_SEGMENT_ADMISSIONbooltrue/falsetrue段准入在每次查询时解析一次统一控制传输批处理、波形成、gRPC 分组与 CPU 状态准入。默认自适应策略在 3-or-16 候选中选择pipeline_budget.rs 中ADAPTIVE_SEGMENT_ADMISSION_CAP 16与DEFAULT_QUERY_DATASET_MAX_CONCURRENCY 16对应RERUN_ADAPTIVE_SEGMENT_ADMISSIONfalse将有效上限保持为 3 同时保留候选遥测有效的RERUN_SEGMENT_ADMISSION_CAP始终优先无效的精确覆盖即使自适应准入开启也失败关闭到 3。若服务端未提供未压缩尺寸旧服务端则回退使用压缩线尺寸——这会低估产生更多而非更少背压。单 chunk 大于整个预算时以警告放行预算临时超支释放后恢复预算的Drop实现会输出生命周期摘要日志便于事后查询分析peak_current、total_released_bytes、total_releases等诊断位于结构体自身。总结一条从事故到架构的演进主线re_datafusion的流式查询架构可以被浓缩为一句话以 latest-at 语义下的 safe horizon 为单一事实来源驱动增量发射与逐 chunk GC并用字节预算 段计数 优先级唤醒 stall-breaker四重机制保证 IO 侧无论网络多快都无法让分区工作集失去边界。而#1736那次 22 分钟的零release死锁正是理解全部设计的钥匙——每一个非目标、每一道预算门、每一次Drop退款都是为了让释放速率 ≥ chunk 到达速率这条不变量在任何调度、取消、错误与乱序面前都能成立。对于想要深入源码的读者建议按此顺序阅读pipeline_budget.rs预算与门控→ segment_chunk_manifest.rshorizon→ cpu_worker.rs发射与 GC→ io_loop.rs并发抓取配合 pipeline_budget/tests.rs 中的test_current_stores_drop_refunds_budget等用例理解不变量如何被验证。【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考