Rerun SDK 微批处理(Micro Batching)机制详解:延迟与吞吐的平衡之道
2026/9/16 19:27:40 网站建设 项目流程

Rerun SDK 微批处理(Micro Batching)机制详解:延迟与吞吐的平衡之道

【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun

Rerun SDK 会在后台线程中自动完成数据的微批处理(micro-batching),将多条日志行聚合为 Chunk 后再交给下游 sink,从而在"数据尽快可见"(低延迟)与"减少元数据开销、提升带宽与 CPU 利用率"(高吞吐)之间寻找最佳平衡点。本文以官方参考文档 docs/content/reference/sdk/micro-batching.md 为主体,结合re_chunkre_sdk的源码实现,完整讲解微批处理的时间/空间双阈值触发机制、三个RERUN_FLUSH_*环境变量、Python 与 Rust 的代码级配置方式,以及批处理后端ChunkBatcher的真实工作流程,帮助你针对实时流式传输、录制文件等不同场景做出正确的调优决策。

为什么需要微批处理

每次log调用都会产生一行(row)数据。如果 SDK 为每一行都单独构造一个数据块(Chunk)并立刻下发,那么每个数据块都要携带一份独立的元数据(timeline 描述、组件描述符、Arrow 数组头等),对于高频、小体积的日志调用而言,元数据开销会显著吃掉带宽,CPU 也被大量浪费在重复的对象构造与序列化上。

微批处理的思想是:在后台线程中把短时间内到达的多行数据暂存起来,累积到一定规模(时间或空间阈值)后再一次性打包成 Chunk 下发。这样既保证了数据不会无限期滞留(由时间阈值兜底),又保证了批量打包带来的吞吐收益(由空间阈值触发)。

正如官方文档所述,这套机制与运行在数据存储侧的 compaction 机制 高度相似、互为映照:SDK 侧用"批处理"减少上行数据的碎片化,存储侧用"压缩合并"减少落盘后数据的碎片化。

核心机制:时间与空间双阈值,谁先触发谁生效

从源码注释可以明确看到(crates/store/re_chunk/src/batcher.rs):

The 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.rs):

RERUN_FLUSH_TICK_SECS —— 时间阈值(周期 tick 时长)

设置触发时间阈值的周期 tick 时长,单位,解析为f64浮点数。

  • 默认值:RERUN_FLUSH_TICK_SECS=0.2(200ms);
  • 例外:当录制流使用网络类 sink(直接向 Viewer 流式推送)时,默认值为RERUN_FLUSH_TICK_SECS=0.008(8ms)——这个值对应源码中专门为实时预览设计的LOW_LATENCY配置,8ms 的周期足以支撑 60Hz 级别的实时数据(源码注释:"We want it fast enough for 60 Hz for real time camera feel")。

RERUN_FLUSH_NUM_BYTES —— 空间阈值(字节数)

设置触发空间阈值的字节上限,单位字节

  • 文档记载默认值:RERUN_FLUSH_NUM_BYTES=1048576(1MiB);
  • 当前仓库源码中ChunkBatcherConfig::DEFAULT的常量值为 2MiB(见 batcher.rs),两处数值存在差异,实际生效值以你所用版本的源码为准

解析时支持两种写法:人类可读的容量字符串(如"10MB",经由re_format::parse_bytes解析)或纯整数(如"1048576"),见apply_envENV_FLUSH_NUM_BYTES分支的实现。

RERUN_FLUSH_NUM_ROWS —— 空间阈值(行数)

设置驱动空间阈值的行数上限,类型为u64整数。

  • 默认值:RERUN_FLUSH_NUM_ROWS=18446744073709551615(即u64::MAX,相当于"行数维度几乎不设限",仅靠字节数与时间 tick 触发)。
环境变量作用默认值解析方式
RERUN_FLUSH_TICK_SECS周期 tick 时长(秒),驱动时间阈值0.2;网络 sink 下0.008f64秒,Duration::try_from_secs_f64
RERUN_FLUSH_NUM_BYTES累积字节数上限,驱动空间阈值1048576(文档)/ 源码常量 2MiB"10MB"格式或纯整数
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_config(batcher.rs)以及recording_stream.rsRERUN_FLUSH_NUM_BYTES=456RERUN_FLUSH_TICK_SECS=456覆盖显式配置的测试,都验证了这一优先级关系。

代码级配置:Python 与 Rust 实战示例

除了环境变量,你还可以在代码中直接构造ChunkBatcherConfig传入录制流。官方 snippet 分别给出了 Python 版本 与 Rust 版本,两者语义完全等价于设置了如下环境变量:

  • RERUN_FLUSH_NUM_BYTES=<+inf>(无限大,禁用字节数触发)
  • RERUN_FLUSH_NUM_ROWS=10
  • RERUN_FLUSH_TICK_SECS=10

即:既不打字节数主意,也不靠 tick 兜底,而是攒够 10 行才冲刷一次——效果是下面 10 次log调用被保证打包进同一个 Chunk。

Python

from datetime import timedelta import rerun as rr # Equivalent to configuring the following environment: # * RERUN_FLUSH_NUM_BYTES=<+inf> # * RERUN_FLUSH_NUM_ROWS=10 # * RERUN_FLUSH_TICK_SECS=10 config = rr.ChunkBatcherConfig( flush_num_bytes=2**63, flush_num_rows=10, flush_tick=timedelta(seconds=10), ) rec = rr.RecordingStream("rerun_example_micro_batching", batcher_config=config) 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(f"log #{i}"))

