Rerun Chunk Processing 完整指南:用 Reader + Lens 管道将原始数据转换为优化的 RRD
Rerun Chunk Processing 完整指南用 Reader Lens 管道将原始数据转换为优化的 RRD【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun本指南系统讲解 Rerun 的 Chunk Processing APIrerun.experimental——介于原始数据与 RRD 文件之间的处理层reader 负责产生Chunklens 负责流式整形终端调用负责真正执行。无论你是把 MCAP 转成 Rerun 格式、用数据集构建 recording、对现有 .rrd 做后处理还是移植旧转换器本文都能帮你掌握“reader lens 管道”的正确姿势避免手工逐行组装的常见误区。阅读前建议先通过rerun-data-modelskill 确定数据应该长成什么样再按数据源选择对应的 importer skillrerun-mcap、rerun-urdf、rerun-parquet、rerun-mp4、rerun-lerobot。一、什么是 Chunk Processing数据源到 RRD 之间的管道层Chunk Processing API 位于 Rerun 的rerun.experimental命名空间其核心思想是一条惰性、列式、多线程的管道reader 产生Chunkstream 转换它们终端调用terminal call负责执行。不同数据源对应不同的 reader 与配套 skill原文档给出了一张速查表数据源Reader配套 SkillMCAP 文件ROS2、protobuf、FoxgloveMcapReader(path).stream()rerun-mcapURDF 机器人模型含 joint states → FKUrdfTree.from_file_path(...).stream()rerun-urdfParquet 表轨迹、传感器日志ParquetReader(path).stream()rerun-parquetmp4 摄像头视频Mp4Reader(path).stream()rerun-mp4HDF5 文件groups → entitiesHdf5Reader(path).stream()暂无LeRobot 数据集目录内置 importer 后接RrdReaderrerun-lerobot现有 RRDRrdReader(path)本文Sidecar 文件JSON 标定、元数据Chunk.from_columnsfrom_iter本文由于 API 处于rerun.experimental阶段当行为细节影响你的实现时建议随时核对已安装的接口面python -c from rerun.chunk import LazyChunkStream; help(LazyChunkStream)在仓库中这一整套 API 的 Python 封装位于 rerun_py/rerun_sdk/rerun/chunk/核心模块包括_chunk.py、_lazy_chunk_stream.py、_lens.py、_selector.py、_rrd_reader.py、_chunk_store.py、_lazy_store.py、_optimization_profile.py等底层均通过rerun_bindings调用 Rust 实现。二、决策规则每个组件从哪里来写任何转换代码之前先按下面的优先级过一遍。默认答案是reader 产生 chunkslens 负责整形。大多数“手工构建”的直觉在这里都是错的数据源有 reader 吗用 reader 的.stream()永远不要自己解析再重新 log。MCAP→McapReader、URDF→UrdfTree、parquet→ParquetReader、mp4→Mp4Reader、HDF5→Hdf5Reader、RRD→RrdReader、LeRobot 目录→log_file_from_path。decoder 已经产出 archetype 了吗Foxglove 直接给出Transform3D、Pinhole、VideoStream真实采样字节等现成组件——直接透传不要重新推导。只有自定义 protobuf 主题才会以Name:message形式到达才需要 lens 处理详见rerun-mcapskill。要就地修复现有组件分辨率交换、重新着色、单位换算用MutateLens配output_modeforward_unmatched。要推导新组件/新实体FK→/tf、从 message 提取标量用DeriveLens。若要把一行拆成 N 行一个 joint batch → 每个 joint 一个/tf使用双 lens 组合先用output_modeforward_all的DeriveLens推导 batch保留原始数据如 joint states再用scatterTrueoutput_modedrop_unmatched的第二个DeriveLens只输出散开的行。可参考robot_data_preprocessing示例。真正的 sidecar 数据任何 reader 或 lens 都产不出如 JSON 标定偏移、手工测量的外参、外部元数据用Chunk.from_columnsfrom_iter。最后以LazyChunkStream.merge(...)→.collect(optimizeOptimizationProfile.OBJECT_STORE)→write_rrd(application_id, recording_id)收尾。为什么强调这个顺序这条管道保持惰性、列式、多线程并且能被OBJECT_STORE优化。手工逐行循环或在 lens 之外组装pa.array会把这一切全部丢掉——这正是我们要刻意避免的路径。三、反模式清单看到左边改用右边如果你正在写左边的代码请停下来换成右边for循环逐行构建 rows/components→ 用带Selector(...).pipe(...)PyArrow-compute 回调的 lens。转换过程中每条消息都rr.initrr.log→ 那是实时日志live logging做摄入ingestion应该用 reader 读取 write_rrd。chunk.to_record_batch()pc.filter后重新用Chunk.from_columns重建手工做行稀疏化→ 用stream.drop(content...)、.split(...)或返回过滤后pa.array的MutateLens。在 lens 之外组装pa.array/pa.RecordBatch/np.frombuffer→ 把转换搬进MutateLens/DeriveLens的 selector 回调里。用自定义 parser 手工组装rr.send_columns→ 使用匹配的 reader它直接产出 chunks。用非 Rerun 库解析 MCAP/URDF 后重新 log→ 用McapReader/UrdfTree。对 reader 已经能解码的数据用Chunk.from_columns如Pinhole内参、VideoStream、来自 transforms 主题的Transform3D→ 留在 reader 流里需要修正时用MutateLens。一个实用的诊断信号当你看到一堵pyarrow.compute的 “missing-attribute” 类型错误墙pc.filter、pc.list_element时通常意味着pc.*调用散落在模块级辅助函数里而不是在Selector.pipe的 lens 回调内部。先把代码重构进 lens再考虑压制检查器——这些报错本身就是“手工构建不应存在”的提示。正在移植旧转换器手工构建的转换器诞生于 decoder 改进之前不能当作 ground truth。复制前请对照第 2 步重新验证 decoder 输出并逐条检查每一个Chunk.from_columns/ for 循环是否命中上述反模式清单。四、核心模型惰性 DAG、move 语义与两种 StoreLazyChunkStream 是惰性管道 DAG不是集合构建 filter、lenses、map、split、merge 这些操作时不会读取任何源数据。执行从终端调用开始write_rrd(...)、collect()、to_chunks()或直接迭代 stream。执行是流式、多线程且基本无 GIL 的因此优先使用 stream/lens 操作和 PyArrow compute而不是 Python 行循环。从源码看LazyChunkStream的文档字符串明确区分了两类方法Builder 方法filter、drop、split、map、flat_map、lenses、merge消费输入 stream 并返回新 stream已被消费的 stream 再次使用会抛出ValueError防止同一 stream 在管道中被重复使用。所以每一步之后都要重新赋值。Terminal 方法to_chunks、__iter__、collect、write_rrd不消费stream——它们运行管道并让 stream 保持可用每次调用都是一次全新执行。若重复执行代价过高collect()一次即可。merge的语义值得一提所有输入流并发执行chunk 按产出顺序 yield每个输入内部的 chunk 顺序保持但跨输入的顺序不确定见 _lazy_chunk_stream.py 的merge实现。ChunkStore 与 LazyStore物化 vs 按需ChunkStore是完全物化在内存中的 store通过stream.collect()或ChunkStore.from_chunks构建见 _chunk_store.py。LazyStore基于 manifest 索引按需加载 chunkRrdReader(path).store()、catalog 分段 store 都属于这一类。其 manifest 常驻内存因此schema()、summary()、__len__不加载任何 chunk 数据即可工作chunk 数据仅在请求时加载见 _lazy_store.py。两者都提供schema()、summary()、stream()和write_rrd(...)。summary()输出每 chunk 一行的确定性摘要{entity_path} rows{n} static{True|False} timelines[…] cols[…]很适合做快照测试。五、Stream 组合过滤、映射与合并所有组件都从这一行导入from rerun.chunk import Chunk, LazyChunkStream, OptimizationProfilestream.filter(content, has_timeline, is_static, components)保留每个 chunk 中匹配的部分stream.drop(...)是其补操作使用相同的关键字过滤器。content接受实体路径 glob 或 glob 列表。stream.map(fn)应用Chunk - Chunkstream.flat_map(fn)应用Chunk - Iterable[Chunk]。这是 chunk 级 Python 逻辑的逃生舱列式工作请优先用 lenses。源码注释提示map/flat_map在 Python 中运行受 GIL 约束、顺序执行见 _lazy_chunk_stream.py。stream.split(content, ...)返回(matching, non_matching)两个分支二者共享同一个上游上游只执行一次。语义上等价于(stream.filter(...), stream.drop(...))但不会重复执行上游见 _lazy_chunk_stream.py。LazyChunkStream.merge(*streams)汇入任意数量的数据源。LazyChunkStream.from_iter(chunks)包装手工构建的 chunk 列表。一个典型组合示例来自原文档stream source_stream() # 任意 importer skill stream stream.drop(content/video_raw/**) stream stream.lenses(fix_lens, content/cam/**, output_modeforward_unmatched) merged LazyChunkStream.merge(stream, sidecar_stream) merged.write_rrd(out_path, application_idmy_app, recording_idrecording_id)filter/drop的所有条件以 AND 组合对components而言chunk 会按组件列拆分只保留匹配的组件列时间线与实体路径保留列表给出时任一列匹配即保留OR 语义若 chunk 完全不含所列组件则整个 chunk 被丢弃见 _lazy_chunk_stream.py。六、手工构建 Chunk仅限 Sidecar 数据Chunk.from_columns只能用于任何 reader 或 lens 都无法产出的数据——JSON/CSV 标定、帧偏移、外部元数据。如果某个 readerMcapReader/UrdfTree/ParquetReader/Hdf5Reader能解码该主题、或某个 lens 能推导它那条路径才是惯用法不要在这里手工组装。在robot_data_preprocessing示例中唯一的手工构建 chunk 就是 JSON offsets sidecar相机修复、FK→/tf、网格、重新着色全部由 readers lenses 完成见 examples/python/robot_data_preprocessing/robot_data_preprocessing.py。Chunk.from_columns(entity_path, indexes, columns)镜像rr.send_columns(...)的 API接受相同的 archetype.columns(...)辅助函数indexes为空表示静态数据。源码层面它通过build_column_args构建时间列与组件列参数后交给 Rust 端ChunkInternal.from_columns见 _chunk.py。chunk Chunk.from_columns( /tf_static/robot_offsets, indexes[], # static columnsrr.Transform3D.columns( translationtranslations, quaternionquaternions_xyzw, parent_frameparents, child_framechildren, ), ) sidecar_stream LazyChunkStream.from_iter([chunk])非标准元数据字段可用rr.AnyValues.columns(...)覆盖。Chunk 对象还暴露entity_path、num_rows、is_static、timeline_names、to_record_batch()和format()人类可读表格等属性用于检查见 _chunk.py。Recording 属性唯一的例外Recording 属性是这里的唯一例外不要手工构建它们。Chunk.from_property(name, values)镜像rr.send_property会落一个静态 chunk 到/__properties/name。应当这样做chunk Chunk.from_property( episode, rr.AnyValues(taskpick, duration_sec12.4), )而不是手工拼装同样的 chunkChunk.from_columns( f/__properties/{name}, indexes[], columnsrr.AnyValues.columns(**{name: [value]}), # type: ignore[arg-type] )原因有二/__properties/是 Rerun 的内部布局第二个版本硬编码了它而且动态字段名会让AnyValues.columns()需要# type: ignore[arg-type]——这个抑制注释本身就是切换写法的信号。七、Lenses就地修复与派生新列Lens 在不逐行迭代的前提下整形、修复或派生组件。通过stream.lenses(lenses, output_mode..., content...)应用见 _lazy_chunk_stream.py。MutateLens(component, selector, keep_row_idsFalse)就地修改已有组件。keep_row_idsTrue时保留原始 row IDs否则生成新 row IDs见 _lens.py。DeriveLens(component, output_entityNone, scatterFalse)创建新列可选择写到另一个实体。通过链式.to_component(descriptor, selector)添加每个输出.to_timeline(name, sequence | duration_ns | timestamp_ns, selector)从数据本身提取时间列duration与timestamp可作为别名值仍以纳秒为单位见 _lens.py。scatterTrue将一行展开为 N 行每个列表元素一行。此外还有若干便捷方法to_translation(x, y, z)、to_quaternion(x, y, z, w)xyzw 顺序、to_scale(x, y, z)、to_rotation_axis_angle(axis_x, axis_y, axis_z, angle)、to_scalars(*fields)它们内部基于to_packed_component把若干 struct 字段打包成定长列表并自动 cast 到组件的规范 Arrow 类型见 _lens.py。当同一组件名存在于多个实体下时用content限定作用域。output_mode默认是 drop_unmatchedoutput_mode行为适用场景drop_unmatched默认只有 lens 输出存活纯派生中间流若广泛应用会静默丢弃其余一切forward_unmatchedlens 输出 未被任何 lens 消费的原始组件保留流其余部分的目标性修复forward_alllens 输出 全部原始组件包括DeriveLens的输入派生后仍需保留输入对 mutate lens 与forward_unmatched等价就地修复示例保持 Arrow 类型与长度不变stream stream.lenses( MutateLens( Pinhole:resolution, Selector(.).pipe( lambda res: pa.array( [(h, w) for w, h in res.to_pylist()], typeres.type, ) ), ), content[/external/cam_low, /external/cam_high], output_modeforward_unmatched, )单位换算派生示例PyArrow compute无 Python 循环DeriveLens(schemas.proto.JointState:message, output_entity/joints_deg/waist).to_component( rr.Scalars.descriptor_scalars(), Selector(.joint_positions).pipe(lambda arr: pc.multiply(pc.list_element(arr, 0), 180.0 / math.pi)), )八、Selector 语法jq 风格的 Arrow 查询Selector(query)以 jq 风格导航嵌套的 Arrow 数据见 _selector.py.当前值.fieldstruct 字段[]遍历列表元素[N]索引列表?抑制错误 / 跳过缺失的 optional!断言非空|将表达式管道进另一个表达式.pipe(fn)链接 Python/PyArrow 转换或另一个 Selector。.execute(array)立即执行.execute_per_row(array)保证输出行数与输入一致在必须保持行对齐的 lens 回调中使用。.pipe返回新 Selector 而不修改原对象注意含 Python 可调用对象的 Selector 不可 pickle——需要 pickle 时改用纯字符串 selector 或传给.pipe()一个 Selector见 _selector.py。九、写入 RRD执行与优化stream.write_rrd(path, application_id..., recording_id...)单次流式执行并写出。application_id与recording_id必须显式提供见 _lazy_chunk_stream.py。stream.collect(optimizeOptimizationProfile.OBJECT_STORE).write_rrd(...)先物化、优化 chunk 布局再写出内存随物化的 chunk 规模增长。不传optimize时只做插入时自然发生的单趟压缩见 _lazy_chunk_stream.py。两种预设 profile 在 _optimization_profile.py 中有具体数值参数LIVEOBJECT_STOREmax_bytes393,21612×8×40962,097,1522 MiBmax_rows4,09665,536max_rows_if_unsorted1,0248,192extra_passes5050gop_batchingTrueTruesplit_size_ratioNone10.0fix_keyframeFalseFalseOBJECT_STORE大 chunk面向存储/查询/catalog 服务。split_size_ratio10.0会把字节规模差异超过该因子的 archetype 组拆分到不同 chunk值越大拆分越少1.0强制每个 archetype 独占 chunk避免大列图像、视频、blob与小列标量、变换、文本同处一个 chunk使 viewer 能只拉取小列。同一 archetype 的组件总是同处一个 chunkEncodedImage:blob离不开EncodedImage:media_typeVideoStream:is_keyframe这类始终独占 chunk 的组件会先被拆出。gop_batchingTrue让视频流 chunk 按 GoP关键帧边界重排且绝不跨 GoP 拆分因此长关键帧间隔的流可能产生远超max_bytes的 chunk。fix_keyframeTrue时会丢弃用户提供的VideoStream:is_keyframe并在视频重排期间从编码样本重新推导。LIVE小 chunk面向低延迟 viewer 工作流。当 RRD 要进入 Rerun catalog 或 Hub 时始终使用OptimizationProfile.OBJECT_STORE除非被明确要求其他 profile。多个物理 RRD 共享同一recording_id时构成一条逻辑 recording利用这一点把基础数据、模型/URDF 数据、各分层分开写。十、Chunk API 与 Logging API 的分工Loggingrr.log、rr.send_columns、RecordingStream用于用户代码中的实时日志chunk processing用于摄入、转换和对现有 recording 的后处理。Logging → chunks写一个 RRD 再用RrdReader读回。RrdReader(path)列出recordings()/blueprints()每个都是带kind、application_id、recording_id的StoreEntry见 _store_entry.py.stream(storeentry)用于顺序遍历.store(storeentry)用于索引访问。需要注意没有 footer/manifest 的旧式 RRD 不支持store()此时请用RrdReader(...).stream().collect()见 _rrd_reader.py。Chunks → loggingrerun.send_chunks(chunks, recording...)接受Chunk、LazyChunkStream、LazyStore、ChunkStore或任意 chunk 可迭代对象。源 store 的application_id/recording_id不会被保留以当前活跃 recording 的身份为准。十一、常见坑位Common Gotchaslens 的默认output_mode是drop_unmatched对目标性修复忘记设置forward_unmatched会静默丢掉流的其余部分——这是最常见的坑。不要复用已消费的LazyChunkStream要么重新赋值要么刻意使用split。用content限定 lens 作用域——同一组件名常常出现在许多实体下。在MutateLens转换中保持 Arrow 数组的类型与长度不变。对 catalog 分层layer 的recording_id必须等于 segment id。这是rerun.experimentalAPI升级时核对签名。十二、端到端实战robot_data_preprocessing 示例仓库中的 examples/python/robot_data_preprocessing/robot_data_preprocessing.py 完整演示了 MCAP URDF JSON sidecar 的组合管道README 见 examples/python/robot_data_preprocessing/README.md其流水线正是上述所有概念的浓缩MCAP 摄入McapReader(DATA_DIR / episode.mcap).stream()直接产生可流式处理的 Rerun 组件。JSON sidecarChunk.from_columns构建静态Transform3D的 robot offsets chunk再用LazyChunkStream.from_iter包装成流——这是全示例唯一的手工 chunk。相机修复MutateLens交换Pinhole:resolution的宽高示例 MCAP 中外置相机标定宽高颠倒了content限定到/external/cam_low、/external/cam_highoutput_modeforward_unmatched保留其余数据。FK 派生双 lens 组合先joints_batch_lens用output_modeforward_all从schemas.proto.JointState:message推导rerun.urdf.JointTransformBatch保留原始 joint states再output_transforms_lens用scatterTrueoutput_modedrop_unmatched把 batch 展开成每个 joint 一个/tf的Transform3D行Selector(.[].translation)等逐个字段提取。URDF 流整形对左右两个机器人分别MutateLens修改Asset3D:albedo_factor着色再drop(content/robot_left/wxai/collision_geometries/**)丢弃碰撞网格。分组合并LazyChunkStream.merge把 MCAPoffsets 合为data_stream三个 URDF 流合为urdf_stream。优化写出两个流分别collect(optimizeOptimizationProfile.OBJECT_STORE).write_rrd(...)且共享同一个recording_idepisode——两个物理 RRD 因此构成同一条逻辑 recording基础数据与 URDF 分层可直接在 Rerun viewer 中打开或注册到数据集 catalog。参考与延伸本 skill 原文skills/rerun-chunk-processing/SKILL.mdChunk API 全部 Python 实现rerun_py/rerun_sdk/rerun/chunk/端到端示例examples/python/robot_data_preprocessing/robot_data_preprocessing.py各数据源专精技能skills/rerun-mcap/SKILL.md、skills/rerun-urdf/SKILL.md、skills/rerun-parquet/SKILL.md、skills/rerun-mp4/SKILL.md、skills/rerun-lerobot/SKILL.md【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考