Rerun SDK 微批处理(Micro Batching)机制详解:延迟与吞吐的平衡之道
Rerun SDK 微批处理Micro Batching机制详解延迟与吞吐的平衡之道【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerunRerun SDK 会在后台线程中自动完成数据的微批处理micro-batching将多条日志行聚合为 Chunk 后再交给下游 sink从而在数据尽快可见低延迟与减少元数据开销、提升带宽与 CPU 利用率高吞吐之间寻找最佳平衡点。本文以官方参考文档 docs/content/reference/sdk/micro-batching.md 为主体结合re_chunk与re_sdk的源码实现完整讲解微批处理的时间/空间双阈值触发机制、三个RERUN_FLUSH_*环境变量、Python 与 Rust 的代码级配置方式以及批处理后端ChunkBatcher的真实工作流程帮助你针对实时流式传输、录制文件等不同场景做出正确的调优决策。为什么需要微批处理每次log调用都会产生一行row数据。如果 SDK 为每一行都单独构造一个数据块Chunk并立刻下发那么每个数据块都要携带一份独立的元数据timeline 描述、组件描述符、Arrow 数组头等对于高频、小体积的日志调用而言元数据开销会显著吃掉带宽CPU 也被大量浪费在重复的对象构造与序列化上。微批处理的思想是在后台线程中把短时间内到达的多行数据暂存起来累积到一定规模时间或空间阈值后再一次性打包成 Chunk 下发。这样既保证了数据不会无限期滞留由时间阈值兜底又保证了批量打包带来的吞吐收益由空间阈值触发。正如官方文档所述这套机制与运行在数据存储侧的 compaction 机制 高度相似、互为映照SDK 侧用批处理减少上行数据的碎片化存储侧用压缩合并减少落盘后数据的碎片化。核心机制时间与空间双阈值谁先触发谁生效从源码注释可以明确看到crates/store/re_chunk/src/batcher.rsThe batching process is triggered solely by time and space thresholds -- whichever is hit first.冲刷flush由两类阈值驱动先触达哪一类就按哪一类触发时间阈值后台线程维护一个周期性 tick 定时器每隔flush_tick时长强制冲刷所有已累积的行空间阈值已累积行的字节数达到flush_num_bytes或行数达到flush_num_rows立刻冲刷。值得注意的是空间阈值是达到或超过即触发源码比较因此把阈值设为零会退化为每行一刷同时最终生成的 Chunk 可能比flush_num_bytes更大因为触发冲刷时会把当前累积的所有行全部打包而不会按字节数精确裁切。环境变量配置三个 RERUN_FLUSH_* 变量官方文档给出了三个核心环境变量SDK 在创建批处理器时会通过ChunkBatcherConfig::from_env()/apply_env()读取实现见 crates/store/re_chunk/src/batcher.rsRERUN_FLUSH_TICK_SECS —— 时间阈值周期 tick 时长设置触发时间阈值的周期 tick 时长单位秒解析为f64浮点数。默认值RERUN_FLUSH_TICK_SECS0.2200ms例外当录制流使用网络类 sink直接向 Viewer 流式推送时默认值为RERUN_FLUSH_TICK_SECS0.0088ms——这个值对应源码中专门为实时预览设计的LOW_LATENCY配置8ms 的周期足以支撑 60Hz 级别的实时数据源码注释We want it fast enough for 60 Hz for real time camera feel。RERUN_FLUSH_NUM_BYTES —— 空间阈值字节数设置触发空间阈值的字节上限单位字节。文档记载默认值RERUN_FLUSH_NUM_BYTES10485761MiB当前仓库源码中ChunkBatcherConfig::DEFAULT的常量值为 2MiB见 batcher.rs两处数值存在差异实际生效值以你所用版本的源码为准。解析时支持两种写法人类可读的容量字符串如10MB经由re_format::parse_bytes解析或纯整数如1048576见apply_env中ENV_FLUSH_NUM_BYTES分支的实现。RERUN_FLUSH_NUM_ROWS —— 空间阈值行数设置驱动空间阈值的行数上限类型为u64整数。默认值RERUN_FLUSH_NUM_ROWS18446744073709551615即u64::MAX相当于行数维度几乎不设限仅靠字节数与时间 tick 触发。环境变量作用默认值解析方式RERUN_FLUSH_TICK_SECS周期 tick 时长秒驱动时间阈值0.2网络 sink 下0.008f64秒Duration::try_from_secs_f64RERUN_FLUSH_NUM_BYTES累积字节数上限驱动空间阈值1048576文档/ 源码常量 2MiB10MB格式或纯整数RERUN_FLUSH_NUM_ROWS累积行数上限驱动空间阈值u64::MAX纯整数另一个相关环境变量RERUN_CHUNK_MAX_ROWS_IF_UNSORTED默认 8192用于控制含未排序 timeline 的 Chunk的最大行数超过则强制拆分。它与存储侧同名环境变量保持一致源码注释 Shared with the same env-var on the store side, for consistency.旧名称RERUN_MAX_CHUNK_ROWS_IF_UNSORTED已弃用。环境变量与代码配置的优先级这是一个容易踩坑的关键点环境变量永远覆盖代码中显式指定的配置。在 crates/top/re_sdk/src/recording_stream.rs 的batcher_config()文档注释中明确写道Any environment variables as specified onChunkBatcherConfigwill always override respective settings.源码中resolve_batcher_config()recording_stream.rs的逻辑是如果调用方显式传入了batcher_config则直接使用但apply_env仍会覆盖其字段否则取当前 sink 的默认配置再应用环境变量覆盖。对应的单元测试chunk_batcher_configbatcher.rs以及recording_stream.rs中RERUN_FLUSH_NUM_BYTES456、RERUN_FLUSH_TICK_SECS456覆盖显式配置的测试都验证了这一优先级关系。代码级配置Python 与 Rust 实战示例除了环境变量你还可以在代码中直接构造ChunkBatcherConfig传入录制流。官方 snippet 分别给出了 Python 版本 与 Rust 版本两者语义完全等价于设置了如下环境变量RERUN_FLUSH_NUM_BYTESinf无限大禁用字节数触发RERUN_FLUSH_NUM_ROWS10RERUN_FLUSH_TICK_SECS10即既不打字节数主意也不靠 tick 兜底而是攒够 10 行才冲刷一次——效果是下面 10 次log调用被保证打包进同一个 Chunk。Pythonfrom datetime import timedelta import rerun as rr # Equivalent to configuring the following environment: # * RERUN_FLUSH_NUM_BYTESinf # * RERUN_FLUSH_NUM_ROWS10 # * RERUN_FLUSH_TICK_SECS10 config rr.ChunkBatcherConfig( flush_num_bytes2**63, flush_num_rows10, flush_ticktimedelta(seconds10), ) rec rr.RecordingStream(rerun_example_micro_batching, batcher_configconfig) rec.spawn() # These 10 log calls are guaranteed be batched together, and end up in the # same chunk. for i in range(10): rec.log(logs, rr.TextLog(flog #{i}))Rust// Equivalent to configuring the following environment: // * RERUN_FLUSH_NUM_BYTESinf // * RERUN_FLUSH_NUM_ROWS10 // * RERUN_FLUSH_TICK_SECS10 let mut config rerun::log::ChunkBatcherConfig::from_env().unwrap_or_default(); config.flush_num_bytes u64::MAX; config.flush_num_rows 10; config.flush_tick std::time::Duration::from_secs(10); let rec rerun::RecordingStreamBuilder::new(rerun_example_micro_batching) .batcher_config(config) .spawn()?; // These 10 log calls are guaranteed be batched together, and end up in the same chunk. for i in 0..10 { rec.log(logs, rerun::TextLog::new(format!(log #{i})))?; }Rust 端惯用ChunkBatcherConfig::from_env().unwrap_or_default()起步再从环境配置基础上微调字段兼顾环境变量可覆盖与代码默认值兜底。源码级实现ChunkBatcher 与后台批处理线程微批处理的完整实现位于 crates/store/re_chunk/src/batcher.rs核心是ChunkBatcher结构体及其专用后台线程batching_threadbatcher.rs。线程模型与消息管线ChunkBatcher可廉价克隆并跨任意线程使用内部所有操作被线性化进一条命令管线re_quota_channel通道同一线程发送的操作按其发送顺序生效多线程之间没有定义的全局顺序调用flush_blocking()可以保证调用线程此前发送的所有数据都已完成批处理并送入chunks()通道不多也不少批处理器只能通过丢弃全部实例来关闭关闭时自动冲刷管线内残留数据且关闭过程绝不阻塞。后台线程的主循环batching_thread使用crossbeam::select!同时监听两类事件命令事件AppendRow/AppendChunk/Flush/UpdateConfig/ShutdownAppendRow把行追加到对应实体路径EntityPath的累加器Accumulator中累加器按实体路径隔离因此一个 Chunk 永远不会包含多个实体路径的数据每次追加后检查空间阈值pending_rows.len() flush_num_rows或pending_num_bytes flush_num_bytes命中即以rows/bytes为原因冲刷并置skip_next_tick标记避免紧接着的 tick 空转Flush以manual为原因冲刷全部累加器并通过 oneshot 通道回执供flush_blocking等待UpdateConfig动态更新阈值并重建 tick 定时器max_bytes_in_flight例外——它必须在创建时固定运行期修改只会记录一条警告因为输入/输出配额通道在new()时已按它分配。tick 事件定时器到期以tick为原因冲刷全部累加器——这就是时间阈值的落地实现。线程结束时收到Shutdown或所有命令发送端关闭以shutdown为原因做最后一次冲刷随后关闭输出通道。行到 Chunk 的组装与拆分规则冲刷时PendingRow::many_into_chunks()batcher.rs负责把一批PendingRow组装成一个或多个合法 Chunk。为保证 Chunk 满足数据模型约束它会按以下顺序进行拆分先按RowId排序——这是全局顺序无论数据来源如何都成立按timeline 集合分组一个 Chunk 内所有行必须拥有相同的 timeline 集合通过TimePoint的确定性哈希分组TimePoint底层是BTreeMap遍历顺序确定在每组内再按组件数据类型集合分组同一个组件不能混用多种 Arrow 数据类型若某 timeline 未排序且累积行数达到chunk_max_rows_if_unsorted默认 8192则强制再拆分。对应测试batcher.rs覆盖了这些拆分条件simple验证同 timeline、同类型的行会合并进同一个 Chunksimple_static验证无 timeline 的静态数据也能批量合并simple_but_hashes_might_not_match与intmap_order_is_deterministic则验证组件插入顺序对数据类型哈希分组的影响。内置配置预设除默认值外ChunkBatcherConfig还提供四个语义清晰的预设batcher.rs可直接用于不同场景预设flush_tickflush_num_bytesflush_num_rows适用场景DEFAULT200ms2MiBu64::MAX大多数常规用例LOW_LATENCY8ms同 DEFAULT同 DEFAULT直连 Viewer 流式预览60Hz 实时感ALWAYS_TEST_ONLYDuration::MAX00测试专用每行一个 Chunk警告生产环境勿用NEVERDuration::MAXu64::MAXu64::MAX绝不自动冲刷仅手动触发两个需要警惕的配置陷阱ALWAYS_TEST_ONLY配文件 sink 会撑爆内存文件 sink 必须在进程退出时才能写 footer因此每个 Chunk 的元数据都要在内存中滞留整个录制周期。每行一刷意味着海量小 Chunk 常驻内存。recording_stream.rs的warn_if_problematic_file_sink_configrecording_stream.rs会针对这种组合发出warn_once提示生产环境请改用默认配置或LOW_LATENCY。NEVER需要显式冲刷将flush_num_bytes/flush_num_rows都设为u64::MAX、tick 设为Duration::MAX后数据只会在手动调用flush_blocking()/flush_async()对应 Rust 端RecordingStream::flush_blocking与flush_asyncPython 端为rec.flush()或 SDK 进程退出时才会下发。手动冲刷与生命周期RecordingStream层提供完整的手动冲刷 APIrecording_stream.rsflush_async()发起冲刷并立即返回不等待传播完成flush_blocking()阻塞直到冲刷完成等价于带Duration::MAX超时的flush_with_timeoutflush_with_timeout(timeout)指定超时超时或通道关闭时返回错误。冲刷分两阶段执行先同步冲刷ChunkBatcher直到所有累积行变成 Chunk 进入 chunk 通道再异步把 Chunk 冲刷进底层 sink。sink 内部也可能有独立缓冲例如 gRPC 网络 sink 有自己的flush_timeout语义。录制流被丢弃Drop时同样会先flush_blocking(Duration::MAX)再关闭保证进程退出前数据不丢失。调优建议小结实时预览、低延迟优先让 SDK 走网络 sink 的默认 8ms tickLOW_LATENCY或显式设置RERUN_FLUSH_TICK_SECS0.008录制文件、吞吐优先保持默认的字节数阈值2MiB 量级必要时用RERUN_FLUSH_NUM_BYTES16MB这类人类可读写法调大批量让每个 Chunk 更饱满减少元数据占比精确控制批量大小例如攒 N 行一起记可像官方 snippet 那样把字节数与 tick 调大、只留行数阈值——但要记住行数阈值只控制触发时机Chunk 仍可能因 timeline/数据类型差异被拆成多个环境变量是最高优先级排查线上行为与代码配置不符时先检查是否存在RERUN_FLUSH_*环境变量覆盖避免每行一刷配置对文件 sink 使用要么用默认配置要么显式手动flush防止元数据常驻内存导致内存暴涨。微批处理是 Rerun SDK 数据管线的第一道合并关口理解它的双阈值触发、配置优先级与 Chunk 拆分规则是写出既低延迟又高吞吐的多模态数据采集与可视化代码的前提。【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考