Rust

// Equivalent to configuring the following environment: // * RERUN_FLUSH_NUM_BYTES=<+inf> // * RERUN_FLUSH_NUM_ROWS=10 // * RERUN_FLUSH_TICK_SECS=10 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_thread(batcher.rs)。

线程模型与消息管线

  • ChunkBatcher可廉价克隆并跨任意线程使用,内部所有操作被线性化进一条命令管线(re_quota_channel通道);
  • 同一线程发送的操作,按其发送顺序生效;多线程之间没有定义的全局顺序;
  • 调用flush_blocking()可以保证:调用线程此前发送的所有数据都已完成批处理并送入chunks()通道,不多也不少;
  • 批处理器只能通过丢弃全部实例来关闭,关闭时自动冲刷管线内残留数据,且关闭过程绝不阻塞。

后台线程的主循环

batching_thread使用crossbeam::select!同时监听两类事件:

  1. 命令事件AppendRow/AppendChunk/Flush/UpdateConfig/Shutdown):
    • AppendRow:把行追加到对应实体路径(EntityPath)的累加器(Accumulator)中;累加器按实体路径隔离,因此一个 Chunk 永远不会包含多个实体路径的数据
    • 每次追加后检查空间阈值:pending_rows.len() >= flush_num_rowspending_num_bytes >= flush_num_bytes,命中即以"rows"/"bytes"为原因冲刷,并置skip_next_tick标记,避免紧接着的 tick 空转;
    • Flush:以"manual"为原因冲刷全部累加器,并通过 oneshot 通道回执,供flush_blocking等待;
    • UpdateConfig:动态更新阈值并重建 tick 定时器;max_bytes_in_flight例外——它必须在创建时固定,运行期修改只会记录一条警告(因为输入/输出配额通道在new()时已按它分配)。
  2. tick 事件:定时器到期,以"tick"为原因冲刷全部累加器——这就是时间阈值的落地实现。

线程结束时(收到Shutdown或所有命令发送端关闭),以"shutdown"为原因做最后一次冲刷,随后关闭输出通道。

行到 Chunk 的组装与拆分规则

冲刷时,PendingRow::many_into_chunks()(batcher.rs)负责把一批PendingRow组装成一个或多个合法 Chunk。为保证 Chunk 满足数据模型约束,它会按以下顺序进行拆分:

  1. 先按RowId排序——这是全局顺序,无论数据来源如何都成立;
  2. timeline 集合分组:一个 Chunk 内所有行必须拥有相同的 timeline 集合(通过TimePoint的确定性哈希分组,TimePoint底层是BTreeMap,遍历顺序确定);
  3. 在每组内再按组件数据类型集合分组:同一个组件不能混用多种 Arrow 数据类型;
  4. 若某 timeline 未排序且累积行数达到chunk_max_rows_if_unsorted(默认 8192),则强制再拆分。

对应测试(batcher.rs)覆盖了这些拆分条件:simple验证同 timeline、同类型的行会合并进同一个 Chunk;simple_static验证无 timeline 的静态数据也能批量合并;simple_but_hashes_might_not_matchintmap_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.rswarn_if_problematic_file_sink_config(recording_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_blockingflush_async,Python 端为rec.flush())或 SDK 进程退出时才会下发。

手动冲刷与生命周期

RecordingStream层提供完整的手动冲刷 API(recording_stream.rs):

  • flush_async():发起冲刷并立即返回,不等待传播完成;
  • flush_blocking():阻塞直到冲刷完成,等价于带Duration::MAX超时的flush_with_timeout
  • flush_with_timeout(timeout):指定超时,超时或通道关闭时返回错误。

冲刷分两阶段执行:先同步冲刷ChunkBatcher直到所有累积行变成 Chunk 进入 chunk 通道,再异步把 Chunk 冲刷进底层 sink。sink 内部也可能有独立缓冲,例如 gRPC 网络 sink 有自己的flush_timeout语义。录制流被丢弃(Drop)时同样会先flush_blocking(Duration::MAX)再关闭,保证进程退出前数据不丢失。

调优建议小结

  1. 实时预览、低延迟优先:让 SDK 走网络 sink 的默认 8ms tick(LOW_LATENCY),或显式设置RERUN_FLUSH_TICK_SECS=0.008
  2. 录制文件、吞吐优先:保持默认的字节数阈值(2MiB 量级),必要时用RERUN_FLUSH_NUM_BYTES=16MB这类人类可读写法调大批量,让每个 Chunk 更饱满,减少元数据占比;
  3. 精确控制批量大小:例如"攒 N 行一起记",可像官方 snippet 那样把字节数与 tick 调大、只留行数阈值——但要记住,行数阈值只控制触发时机,Chunk 仍可能因 timeline/数据类型差异被拆成多个;
  4. 环境变量是最高优先级:排查线上行为与代码配置不符时,先检查是否存在RERUN_FLUSH_*环境变量覆盖;
  5. 避免每行一刷配置对文件 sink 使用:要么用默认配置,要么显式手动flush,防止元数据常驻内存导致内存暴涨。

微批处理是 Rerun SDK 数据管线的第一道"合并关口",理解它的双阈值触发、配置优先级与 Chunk 拆分规则,是写出既低延迟又高吞吐的多模态数据采集与可视化代码的前提。

【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询