首页 / 资讯中心 / 文章详情

Daft、Ray、Lance三件套:构建AI数据管道的现代方案

Daft、Ray、Lance三件套:构建AI数据管道的现代方案 ★ FEATURED ARTICLE
1. 这次课程为什么要把 Daft、Ray、Lance 放在一起讲做数据处理这一行的人多数时间都在跟三类东西较劲计算引擎怎么选、任务怎么调度、数据落盘之后怎么保证还能快读快查。过去我们习惯把这几个问题分开解决——用 Spark 管计算用 Airflow 管调度用 Parquet 管存储。这套组合拳本身没问题但在 AI 场景里越来越显得别扭。特征工程要上向量检索训练样本要频繁版本回溯数据管道要支持任意语言写的 UDF传统数仓工具链越拼越长维护成本直线上升。所以我重新整理这套专题课的时候第一反应就是把 Daft、Ray、Lance 塞进同一个教学框架里。三个工具不是硬凑的它们分别卡在现代数据基建的三个关键位置Daft 负责分布式 DataFrame 计算Ray 负责分布式任务编排和状态管理Lance 负责列式存储和向量索引。单独学任何一个都能讲得热闹但只有把它们拼起来学员才能看到一条完整的数据管道长什么样也才能理解为什么这三样东西在 2025 年会同时被频繁提起。这期课程面向的读者主要有三类一是被 Pandas 大数据量处理折磨的同学二是想给团队搭建 AI 数据平台但不想再引入一堆开源组件的大夫三是对 Lance 这类新一代存储格式好奇、想把它用在生产里的工程师。你不需要事先精通 Ray 或 Rust但最好对 Python 数据处理有手感知道 DataFrame 是什么否则第一部分会稍微吃力一些。课程设计的核心思路其实就一句话让每一个工具做它最擅长的事再用一个真实项目把它们串起来。Ray 不抢 Daft 的活Daft 也不抢 Lance 的活三者的边界划分清楚之后你会发现整个数据管道的复杂度瞬间下降了一个档次。2. 核心工具逐个拆解Daft 的计算模型与 Ray 的调度体系2.1 Daft面向数据科学家的分布式 DataFrame 引擎Daft 给我的第一印象是一个字快。它和 Spark 最大的区别在于底层用 Rust 实现同时借用了 Arrow 列式内存格式所以在单机多核的场景下表现相当惊艳。我做过一个测试处理 1 亿行 CSV用 Pandas 直接读的话内存直接爆掉用 Daft 的read_csv配合分区读取几秒就能完成全集扫描内存占用还被控制在合理范围。Daft 的 API 设计非常贴近 Pandas我之前带过一个完全没用过 Spark 的学员切到 Daft 之后几乎没有学习成本。核心概念可以理解为一个支持惰性求值的分布式 DataFrame你写df.filter(...).select(...)它不会立刻执行而是构建一棵逻辑计划树等真正触发动作时才交给优化器做谓词下推、列裁剪和分区剪枝。需要注意 Daft 并不是要取代 Ray它自己也有分布式执行能力但团队一直在强调 Daft 本身是计算引擎而非调度平台。在这个专题课里我们是把 Daft 当成“被调用方”来用的Ray 负责切分任务、分发数据Daft 在每一个工作节点上处理分配给它的那批数据。这样分工的好处是每个组件的职责单一调试问题的时候不会牵扯进一堆互相纠缠的配置。2.2 Ray从任务编排到分布式对象存储Ray 这个框架自从在 OpenAI 内部被大量使用之后已经逐渐成为 Python 分布式生态的事实标准。它不是单纯的任务队列而是提供了一整套并行原语远程函数通过ray.remote装饰器声明调用后立刻返回一个ObjectRef这个引用指向一个分布式的对象存储——确切地说是每个 Ray 节点本地内存里的共享内存对象。数据在节点间通过这把“分布式钥匙”流转避免了反复写入磁盘带来的 I/O 开销。课程里我会用很大篇幅来讲 Ray 的 Actor 模型。简单说ray.remote修饰一个类就成了 Actor它维护独立的状态可以并发接收多个调用。这个特性在处理有状态的数据流时很有用比如某个特征工程模块需要缓存上一批结果直接用 Actor 包一层就行。Ray 生态里还有 Ray Data、Ray Train、Ray Serve 等组件但在专题课里我们只深入使用 Ray Core。原因很直接Ray Data 目前对自定义数据源的支持还没有 Daft 灵活Ray Train 更适合深度学习场景而非通用数据处理Ray Serve 则属于服务化部署范畴单开一门课都不嫌多。我们课程的边界定在数据处理和存储上所以 Ray Core 刚刚好。2.3 课程中如何安排两套系统协作如果你去看官方文档Daft 也支持直接用 Python runner 跑多线程Ray 更像是一个可选的分布式执行后端。我们在课程里做了一个明确的选型决策默认使用daft.set_runner_ray()把 Daft 的底层执行切换到 Ray 集群上。这样 Daft 的数据分区、聚合、join 等算子会通过 Ray 分发到不同节点而我们编写的数据处理主逻辑仍然保持 Pandas 风格的简洁调度。但这中间有一个绕不开的痛点——数据序列化。Ray 的远程函数参数和返回值默认通过 Arrow 格式传递如果 Daft DataFrame 的内部表示没有走 Arrow就会触发 Python pickle 序列化导致性能断崖下跌。所以课程会专门教大家检查每一层数据流确保中间产物尽量是 Arrow Table 或 Lance 的 native 结构而不是 Python list 或 dict。这个细节在本地跑不出来区别一旦集群规模上了 4 节点以上差距立刻明显。3. Lance 文件组织结构深度解析课程重头戏3.1 一句话搞清楚 Lance 是什么Lance 是一种面向 AI 工作负载的列式存储格式最直白的理解是“升级版 Parquet”它支持列存压缩、向量索引、快速随机访问和多版本管理。Parquet 的强项是静态数据的高压缩比但一旦你要频繁追加数据、增量更新、或者在上面做语义化版本切换Parquet 会表现得非常笨重。Lance 从设计上就把这些场景作为第一优先级。我记得第一次看到 Lance 的目录结构时愣住了——它跟普通文件格式很不一样不是单个.lance文件而是一个包含多个子文件和清单文件的文件夹。这对刚接触的同学来说需要调整认知Lance 的“文件”其实是一个“数据集目录”目录本身内置了版本、索引和统计信息。3.2 Lance 目录与文件结构逐层拆解来落地一下假设你创建了一个 Lance 数据集并写了几个批次的数据目录长这样my_dataset.lance/ ├── _latest.manifest ├── _versions/ │ ├── manifest.0 │ ├── manifest.1 │ └── manifest.2 ├── data/ │ ├── fragment-0.lance │ └── fragment-1.lance ├── indices/ │ ├── index-0.lance │ └── index-1.lance └── _deletions/ └── deletion-0.lance打开这个目录先看_latest.manifest它是指向当前版本清单的软指针。每次写入或删除数据Lance 不会去修改旧数据文件而是生成新版本清单然后更新_latest.manifest。这个机制保证了读取者永远能拿到一致的快照——这正是多版本控制里“写时复制”思想在文件层面的实现。data/目录下存放的是真正的数据分片Lance 内部叫 fragment。每个 fragment 是一组行组成的不可变文件按列式布局存储列数据内部又分很多个 pagepage 是最小 IO 单位。我建议学员重点理解 fragment 的意义它决定了查询时要扫多少个文件、做多少次随机 IO。理想情况下一次过滤查询只需要基于统计信息跳过大部分 fragment击中少量 fragment 就能完成。indices/目录存放向量索引和标量索引。Lance 的索引文件也是基于列的内部用 page 组织查询走的路径是先查索引页拿到候选行号再回数据文件取行。_deletions/目录是删除标记文件记录哪些行被逻辑删除配合 manifest 为多版本服务——这比直接物理删除数据效率高得多。3.3 manifest 与版本管理的实现机制很多同学第一次接触 Lance 的版本管理时会有疑问这不就是 Git 吗确实原理有相似之处但没有提交历史对象树那么复杂。Lance 的每个 manifest 文件都很轻量内部只是记录了该版本包含哪些 fragment、哪些索引、哪些删除标记以及一些统计信息。因为 manifest 本身很小所以版本切换只需要切换_latest.manifest的指向毫秒级完成。课程里我会手把手带学员做一次版本回溯往数据集里写入三个批次记下每个批次的版本号然后用dataset.checkout(version1)读一遍数据看看它如何只加载第一个版本的 fragment 而完全忽略后面的数据。实操下来你会有种直觉——这就像手上多了一台时光机再也不怕训练数据更新后想回头复现实验结果却找不到原始数据集了。一个特别容易踩的坑是多个进程同时写入同一个 Lance 数据集时manifest 的更新不是完全无锁的。Lance 内置了简单的并发控制但性能在高频写入下会受影响。我建议在管道设计里尽量避免多写入方争抢同一个数据集而是用“主写者 从读者”的架构。3.4 向量索引在文件里是怎么落的Lance 对向量检索的支持是它区别于 Parquet 的最大亮点。它支持两种主流索引算法IVF 和 HNSW。课程里我两种都会讲但重点放在 HNSW 上因为它在召回率和查询延迟之间表现更稳定尤其在数据维度高的时候。写索引的过程大体是把向量列按行取出来构建近邻图然后把图的节点和邻居关系序列化到indices/目录下的索引文件中。每次查询时Lance 会根据查询向量遍历图找出最相近的候选集然后通过 row id 去数据文件里取结果。实际操作中有一个重要参数是num_partitionsIVF或ef_constructionHNSW这两个值越大索引质量越高但构建时间也越长。我的经验是几千万行级别、100 维左右的向量HNSW 的ef_construction设 200查询延迟基本能稳定在个位数毫秒。还有一点要提醒索引写入之后如果继续往数据集里追加数据索引不会自动更新。Lance 会标记索引为“stale”此时查询会退化成暴力扫描或需要重建索引。生产策略通常是定期重建索引或者把增量数据单独建索引再走合并流程。专题课里有一节专门演示这个流程这是很多网上教程不会讲到的。4. 端到端实操用 Ray 编排、Daft 处理、Lance 存储4.1 准备环境与依赖版本开工之前先把环境摸清楚。这个专题课所有实验代码基于 Python 3.10依赖版本我锁在这些组合上经过实测不会打架。但如果你要照抄建议尊重这个矩阵别用最新版轮子莽撞跑。组件版本建议备注Python3.10 - 3.113.12 下部分 Daft 旧版本有兼容问题daft 0.3.0使用set_runner_ray需要这个版本以上ray 2.30用默认版本即可注意和 daft 的兼容表pylance 0.20Lance 的 Python SDK支持读写pyarrow 15.0Lance 底层依赖版本太旧会报错安装命令很简单pip install daft ray pylance pyarrow装完先做一个快速验证import daft; daft.set_runner_ray()不报错说明环境没问题。如果这里报缺包或者找不到动态库十有八九是 pyarrow 版本冲突先把 pyarrow 升级到 15 以上再看。4.2 代码实现全流程整个实验项目模拟一个常见的 AI 数据场景有 100 个 CSV 文件散落在本地目录每个文件包含约 50 万行用户行为日志结构里有用户 ID、时间戳、行为类型、行为分数还有一个 128 维的 embedding 向量。任务分三步清洗特征、把 embedding 写入 Lance、在 Lance 上建向量索引并跑一个相似度查询。第一步用 Ray 把文件列表并行派发出去让每个远程函数调用 Daft 读取一个 CSV。这里的关键技巧是不要让 Daft 一次性读全部文件——即使 Daft 支持自动分区最好还是由上层控制任务粒度因为 Ray Dashboard 上能看到每个 task 的执行时间粒度太粗不利于定位问题。import ray import daft from daft import col, lit daft.set_runner_ray(num_cpus8, num_threads_per_cpu4) ray.remote def process_file(path: str): df daft.read_csv(path) # 基本清洗过滤分数为空的记录时间戳转成日期 df df.filter(col(score).not_null()) df df.with_column(date, col(timestamp).str.strptime(%Y-%m-%d)) # 将 embedding 字符串列解析为 float 数组 df df.with_column(embedding_arr, col(embedding).str.split(,).cast(daft.DataType.list(daft.DataType.float64()))) # 只保留需要的列去掉原始字符串 return df.select([user_id, date, behavior, score, embedding_arr]) paths [fdata/log_{i}.csv for i in range(100)] futures [process_file.remote(p) for p in paths] result_dfs ray.get(futures)ray.get拿到的是一个包含了每个分区 DataFrame 的列表这时候 Daft 还没有真正执行任何计算因为整个链都是惰性的。需要把多个 DataFrame 合并成一个然后再统一触发执行。但这里有一个坑daft.concat在底层需要所有分区的数据结构完全对齐如果不同process_file返回的 schema 有细微差异比如列顺序不同concat 会报错。我的做法是在每个远程函数内部先做一次df df.select(...).collect()转成 Arrow Table确保数据被物化然后外层再用 pyarrow 的combine_chunks合并这是最稳的组合姿势。第二步把合并后的 Arrow Table 写入 Lance。这里要注意pylance的 API 设计和很多数据库 connector 不一样它需要先定义 schema 再写入。拿上面的例子里 embedding_arr 这个 float64 列表列来说Lance 默认会把这种嵌套列表自动识别为固定维数的向量但前提是你的数据里每一行的数组长度都一致。import lance import pyarrow as pa schema pa.schema([ pa.field(user_id, pa.string()), pa.field(date, pa.date32()), pa.field(behavior, pa.string()), pa.field(score, pa.float64()), pa.field(embedding_arr, pa.list_(pa.float64(), 128)), ]) lance.write_dataset(arrow_table, output/user_logs.lance, schemaschema, modeappend)写入完成之后马上验证目录结构你会看到data/下多了几个 fragment 文件_latest.manifest也更新了。这时候用lance.dataset(output/user_logs.lance)重新读取测试一下随机查询性能和版本回溯。第三步是建向量索引。用pylance的create_index函数指定列名和索引类型我这边用的 HNSW 参数是m16, ef_construction200。然后做一个 KNN 查询从一个用户向量出发找最相近的 10 条记录dataset lance.dataset(output/user_logs.lance) dataset.create_index( embedding_arr, index_typeHNSW, metricL2, m16, ef_construction200, ) query_vector [0.1] * 128 result dataset.search(query_vector, vector_column_nameembedding_arr).limit(10).to_arrow() print(result.to_pandas())这个地方是最容易出问题的如果写入的数据里 embedding 列存的是 Python list 而不是 Arrow listcreate_index就会报“unsupported data type”。所以写完数据后第一次建索引前建议先dataset.schema看一眼列类型不放心就重写一遍数据。真实项目中我至少为这个错误排查过半天。4.3 关键参数选型说明上面代码里的参数不是随便拍的背后都有考量Ray 侧num_threads_per_cpu4的设定是为了在 CPU 密集的解析和向量化操作中平衡线程切换开销。设太低比如 1会让单节点多核利用不起来设太高又会因为来回切换线程导致缓存命中率下降。四是一个比较稳妥的中间值。Lance 写入时modeappend而不是overwrite这是课程里特意强调的——如果你反复跑全量写入但不清理旧数据数据集会在后台堆积大量无效 fragment。正确做法是每次批处理开始前做一次 dataset 删除或者直接重建目录。我用 append 模式模拟的是增量日志的持续摄入场景。HNSW 的m参数表示每个图节点最多几根边边越多搜索质量越好但内存也涨得快。千万行以下的数据量m16 完全够用再往大可以试着调到 32。ef_construction决定建图时的搜索宽度我只在训练集上做了几次网格搜索就发现 200 左右已经进入收益递减区间再高只是白白消耗构建时间。5. 实战中的常见问题与避坑清单5.1 Lance 相关的坑第一个高频问题manifest 文件损坏。Lance 的_latest.manifest是文本文件如果不小心被外部程序改坏整个数据集就打不开了。解决办法是去_versions/目录找到最后一个完整可用的 manifest手动拷回来改名为_latest.manifest。我们课堂上专门演练了这个救援流程每个学员都应把这个技能刻进肌肉记忆。第二个坑是旧文件清理。Lance 不主动清理历史文件时间久了磁盘占用会非常夸张。dataset.cleanup_old_versions()可以按保留版本数清理过期数据但要注意正在读取的进程可能持有旧版本引用在线上环境跑清理前务必等读流量低谷期不然会出现读取端报版本不存在的偶发错误。第三个坑是 schema 演化。Lance 支持加列但不是所有的类型变更都能无缝兼容。比如把 float32 列改成 float64旧 fragment 里的数据在读取时需要做类型转换性能锐减。最好在生产环境就约定好 schema 冻结策略有变更先写一个 migration 脚本避免隐式转换。5.2 Daft 相关的坑Daft 虽然 API 像 Pandas但毕竟是一个分布式引擎很多 Pandas 里允许的写法在这里会被优化器拒绝。典型的是 DataFrame 里混入 Python 对象类型Daft 会尝试把它转成 pyarrow 类型转不了就报错。建议所有自定义逻辑统一拆进 UDF明确输入输出类型。另一个常见问题是我前面提到的惰性求值。很多同学写完一个查询链后没有调用collect()以为数据已经处理完了结果后续代码拿到的还是一个空的 logical plan各种奇怪错误。我在专题课里反复强调一个习惯每个 Daft 代码块最后都写一行注释说明触发点。这能避免掉至少一半的 debug 时间。5.3 Ray 相关的坑Ray 的坑集中在反序列化和资源管理上。远程函数里随便传大 DataFrame 是一个经典错误因为 Ray 默认把所有参数转成对象存进分布式内存传一个几百 MB 的 DataFrame 比实际计算还慢。正确方法是把大对象提前存到 Ray 的put(obj)拿到 ObjectRef 后再传给远程函数这样所有 task 共享同一份内存对象高效且省内存。资源控制方面如果 Ray 集群的 CPU 核数有限而你在每个远程函数里又开了 Daft 多线程一不留神就会线程爆炸CPU 超订阅。最简单的方法是让 Ray 每个 task 常驻一个核心Daft 内部只分配给少量线程实验参数按比例放大到整个集群的时候要重新测一下不要四节点跑出一核九用。5.4 版本兼容问题这套技术栈里 Daft 和 pyarrow 的兼容矩阵最磨人。Daft 每次发版都在补强对最新 Arrow 的支持但旧版本的 Daft 遇见新版本 Arrow 可能直接拒绝加载。我的建议是锁定数组库版本项目里放一个requirements-lock.txt明确pyarrow17.0.0这种精确版本。别偷懒用在这套组合里真的会出事。Ray 官方在版本升级时对大版本之间的 API 兼容承诺也不完全到位收到ray.exceptions.RayTaskError的时候先确认是不是版本不匹配。新集群升级 Ray 前最好先在测试环境用同样的 Daft/Lance 代码跑一遍全流程再推进生产。6. 专题课更新计划课程大纲与项目作业改版思路6.1 课程大纲 v2 怎么组织这次专题课更新不是简单地加一个新章节就完事而是重新划分了三个递进阶段第一阶段是工具认知每个工具单独配两到三个小实验确保学员对各自能力边界有感觉。第二阶段是两两协作比如 Daft Lance 解决“清洗完数据直接落盘”或者 Ray Daft 解决“并行处理大批量文件”中间穿插着我们在实战里踩过的坑。第三阶段才是三者合流回到第四节那种端到端管道只是数据量会拉到 5 亿行让学员体会集群资源调度带来的差异。大纲上一共安排了 12 个课时其中 Lance 部分我分配了三节半这也是这次更新的重点。很多教材都轻视存储层其实存储选型决定了整个管道的上限——计算可以靠堆机器解决但存储如果设计不合理多快的引擎也救不回来。6.2 项目作业怎么设计作业的分量不能贪多我们在 v2 版设计了两个核心大作业第一个是让学员基于公开数据集自己设计一个 Ray Daft Lance 管道要求实现增量写入、版本回溯和向量检索三个功能评测标准是延迟和磁盘占用两个指标。第二个是给一个“损坏的 Lance 数据集”让学员通过排查 manifest 和 fragment 文件结构来恢复数据这个作业刻意模拟了线上故障场景做完一遍比背十遍文档都管用。课程最后还会开放一个可选的挑战题把学习到的三件套迁移到一个 GPU 场景里让 Ray 调度 GPU 资源Daft 做数据预处理Lance 存储 embedding并实测 HNSW 索引的查询吞吐量。时间允许的话我会把这个挑战题录成直播边跑边讲每一步的资源观测方式。我个人在编写这门课的时候感触最深的一点是技术选型不是堆砌热门工具而是要理解每个工具背后的设计约束。Daft 为什么用 Rust是因为单节点上 Python 的 GIL 已经卡死了并行性能。Ray 为什么强调分布式对象存储是因为跨节点数据搬运才是分布式系统最大的瓶颈。Lance 为什么做多版本文件结构是因为 AI 实验最需要数据可回溯性。这三件事想通了课程内容自然就串起来了。这份课件我会持续维护下一轮计划把 Ray Data 和 Daft 的 API 差异做一张详细的对照表放进附录共享给所有订阅的学员。
阅读完成 · 觉得有帮助?
咨询建站