
Rerun Catalog 查询性能实战基于 rerun-catalog-queries 的 Python 数据管道调优指南【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun本篇技术指南以 Rerun 仓库中skills/rerun-catalog-queries/SKILL.md为骨架系统讲解如何用 Python 高效查询 Rerun catalog调用链为rerun.catalog.CatalogClient→dataset.reader(...)→ DataFusionDataFrame。读完你将掌握如何以往返次数优先、载荷字节其次的成本模型评估查询性能如何用filter_contents、时间窗口过滤、cache()、跨 segment/实体批量拉取等手段把数十次云上往返压缩到一次以及一套可复现的慢查询调试流程直接用于 per-segment / per-episode 的数据处理管道。一、背景与调用栈Catalog 查询在 Rerun 中的位置Rerun 是一个面向多模态机器人数据的可视化、查询与流式训练平台。当数据以.rrdRerun Recording Data形式注册到 catalog 服务器后Python 侧可以通过rerun.catalog.CatalogClient连接远程 catalog 服务器再通过DatasetEntry派生出DatasetView最终由reader()返回一个 DataFusionDataFrame供下游做过滤、聚合与物化。从源码结构看这条链路的 Python 侧实现位于rerun_py/rerun_sdk/rerun/catalog/_catalog_client.py中的CatalogClient负责连接远程 catalog 服务器提供datasets()、get_dataset()、benchmark()、ctxDataFusionSessionContext等能力。初始化时会校验本地datafusion包版本是否与 FFI 兼容当前兼容集见文件内DATAFUSION_MAJOR_VERSION_COMPATIBILITY_SETS。_entry.py中的DatasetEntry与DatasetViewfilter_segments()、filter_contents()、reader()全部在此实现。_schema.py中的Schema提供index_columns()、component_columns()等 schema 内省能力。_content_filter.py中的ContentFilter不可变的实体路径过滤器构建器。本指南聚焦于 catalog 特有的行为与往返成本round-trip cost—— 这正是很多团队在编写 per-segment / per-episode 管道时容易踩坑的地方。DataFusion 侧的DataFrame/SessionContext/ 表达式 API 细节属于datafusion-python技能的范畴本文只在必要时顺带涉及。二、查询成本模型一句话概括每一次dataset.reader(...)DataFrame 的物化都是一次云上往返round-trip其中大部分成本是网络 解码这对组合而不是计算本身。规划时按往返次数 → 载荷字节数的顺序进行优化。一次典型的 catalog 往返即便结果极小也需数秒。因此30 个 segment × 每个 1 次查询 ≈ 90 秒朴素的 per-segment 循环30 个 segment × 每个 4 次查询 ≈ 6 分钟例如一个 splitter 对起点和终点分别执行countcollect_column1 次覆盖全部 30 个 segment 的查询 ≈ 3 秒。同样的扇出fan-out问题也存在于实体entity维度10 个实体 × 各自执行 1 次filter_contents([one_entity]).reader()≈ 30 秒而一条filter_contents([all_entities]).reader()≈ 3 秒。结论尽可能把工作压进同一次往返 —— 无论是跨 segment 还是跨实体。这一点在源码里也有印证DatasetEntry.reader()_entry.py内部会先执行self.filter_contents([/**])再构造DatasetView而DatasetView.reader()最终调用self._internal.reader(...)Rust 侧实现一次调用对应一次服务端数据读取。你写多少个reader(...)终端调用就产生多少次云上往返。三、最大杠杆.reader(...)之前先应用filter_contents与时间窗口过滤这是唯一最重要的优化手段filter_contents([entity_globs])限定 reader 生成的实体路径列。不写它数据集里每个实体都会被读到。对 Scalars 类型的列这一过滤还能把数组嵌套深度从listlistdouble降为listdouble。时间窗口过滤df.filter(col(index).cast(int64) start).filter(col(index).cast(int64) end)会被下推到存储层大幅减少扫描的字节数。顺序很关键先过滤再做 reader 绑定后的投影projection顺序不能颠倒。组合起来的标准形态dataset.filter_segments(seg).filter_contents(entities).reader(...).filter(in_window).select(...)源码侧佐证DatasetView.filter_contents()接受带通配符的实体路径表达式/points/**匹配/points下所有实体-/text/**排除/text下所有实体也接受ContentFilter流式构建器见_entry.py与_content_filter.py。ContentFilter.everything()等价于filter_contents(/**)ContentFilter.nothing()等价于filter_contents([])__properties子树默认被自动排除需要时用include_properties()显式纳入。另外DatasetView.reader()的 docstring_entry.py明确说明返回的 DataFrame支持对rerun_segment_id和索引timeline列的服务器端过滤例如dataset.reader(indexreal_time).filter(col(rerun_segment_id) aabbccddee)或(col(rerun_segment_id) ...) (col(real_time) ...)都会被 Rerun 服务器在远端执行而非拉回客户端。这意味着把过滤条件写成可下推的形式本身就是在省钱。四、df.cache()重复探测时的好朋友当同一个物化结果被多个下游 filter/count/collect 调用复用时先用DataFrame.cache()物化一次再在缓存帧上操作cached ( dataset .filter_segments(seg) .filter_contents([entity]) .reader(indexindex_col) .select(col(index_col).cast(pa.int64()).alias(index_col), value.alias(v)) .cache() # 一次网络往返物化为内存中的 batches ) starts cached.filter(col(v) start_val).collect_column(index_col) stops cached.filter(stop_pred(col(v))).collect_column(index_col)没有cache()时每次count()/collect_column()都会重新执行整条 reader 链路即再次发起云上往返。什么时候不要 cache。cache()强制把数据物化成 Arrow batches会打破惰性laziness。如果下游还在继续叠加 DataFusion 算子join、window、更多过滤并且只在最后物化一次那么在管道中途 cache 会把一次执行变成两次执行还会抢先阻断引擎跨越边界可能做的物理计划优化。cache()的适用场景是消费者是终端操作count()、collect_column()、to_arrow_table()而不是消费者本身又是一个惰性DataFrame。五、跨 segment 批处理丢弃filter_segments按rerun_segment_id分组对于需要在很多 segment 上执行同一查询的管道干脆不要调用filter_segments(...)直接拉一张跨 segment 的完整表。每一行 reader 输出都带一个rerun_segment_id列 —— 在本地分组即可df dataset.filter_contents(entities).reader(indexindex_col) cached df.select( rerun_segment_id, col(index_col).cast(pa.int64()).alias(index_col), value.alias(v), ).cache() # 现在 N 次 filter/aggregate 调用都是本地的不再走网络。 starts cached.filter(col(v) start_val).select(rerun_segment_id, index_col).to_arrow_table()Trigger / event 这类列通常足够小一次拉取全部 segment 相比逐个 segment 循环性能可提升一个数量级。源码侧佐证filter_segments本身也可以接受一个 DataFusion DataFrame含rerun_segment_id列在_entry.py中会先select(rerun_segment_id).to_pydict()取出 id 列表再构造 view —— 这常用于先用 segment 元数据表筛出好的 segment再据此查询的场景。而跨 segment 单次拉取则完全绕开这个 id 列表把过滤职责交给本地 DataFrame是本文推荐的高吞吐路径。六、segment 内的 per-entity 扇出这是跨 segment 批处理的对称问题只不过沿实体轴展开。如果只需要单个 segment 里 N 个实体的数据不要循环# 反模式N 次 reader 构建N 次往返。 for entity in entities: df dataset.filter_segments(seg).filter_contents([entity]).reader(indexix) ...正确做法是一次性把 N 个实体全部拉下来然后在本地按实体做投影。依据Reader 行布局见下文第八节每一行只携带一个实体的数据其余实体的列在该行上是 NULL所以col(entity:archetype:component).is_not_null()就是天然的 per-entity 过滤器shared ( dataset .filter_segments(seg) .filter_contents(sorted(set(entities))) .reader(indexix) .filter(col(ix).cast(pa.int64()).between(start_ns, end_ns)) ) # 每个下游消费者惰性地收窄到自己的实体行。 src_a shared.filter(col(f{ent_a}:{comp_a}).is_not_null()).select(ix, f{ent_a}:{comp_a}) src_b shared.filter(col(f{ent_b}:{comp_b}).is_not_null()).select(ix, f{ent_b}:{comp_b})DataFusion 在构建物理计划时可以跨这些 per-entity 投影共享底层扫描所以即使有 N 个逻辑消费者这仍然只是一次 catalog 往返。该模式可以无缝放进按数据源逐个构建 DataFusion 计划的生成器里重采样、bracket 查找、nearest-in-time join 等—— 只消除网络扇出不改变 per-source 逻辑。重构时的一个陷阱如果下游查询使用 reader 列的全限定名col(f{entity}:{archetype}:{component})你不需要在shared里给列起别名 ——sharedreader 的输出 schema 保留了原生列名既有的 per-entity 投影辅助函数可以原封不动地继续工作。七、count()并不免费反直觉的一点df.count()和df.aggregate([], [F.count(col)])并不总能下推。聚合计划可能迫使引擎在服务端物化底层列数据然后在客户端计数。F.count(col)作用于宽实体列时甚至可能把完整的 struct 或 blob 载荷都搬过来只为统计非空性。针对这里到底有没有数据这类问题按优先级选择需求做法该过滤条件下是否有任何行df.select(col(index)).limit(1).to_arrow_table().num_rows 0—— 服务器在首个匹配处短路极小时间窗口内的行数时间过滤后df.filter(window).select(col(index)).count()宽查询中每个实体的计数对每个实体做limit(1)探测并并发执行 ——不要用一个大的count(col)聚合一个曾坑过作者的陷阱以为bool_or(col.is_not_null())只需要 nullity buffer。事实并非如此 —— 在大多数物理计划上该算子仍然会触碰载荷数据。八、using_index_valuesfill_latest_at适合重采样不适合存在性判断.reader(indexindex_col, using_index_valuestargets, fill_latest_atTrue)每个目标时间戳返回一行各实体列携带截至该目标时刻最近的非空值latest-known value。非常适合最近前值重采样nearest-prior resampling且无需 DataFusion。不要用它做存在性presence判断原因有二语义是T 之前曾经发出过而不是在 [start, T] 窗口内发出过服务器为了计算 latest-known value仍然会为每个实体传输完整的 struct/blob 载荷下游的is_not_null()投影发生在传输完成后不会减少线上字节数。严格窗口内存在性检查推荐对时间过滤后的 reader 做 per-entitylimit(1)探测并并发执行。源码侧佐证DatasetView.reader()的using_index_values参数_entry.py支持三种形态 ——普通数组仅应用于索引范围覆盖该值的 segment需要扫描 segment 表做映射segment 多时开销大、dictkey 为 segment idvalue 为该 segment 的采样值客户端已知时优先、DataFrame必须含rerun_segment_id和索引列。docstring 同时提醒未知 segment id 会被静默忽略普通索引切片不要用它改用索引列上的 DataFusion filter如(col(real_time) lit(t0)) (col(real_time) lit(t1))。九、Schema 内省很便宜探测之前先用它schema dataset.filter_segments(seg).schema() available {(c.entity_path, c.component) for c in schema.component_columns()}这只是一次往返就能告诉你该 segment 注册了哪些(entity, component)组合。如果某列不在 schema 里你就可以直接把它从 manifest 里剔除不需要再发起任何云上查询。在很多场景下这完全可以替代该实体是否有事件的探测。⚠️ 注意schema 存在 ≠ 有事件。MCAP 即使对未使用的 topic 也会记录 topic schema。如果你的管道需要区分已注册但从未发出与已注册且有事件那就必须探测 —— 参见上文第七节的 is anything here 模式。源码侧佐证Schema_schema.py提供component_columns()、index_columns()、column_for(entity_path, component)、columns_for(...)可按实体路径 / archetype / component type 过滤、entity_paths()、archetypes()等方法DatasetView.schema()会反映已应用的 content filters所以先filter_contents再看 schema得到的是裁剪后的列集合。十、Reader 行布局实体是列不是行dataset.filters.reader(...)返回的每一行对应单个实体上的单个事件其他实体的列在该行上是 NULL。推论select(rerun_segment_id, entity:archetype:component)是合法的 —— 在 SQL 中使用实体列时需要加引号。没有rerun_entity_path行属性。要把行归属到实体要么一次过滤一个实体要么选一个 per-entity 列例如:McapChannel:id用它的非空模式来标识来源。df.count()返回的是跨所有实体的总事件数而不是 per-entity 计数。这也是第六节 per-entity 扇出方案能够成立的根本原因 ——is_not_null()过滤正是利用了一行只有一个实体的值这一布局。十一、null struct 上的字段访问返回0.0而不是 null这是一个会咬到读取 struct 消息管道的 DataFusion 陷阱col(/some/entity:msg.MyType:message)[0][sub][x] # 当父 struct 在该行上为 null 时这个表达式求值为 0.0 # 字符串则是 而不是 null。给 struct 遍历投影包一层 null 保护parent col(/some/entity:msg.MyType:message) leaf parent[0][sub][x] guard parent.is_null() | parent[0].is_null() | parent[0][sub].is_null() expr F.when(guard, lit(None)).otherwise(leaf)只对 struct 数据源应用此保护。对标量 / blob 列Scalars:scalars、EncodedImage:blob等这层包装最好也只是空操作最坏会以改变下游 join 行为的方式重写计划。是否加保护取决于数据源是否真的穿越了 struct 边界。十二、常见调试配方Debug Recipe当某个查询阶段比预期慢时按以下顺序排查数往返次数。给每个to_arrow_table()/collect_column()/count()包上time.perf_counter()。每一个都是一次往返。如果看到 N 次 segment 查询那么最低成本就是 N × 数秒。拆分构建与物化计时。惰性 DataFrame 的构建时间与终端的to_arrow_table()调用时间要分开测量。如果构建耗时数秒说明内部有东西在急切物化 —— 可能是 join 辅助函数里的 Arrow 往返、生成器里的cache()、或链式 join 工具里隐藏的.collect()。一个正确的惰性计划无论结果多大构建都应在毫秒级。测量字节数。to_arrow_table()后的tbl.nbytes能揭示我以为做了is_not_null()投影实际却传输了数 MB 的情况。如果投影很小但字节很大说明算子没有下推。先跨 segment再跨实体。如果 per-segment 查询本质上相同只是按 id 限定范围丢掉filter_segments并在本地按rerun_segment_id分组如果 segment 内还有 per-entity 循环同样折叠一条filter_contents([all])reader 下游 per-entityis_not_null()过滤。终端复用前先 cache。如果两个count()/collect_column()调用共享同一个 reader在它们之间加df.cache()但如果消费者本身是还要继续组合的惰性 DataFrame不要cache —— 缓存会破坏计划级优化。窗口优先。任何触碰载荷列的投影或聚合之前永远先推时间过滤。这条流程与上文各节一一对应第 1 条对应第二节成本模型第 4 条对应第五、六节第 5 条对应第四节第 6 条对应第三节。十三、速查清单Cheat Sheet场景正确姿势关键收益30 个 segment 各查一次1 次跨 segment reader 本地按rerun_segment_id分组90s → 3sN 个实体各查一次1 次filter_contents([all]) per-entityis_not_null()N 次往返 → 1 次多个终端消费同一 reader先.cache()再 count/collect避免整链重执行是否有数据select(col(index)).limit(1).to_arrow_table().num_rows 0服务端短路重采样到目标时间戳using_index_valuestargets, fill_latest_atTrue免 DataFusion时间窗口切片索引列 cast int64 后 filter先于投影下推到存储层判断列是否存在schema.component_columns()一次往返省去无效探测null struct 字段访问F.when(guard, lit(None)).otherwise(leaf)避免 0.0/ 假值十四、延伸阅读DataFusion 侧的 DataFrame API、SQL 等价性、表达式构建与常见陷阱布尔运算符、不可变性等属于datafusion-python技能的范围。若未安装可通过npx skills add apache/datafusion-python全局安装后加载。若想深入了解本文涉及的 Python API 完整签名与 docstring可继续阅读仓库内rerun_py/rerun_sdk/rerun/catalog/_entry.pyDatasetEntry/DatasetView/reader/using_index_values、rerun_py/rerun_sdk/rerun/catalog/_catalog_client.pyCatalogClient/benchmark、rerun_py/rerun_sdk/rerun/catalog/_schema.py与rerun_py/rerun_sdk/rerun/catalog/_content_filter.py。【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考