news 2026/10/8 3:45:08

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Daft、Ray、Lance三件套:构建AI数据管道的现代方案

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 是一组行组成的不可变文件,按列式布局存储,列数据内部又分很多个 page,page 是最小 IO 单位。我建议学员重点理解 fragment 的意义:它决定了查询时要扫多少个文件、做多少次随机 IO。理想情况下,一次过滤查询只需要基于统计信息跳过大部分 fragment,击中少量 fragment 就能完成。

indices/目录存放向量索引和标量索引。Lance 的索引文件也是基于列的,内部用 page 组织,查询走的路径是先查索引页拿到候选行号,再回数据文件取行。_deletions/目录是删除标记文件,记录哪些行被逻辑删除,配合 manifest 为多版本服务——这比直接物理删除数据效率高得多。

3.3 manifest 与版本管理的实现机制

很多同学第一次接触 Lance 的版本管理时会有疑问:这不就是 Git 吗?确实原理有相似之处,但没有提交历史对象树那么复杂。Lance 的每个 manifest 文件都很轻量,内部只是记录了该版本包含哪些 fragment、哪些索引、哪些删除标记,以及一些统计信息。因为 manifest 本身很小,所以版本切换只需要切换_latest.manifest的指向,毫秒级完成。

课程里我会手把手带学员做一次版本回溯:往数据集里写入三个批次,记下每个批次的版本号,然后用dataset.checkout(version=1)读一遍数据,看看它如何只加载第一个版本的 fragment 而完全忽略后面的数据。实操下来你会有种直觉——这就像手上多了一台时光机,再也不怕训练数据更新后想回头复现实验结果却找不到原始数据集了。

一个特别容易踩的坑是:多个进程同时写入同一个 Lance 数据集时,manifest 的更新不是完全无锁的。Lance 内置了简单的并发控制,但性能在高频写入下会受影响。我建议在管道设计里尽量避免多写入方争抢同一个数据集,而是用“主写者 + 从读者”的架构。

3.4 向量索引在文件里是怎么落的

Lance 对向量检索的支持是它区别于 Parquet 的最大亮点。它支持两种主流索引算法:IVF 和 HNSW。课程里我两种都会讲,但重点放在 HNSW 上,因为它在召回率和查询延迟之间表现更稳定,尤其在数据维度高的时候。

写索引的过程大体是:把向量列按行取出来,构建近邻图,然后把图的节点和邻居关系序列化到indices/目录下的索引文件中。每次查询时,Lance 会根据查询向量遍历图,找出最相近的候选集,然后通过 row id 去数据文件里取结果。实际操作中有一个重要参数是num_partitions(IVF)或ef_construction(HNSW),这两个值越大索引质量越高,但构建时间也越长。我的经验是:几千万行级别、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_cpus=8, num_threads_per_cpu=4) @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 = [f"data/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", schema=schema, mode="append")

写入完成之后马上验证目录结构,你会看到data/下多了几个 fragment 文件,_latest.manifest也更新了。这时候用lance.dataset("output/user_logs.lance")重新读取,测试一下随机查询性能和版本回溯。

第三步是建向量索引。用pylance的create_index函数,指定列名和索引类型,我这边用的 HNSW 参数是m=16, ef_construction=200。然后做一个 KNN 查询,从一个用户向量出发找最相近的 10 条记录:

dataset = lance.dataset("output/user_logs.lance") dataset.create_index( "embedding_arr", index_type="HNSW", metric="L2", m=16, ef_construction=200, ) query_vector = [0.1] * 128 result = dataset.search(query_vector, vector_column_name="embedding_arr").limit(10).to_arrow() print(result.to_pandas())

这个地方是最容易出问题的:如果写入的数据里 embedding 列存的是 Python list 而不是 Arrow list,create_index就会报“unsupported data type”。所以写完数据后第一次建索引前,建议先dataset.schema看一眼列类型,不放心就重写一遍数据。真实项目中我至少为这个错误排查过半天。

4.3 关键参数选型说明

上面代码里的参数不是随便拍的,背后都有考量:

Ray 侧num_threads_per_cpu=4的设定,是为了在 CPU 密集的解析和向量化操作中平衡线程切换开销。设太低(比如 1)会让单节点多核利用不起来;设太高又会因为来回切换线程导致缓存命中率下降。四是一个比较稳妥的中间值。

Lance 写入时mode="append"而不是overwrite,这是课程里特意强调的——如果你反复跑全量写入但不清理旧数据,数据集会在后台堆积大量无效 fragment。正确做法是每次批处理开始前做一次 dataset 删除或者直接重建目录。我用 append 模式模拟的是增量日志的持续摄入场景。

HNSW 的m参数表示每个图节点最多几根边,边越多搜索质量越好但内存也涨得快。千万行以下的数据量,m=16 完全够用,再往大可以试着调到 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,明确pyarrow==17.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 差异做一张详细的对照表,放进附录共享给所有订阅的学员。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/8 3:45:05

华为ICT大赛网络赛道国赛ESNP实验全流程拆解与避坑指南

简介:这份资源是华为ICT大赛2022-2023网络赛道国赛实验ESNP的真题与完整详细答案合集,面向备战华为ICT大赛的高校学生、网络技术学习者及指导教师,帮助解决国赛实验环境搭建难、配置思路不清晰、缺少权威参考答案等问题。压缩包共34个文件&am…

作者头像 李华
网站建设 2026/10/8 3:44:47

Android ScrollView 与 HorizontalScrollView 的滚动与嵌套避坑指南

今天想认真聊聊 Android 里两个被用烂却很容易踩坑的滚动容器:ScrollView和HorizontalScrollView。它们一个管纵向滚动,一个管横向滚动,几乎所有“内容比屏幕大”的页面都会遇到,但很多开发者在真正上手时会碰上测量异常、事件冲突…

作者头像 李华
网站建设 2026/10/8 3:44:19

AI对话历史管理:用Markdown和日历视图打造本地记忆系统

1. 为什么AI对话历史记录让人抓狂用AI写代码、查资料、做方案的人应该都有同感:对话列表越拉越长,想找三天前让AI帮忙改的那段正则表达式,得在侧边栏里翻半天。翻到了还好,翻不到就只能重新问一遍,而重新问的结果往往和…

作者头像 李华
网站建设 2026/10/8 3:43:30

DSec:面向DeepSeek智能体训练的弹性计算沙箱基础设施

大型语言模型的智能体训练,往往不是模型训练本身有多难,而是“一批智能体同时跑起来”这件事,能把一个普通开发机折腾到怀疑人生。我在把多个DeepSeek驱动的智能体放出去做环境交互、工具调用和多轮自我博弈时,第一波遇到的就是资…

作者头像 李华
网站建设 2026/10/8 3:42:31

pstack-claude:本地可信AI编程助手的进程级实现原理

1. 项目概述:pstack-claude 是什么,它解决的是哪类真实开发痛点?pstack-claude 这个名字乍看像一个工具组合词,但拆开来看,“pstack”是 Linux 系统中一个真实存在的诊断命令,用于打印指定进程的调用栈&…

作者头像 李华