具身智能数据平台三大件:清洗、数据集管理与评测分析怎么落地
具身智能数据平台三大件:清洗、数据集管理与评测分析怎么落地
本文是 具身智能的数据地基:从传感器带宽到湖仓、清洗、评测与回流闭环 的落地实操篇——那篇讲全景与选型,这篇把三大件怎么写代码、怎么衔接拆透。
具身智能的数据工程,说到底就是三件事:把脏数据洗干净(Spark + Daft)、把干净数据管起来(Iceberg + LeRobotDataset)、让数据可查可分析(Doris / StarRocks)。这三个环节分别对应数据生命周期的不同阶段,解决的是完全不同性质的问题——如果混淆了层次,很容易在选型上踩坑:比如试图用 Spark 直接管理数据集版本,或者让 OLAP 引擎承担数据清洗的重活。
这篇文章把这三个模块逐一拆开,讲清楚每个工具在具身场景下的定位、用法和它们之间的衔接方式。
一、Spark + Daft:具身数据清洗管线怎么写
1.1 为什么需要 Daft,Spark 自己不够吗
Spark 擅长海量结构化数据的 ETL,但具身智能的数据主体是多模态非结构化数据——视频、图像、点云、张量。用原生 Spark 处理视频,你需要自己写 UDF 管理解码器、处理分片、对齐时间戳,工程成本很高。
Daft 是专为 AI 时代设计的分布式 DataFrame 引擎,底层 Rust 实现,对外提供 Python API。它的核心设计是把 Image、Video、Audio、Tensor 当作一等公民类型——一行表达式就能完成分布式视频解码、帧提取、图像变换,不需要自己管理分片调度。
它和 Spark 不是竞争关系,而是分工关系。在阿里云 EMR Serverless Spark 上,Daft 可以作为引擎直接运行,与 Spark 任务共用同一套资源与调度;也可以本地单机跑小规模验证。
对具身场景来说,Daft 最大的价值是数据处理即 AI 推理:通过 ai_query、ai_embedding 等 AI Function,把 VLM 打标、多模态向量化直接写进 DataFrame 表达式——模型调用不再需要单独的脚本去拆数据、对齐结果。这一点在后面打标环节会具体展开。
1.2 一条完整的具身数据清洗管线
一条机器人轨迹从湖里读出来,到成为可训练样本,通常经过五个步骤。下面用代码走一遍。
第一步:读取与解码。 从对象存储的 glob path 加载视频文件:
1 | import daft |
第二步:分布式抽帧与图像变换。 Daft 的图像操作都是表达式级的,集群自动并行:
1 | # 抽帧:按时间间隔提取关键帧 |
一个工程细节值得注意:image.decode() 支持 on_error="null" 参数,坏帧不会让整个任务挂掉,而是置空后由后续步骤过滤。具身采集现场环境复杂,镜头脏污、码流损坏时有发生,这个容错设计在量产管线上是刚需。
第三步:脏数据过滤。 静止帧、模糊帧、尺寸异常帧的过滤:
1 | df = df.with_column("width", df["image"].image.width()) |
第四步:VLM 打标。 这是 Daft 区别于传统 Spark 的关键能力——给每帧打物体类别、场景、任务阶段标签,直接在 DataFrame 链路里完成:
1 | from daft.ai import openai # 或通义千问 |
如果用传统方案,这一步需要先把数据导出、调模型服务、再把结果 join 回来,中间的序列化和对齐开销很大。Daft 把模型调用内嵌进表达式,引擎负责并发调度和失败重试。
第五步:写回数据湖。 清洗后的帧、标签、元数据写回 Parquet/Iceberg,供训练侧流式读取。
整个管线的逻辑可以概括为:粗筛(结构化过滤)→ 细筛(图像质量)→ 增强(打标/向量化)→ 落湖。每一步都是声明式表达式,底层并发、容错、血缘由引擎接管。
二、Iceberg + LeRobotDataset:训练数据集管理
2.1 两个工具解决的是两个层次的问题
这是这个环节最容易混淆的地方。LeRobotDataset 定义"机器人轨迹长什么样"(逻辑格式),Iceberg 定义"这些文件怎么被事务性管理"(物理管理层)。理解了这层关系,两者的结合方式就清晰了——它们是互补而非替代。
2.2 LeRobotDataset v3.0 的格式设计
v3.0 相比 v2.1 的核心变化是存储与访问解耦:不再强制一个 episode 一个 Parquet、一个视频文件,而是把多个 episode 合并进大文件,用元数据记录每个 episode 在聚合文件中的精确偏移量。目录结构长这样:
1 | my_dataset/ |
episode 元数据里记录了每个 episode 的定位信息:dataset_from_index / dataset_to_index(在 Parquet 里的行区间)、videos/<key>/from_timestamp / to_timestamp(在视频文件里的时间区间),loader 据此实现零拷贝读取。
这个设计的收益是实打实的:文件数量减少 90%,初始化加载速度提升 3-5 倍,存储占用减少 30%。当数据集规模到 OXE(Open X-Embodiment)那个量级——数百 GB 起——这些差异会直接决定训练任务的启动效率。
2.3 Iceberg 在这之上叠加的能力
裸的 LeRobot 格式是文件集合,没有事务、没有版本、没有并发控制。Iceberg 注册进来后提供四个关键能力:
- ACID 事务 + 快照隔离:每批新 episode 写入产生一个新快照,清洗任务和训练任务并发读写互不干扰;
- Schema Evolution:新增传感器通道(比如从 2 相机升级到 4 相机)只需
ALTER TABLE ADD COLUMN,不重写历史数据; - Time Travel:
SELECT * FROM table TIMESTAMP AS OF '2026-08-01'或直接 pin 一个 snapshot ID,训练结果可精确复现; - 分区演进:可以按机器人 ID / 任务类型 / 采集日期重新分区,旧查询不受影响。
其中 Time Travel 对具身训练尤其关键。训练实验的可复现性依赖数据版本的确定性——两周前那次实验效果好,今天用"同一份数据"重跑效果变差,如果数据集没有版本管理,你根本无法定位是模型问题还是数据被增量污染了。Iceberg 的 snapshot 机制把这个问题从流程约定变成了技术保证。
2.4 实际工作流
结合两者的典型数据集管理流程:
1 | 写入 新一批 episode(Parquet + MP4)由采集端产出 |
一句话概括这个分层:LeRobot 格式管数据长什么样,Iceberg 管数据怎么被安全地读写和版本化。
三、Doris 和 StarRocks:评测分析层入门
3.1 先搞清楚 OLAP 是什么
OLAP(Online Analytical Processing,联机分析处理)数据库专门为分析型查询设计:列式存储、海量数据扫描、多维聚合、毫秒到秒级响应。与之相对的是 OLTP(如 MySQL),面向高频小事务。
具身场景里什么时候需要 OLAP?当你做这些事的时候:评测结果的多维交叉统计(任务类型 × 物体类别 × 场景光照 × 机器人本体的成功率矩阵)、失败模式分布分析、回归测试对比、训练任务产出的指标报表。数据量到亿级以上、要做多维 GROUP BY、MySQL 已经跑不动时,就是引入 OLAP 引擎的时机。
3.2 Doris:架构与核心概念
Apache Doris 是一个 MPP(Massively Parallel Processing,大规模并行处理) 架构的实时分析数据库,只有两类进程,运维极简:
- FE(Frontend,Java 编写):负责元数据管理、SQL 解析、查询规划。内部又分 Leader(写)、Follower(元数据副本、可选举)、Observer(只读扩展查询能力),通过 Raft 协议保证高可用;
- BE(Backend,C++ 编写):负责数据存储和查询执行,多副本机制保证可靠性。
对外的接口是 MySQL 协议——BI 工具、Grafana、Python 的 pymysql 都能直连,学习成本很低。
Doris 建表时需要选数据模型,这是新手最容易困惑的地方:
| 模型 | 语义 | 典型场景 |
|---|---|---|
| Duplicate Key(明细) | 保留所有原始行,Key 仅用于排序 | 日志、行为事件明细存储 |
| Aggregate Key(聚合) | 相同 Key 的行按指定方式聚合(SUM/MAX/…) | PV/UV 汇总、指标预聚合 |
| Unique Key(主键) | 相同 Key 的行只保留最新值 | 订单状态、用户画像、评测结果 |
其中 Unique Key 模型在 Doris 2.0+ 引入了 Merge-on-Write(写时合并):写入阶段就完成主键合并、旧数据打删除标记,查询时直接读最终生效数据,性能接近普通明细表,更新毫秒级可见——这是它适合"实时更新 + 高频查询"场景的根本原因。具身评测场景里,评测结果表需要持续 Merge 新数据,这个模型正好对口。
3.3 StarRocks:和 Doris 什么关系、强在哪
StarRocks 早期从 Doris(Palo)分支出来深度重构,架构同源(FE/BE、MySQL 协议、MPP),但在湖仓查询方向走得更深。
它对具身场景最有价值的能力是 External Catalog:可以直接把 Iceberg 数据湖挂载进来当原生表查,无需数据迁移。查询流程是先读 Iceberg 的 Snapshot 定位 Manifest、做分区裁剪,再分发到计算节点执行——本质上是用数仓级优化器去加速湖上查询。
4.0 版本更进一步,把数据治理理念引入写入端:写入前智能路由避免冲突、后台 Compaction 服务持续合并小文件,让 Iceberg 表"落地即可查"。这与第二节讲的 Iceberg 数据集管理形成了自然衔接——训练数据在湖上被 Iceberg 管理,评测分析通过 StarRocks 直接查湖,不需要额外的同步链路。
1 | 训练数据湖 (Iceberg 管 LeRobot Episode) |
3.4 选型逻辑
两者功能高度重合,务实的选择逻辑是:
- 数据主要在湖上(Iceberg/Hive),要做联邦查询、交互式报表 → StarRocks 的 External Catalog 和湖仓优化更成熟;
- 数据直接写入引擎、需要高频实时更新(如评测结果持续 Merge)、追求极简运维 → Doris 的 Merge-on-Write 和内嵌 Raft 元数据更省心;
- 已有 Doris 存量、团队能力沉淀在这边 → 没必要迁移,Doris 同样支持 Iceberg 外表查询,只是湖仓方向的迭代节奏略慢于 StarRocks。
四、三个模块怎么串成一条链路
把三部分连起来,具身数据平台的完整数据流是:
1 | 采集端产出的原始轨迹 |
每个环节的职责边界很清楚:清洗层管"数据质量",数据集管理层管"版本与一致性",分析层管"可观测与决策支持"。当其中某个环节成为瓶颈时,才考虑引入更重的方案——比如数据量没到 PB 级之前,对象存储 + Parquet + LeRobot 格式就够用,Iceberg 的引入时机是并发读写或多引擎访问成为实际痛点的时候。
工具选型的原则始终是:先想清楚这一层要解决什么问题,再挑手合适的家伙。
参考
- HuggingFace LeRobot Documentation(LeRobotDataset v3.0)
- Apache Iceberg 表格式规范
- Apache Doris / StarRocks 官方文档
- Daft 分布式 DataFrame 引擎
- 本文姊妹篇:具身智能的数据地基:从传感器带宽到湖仓、清洗、评测与回流闭环