
1. 为什么是这三件套大模型时代被重新定义的数据基础设施先聊聊这次更新课程的动机。去年我在线下分享会里做过一次小调查问现场做机器学习平台的同学你们的数据处理链路用什么一半人回答还是Pandas加SQL另一半人已经开始聊Lance和Ray了。这个变化很有代表性——大模型带来的多模态数据、向量检索、大规模微调正在把过去那套CSV灌进DataFrame、算完特征再存回数据库的流程彻底冲垮。Daft、Ray、Lance这三个名字放在一起其实分别回答了三层问题Daft是算得更快的答案一个原生分布式设计的DataFrame引擎Rust核心让数据分析能直接吃下远超单机内存的数据集Ray是调得更顺的答案一套把Python应用从单进程扩展到上千节点的分布式运行时任务调度、状态管理、弹性伸缩都替你处理好了Lance是存得更聪明的答案专为ML和检索场景设计的列式存储格式既能存原始数据又能存向量索引文件组织方式天生支持随机访问和增量更新。这门专题课要更新的核心就是把这套存储-计算-调度三层结构讲透。课程面向三类人第一类是数据工程师你们已经在生产环境里维护Spark或Flink想看看新一代工具在机器学习场景下能省多少事第二类是算法工程师你们训练模型之前要花大量时间做特征工程和数据预处理受够了Pandas的OOM和慢查询第三类是平台开发你们想做一套面向AI应用的数据基础设施需要理解文件格式和分布式调度这些底层机制。课程本身不做纯理论推演——每节都配了可运行的代码案例和数据集。这次更新的一个重要原因就是过去几个月里Daft、Ray、Lance都有了大版本迭代很多API用法变了网上能找到的资料大多过时了我们干脆把所有代码重新跑了一遍、踩过的坑全部写进课程里。说句实话这三个工具单独用起来都不难难的是把它们组装成一套生产可用的管线。这个组装过程才是专题课真正的核心价值。2. Daft一个从设计之初就为分布式而生的DataFrame引擎先讲Daft因为它最颠覆直觉。很多人第一次接触Daft会问它和Polars有什么区别和Dask又有什么区别市面上的分布式DataFrame不少Daft的特别之处在于它不是先做好了单机引擎再想办法分布式化的修补方案而是从执行引擎的底层设计上就是分布式的。2.1 为什么Pandas在大数据场景下注定力不从心理解这一点要先搞清楚Pandas慢在哪。Pandas的DataFrame把数据全加载到内存然后靠Python对象逐行处理。单机内存上限决定了数据集规模上限Python解释器的逐行循环又决定了处理速度上限。我在课程里专门做了一个对比实验同样是对10亿行日志做groupby-countPandas直接内存溢出Daft在8核机器上跑出了接近150MB/s的扫描吞吐。Daft的架构大致分三层最底层是Rust实现的原生列式内存格式利用SIMD指令做向量化计算这一步和Polars的思路一致中间层是执行引擎把用户写的查询转换为分布式执行计划支持分区、并行、流水线优化最上层才是用户看到的DataFrame API和表达式系统。关键是中间这层Daft把一个DataFrame按行做水平切分分成若干个Partition每个Partition里的计算任务可以分发到不同机器或不同CPU核心上并行执行。数据量不够大时它就在单机上用多线程跑数据量大了可以无缝切换到Ray集群或者本地多机环境。这种设计让用户不需要像学Spark那样预先理解一堆集群概念——同一套API小数据量时当Polars用大数据量时当Spark用。2.2 Daft表达式API的实操笔记能写出更优的执行计划Daft的API设计明显参考了Polars的表达式风格但用起来又更接近Pandas的习惯。从一个实际例子看import daft # 读取本地数据集schema自动推断 df daft.read_csv(s3://bucket/events/*.csv) # 构建表达式筛选 派生新列 聚合 result ( df.where(df[event_type] click) .with_column(hour, df[timestamp].dt.hour()) .groupby(hour) .agg(daft.col(user_id).count().alias(click_count)) )这段代码看起来平平无奇真正有价值的是Daft拿到这个查询之后做的事情。它会先分析整棵表达式树做谓词下推——把event_type click这个过滤条件下推到数据读取阶段实际读进内存的行数大幅减少再做列裁剪——只加载查询涉及到的列跳过其余列这在列式文件上效果尤其显著。如果查询涉及了多个文件或分区Daft会自动规划任务分发顺序尽量让数据本地性更优的任务先执行。数据量不大时Daft在体验上和Pandas持平甚至更快数据量到了几百GB或者几个TBDaft依然不需要你改写代码逻辑。能够在异构环境之间无缝伸缩这是我认为它比重新学一遍Spark更值得投入的原因。2.3 Daft在专题课里的定位非结构化与多模态数据的统一入口大模型应用带来了大量非结构化数据——图片、PDF、音视频、JSON日志。Pandas处理这些数据时通常要先写一堆解析函数然后批处理成结构化表格Daft则把多模态数据的一等公民支持做进了引擎本身。比如读取一个图片目录Daft可以直接生成带image类型列的DataFramedf daft.read_images(s3://bucket/product-images/, limit1000)读进来之后可以配合表达式做缩放、裁剪、特征提取还能直接调用Python UDF把CLIP或者GPT-4V的输出作为新列挂到原表上。这种原始文件直接变成可分析DataFrame的能力在处理多模态训练数据时特别实用。我平时给客户做数据管线很多时间花在写各种数据加载和清洗脚本上用Daft之后这类脚本的代码量至少砍掉一半。Daft在课程中承担的定位是数据入口和计算引擎我们从各种存储源读入数据在Daft里完成清洗、变换、特征工程再把处理好的数据写出去。3. Ray从多进程到分布式运行时的思维升级如果说Daft解决的是数据怎么算得快Ray解决的就是计算任务怎么编排得稳。Ray这两年几乎成了Python分布式系统的事实标准从强化学习RLlib到模型服务Serve从数据预处理Data到超参搜索Tune生态铺得很广。3.1 Ray Core的核心模型一切皆Task与Actor很多从Python多进程转过来的同学一开始会不太适应Ray的写法。Ray为我们提供了一个按函数粒度拆分任务的模型——把计算逻辑定义成远程函数然后用ray.remote装饰器标记通过.remote()方法提交这些任务会被自动调度到集群的任意可用节点上执行。import ray ray.remote def process_batch(batch): # 模拟数据处理逻辑 cleaned clean(batch) return feature_extract(cleaned) # 提交100个任务Ray自动处理调度 futures [process_batch.remote(data[i]) for i in range(100)] results ray.get(futures)这段代码背后发生了什么Ray把所有远程函数和对象都登记到全局分布式调度器调度器根据各节点的资源余量CPU、GPU、内存、数据本地性、队列状态将任务派发到最优节点对象存储是一个分布式内存存储层任务间通过引用传递数据。理解了Task之后再理解Actor就简单了——Actor就是带状态的远程对象可以持续在某个节点上驻留供多个任务读写。状态管理比如模型参数、连接池、会话状态的本质需求在分布式系统里经常出现用Actor能优雅地解决。3.2 Ray Data为AI场景设计的数据管道层单独用Ray Core做数据任务往往过于底层课程里更多会用到Ray Data这是一套建立在Ray分布式运行时之上的数据集API处理任务、文件读取、数据混洗等细节都替你封装好了。从不同来源读入数据做变换再输出是Ray Data的主打场景对比Spark DataFrameRay Data在同一套Python生态里工作没有JVM和Python之间的序列化鸿沟和PyTorch/TensorFlow的衔接也更顺滑。Ray Data内置了数据集分块机制处理大数据时自动把数据按大小切分还支持流式执行前一个算子还没跑完就拉起下一个算子管道式负载比逐阶段等待高得多。3.3 分布式调试与资源调度的踩坑记录这里多聊一点专题课更新过程中反复踩坑的经验也是我认为Ray最难掌握的地方。第一个坑是对象存储的内存溢出。Ray的ray.get会把分布式对象拉到本地内存如果任务产出的对象太大本地内存会被打爆。我们在课程实验里遇到过一次处理一个2GB的DataFrame直接ray.get整个对象运行到一半节点OOM。正确做法是使用ray.data或把任务返回值拆成细粒度对象分批获取。第二个坑是资源参数的设置。很多人写ray.remote不指定资源需求结果大量任务被挤到一个节点上。我们需要用num_cpus和num_gpus显式声明每个任务需要多少资源调度器才能合理分发。ray.remote(num_cpus4, num_gpus1) def train_model(config): ...第三个坑是死锁问题。Actor里如果还在等待另一个Actor的结果而对方卡在等待当前Actor的返回值常见的环形依赖局面就出现了。Ray不会自动检测这种循环等待任务会一直挂到超时。这些坑在官方文档里很少被提到都是生产环境里踩出来的。专题课专门加了一节Ray排错手册把这几个月遇到的问题按症状、原因、解法整理成表。4. 深挖Lance的文件组织结构理解列式存储的最新范式最近在社区里被问得最多的技术点就是Lance的文件组织结构。这也是专题课更新时花了最多时间研究的部分。Lance的数据格式设计理念一句话概括既要Parquet列式存储的高压缩率和查询性能又要支持随机访问、增量写入和向量索引。这决定它的文件组织和传统Parquet文件很不一样。4.1 顶层目录一坐下来就得看懂的五个部分一个标准的Lance数据集目录大概是这样的my_dataset.lance/ ├── data/ │ ├── 000000_000100.fragment │ ├── 000001_000200.fragment │ └── ... ├── index/ │ ├── 000000_ivf_pq.index │ └── ... ├── _latest.manifest ├── _transactions/ │ ├── 00000000000000000001.operation │ ├── 00000000000000000002.operation │ └── ... └── version/ ├── 00000000000000000001.manifest ├── 00000000000000000002.manifest └── ...逐层解释。data/目录存放实际数据文件。Lance不把整个表写进一个大文件而是拆成多个fragment片段每个fragment是若干行数据集合独立存储。fragment命名格式是{fragment_id}_{version}.fragment同一个fragment在后续版本被修改时会生成新版本文件旧文件保留这是实现时间旅行和增量写入的基础。index/目录存放向量索引文件。Lance原生支持IVF_PQ、HNSW这类ANN索引索引文件同样按版本管理。这解决了旧方案里向量存一个系统、原始数据放另一个系统的割裂问题。_latest.manifest一个指向当前最新版的manifest文件的软链式入口。每次commit成功后就会更新这个文件查询时先读它定位数据集当前状态。version/目录存放各版本manifest文件。每个版本号是一个20位十进制数字例如00000000000000000001递增规则和其它系统一致。每个manifest文件记录了该版本数据集包含哪些fragment、schema是什么、索引状态如何。_transactions/目录新版本引入的事务日志。写入操作不是直接改数据文件而是先记一条操作比如Append、Overwrite、Deletecommit时再生成对应manifest。这样任何一次写入都能被回溯和审计。4.2 fragment与column groups数据组织的两个关键概念Fragment可能是Lance文件组织里最核心的概念。一个fragment相当于Parquet的row group是逻辑上连续的行子集。但Lance在fragment内部又引入了column groups的概念这是它区别于Parquet的一个巧妙设计。默认情况下表的所有列存放在同一个column group里如果某些列经常被一起查询可以分成多个column groups每个group内的列数据连续存储。这大大提升了列裁剪的效率——查询只需要把用到的column groups加载进内存完全用不到的那些列连碰都不用碰。举个例子一张用户画像表的schema包含基础属性列user_id、age、country行为序列列recent_clicks、recent_views向量列embedding。如果你经常跑的是向量检索场景——按embedding做近邻搜索再回查基础属性——可以把向量列单独放在一个column group基础属性放在另一个group。搜索时只加载向量group回查时只加载属性group整个过程的I/O开销比把整张表读进来少一个数量级。fragment级别还支持细粒度的行级过滤优化。Lance的manifest里记录了每个fragment的一些统计信息比如某列的min/max值查询时执行器可以先跳过明显不满足过滤条件的fragment只扫描可能命中的数据段。4.3 写入、Commit与快照Lance的版本管理机制Lance的写入机制设计得相当讲究。大体流程是这样的第一步写入数据以append方式创建新的fragment文件新文件写入data/目录第二步在_transactions/目录下记录本次操作的类型和涉及的文件第三步生成新的manifest文件内容包含最新fragment列表、列组信息、schema和索引元数据第四步把_latest.manifest原子地更新指向新manifest。这四个步骤确保了并发安全。传统方案在写入时锁表或用全局状态Lance用原子文件替换配合版本号递增多个写入者可以并发操作而不会互相覆盖。写完一个版本后旧版本的数据文件并不会被物理删除而是作为历史版本保留——这就是Lance的时间旅行能力。查询时指定versionN系统就加载对应版本的manifest按当时的fragment列表扫描数据默认查询加载_latest.manifest指向的最新版本。4.4 从文件组织反推设计理念为什么Lance能又快又省理解Lance文件组织之后回看整个设计能看出几条明确的设计原则。第一条是不可变文件版本指针。数据文件一旦写入就不再变更修改通过新增fragment和更新manifest完成。这让并发读写变得安全也让备份和恢复变得简单——只需要快照manifest和对应数据文件即可。第二条是索引与数据同构。向量索引和原始数据放在同一个文件系统结构里由一个引擎统一管理。数据更新时索引同步更新避免了数据在MySQL、向量在Milvus、日志在ES这种多系统维护的噩梦。第三条是按访问模式组织。fragment和column groups都是为了减少I/O尽量让查询只触碰必要的数据。这套指导思想在数据量达到TB级之后尤其重要。我实际测试过一个200GB的图片特征数据集Parquet格式做随机读取需要遍历整列文件才能定位目标行Lance按fragment和索引定位后在相同硬件上的响应时间低了接近一个量级。专题课里专门安排了两节课讲Lance一节讲列式存储基础一节专门拆解文件组织。后一节的作业是让学生动手读_latest.manifest里的二进制信息自己画出某个版本对应的fragment树。5. 三件套串起来手工搭建一条端到端管线理论讲完得来点实操。专题课更新到一半时我用Daft、Ray、Lance三个工具重新搭建了一遍课程配套的演示场景从一批公开的商品评论数据50GB中提取文本向量建立检索索引并做基础的分布式统计分析。这条管线完整覆盖了三个工具各自的主要能力。5.1 管线的整体流程设计管线分四段读取用Daft读取原始评论数据清洗文本做分片向量化用Ray分布式调用SentenceTransformer模型批量生成embedding写入索引把原文和embedding写入Lance数据集并创建ANN索引检索分析用Daft查询Lance做统计分析用向量检索做语义搜索Demo。5.2 关键环节的代码拆解向量化这一步是重点。50GB的数据按分片提交到Ray集群每个分片调用一次embedding模型生成向量产物是不需要全部拉回本地的——直接用Ray对象引用传给下一阶段。import ray from sentence_transformers import SentenceTransformer ray.remote(num_gpus1) class Embedder: def __init__(self, model_nameall-MiniLM-L6-v2): self.model SentenceTransformer(model_name) def embed(self, texts): return self.model.encode(texts, batch_size64).tolist() # 每个数据分片对应一个Embedder Actor的embed调用为什么要用Actor而不是普通Task因为模型加载是有状态的资源——加载一次模型要几百MB显存如果每个任务都重新加载一次GPU会被浪费。Actor常驻节点模型只加载一次后续任务复用。这是Ray Actor在AI管线里的典型用法。写入Lance的代码也很直接import lance import pyarrow as pa def write_to_lance(table: pa.Table, uri: str): lance.write_dataset(table, uri, modeoverwrite) dataset lance.dataset(uri) dataset.create_index( embedding, index_typeIVF_PQ, num_partitions16, num_sub_vectors48, )这里有个容易踩的坑第一次写入用modeoverwrite后续增量更新必须用modeappend。但需要注意append模式下如果schema有变化会失败必须保证新数据和已有数据的schema完全一致。我们更新课程时就遇到过一批embedding维度不一致的数据写入直接报错——排查了半天才发现是数据源里混进了一批旧版本的模型输出。5.3 管线效果和性能数据整条管线在4台8核GPU机器上跑完用了约37分钟其中向量化占了大头写入Lance并构建IVF_PQ索引大约花了6分钟。对比之前用Pandas加faiss的方案实测同数据集流程从约2个小时加各种手动内存管理降到不到40分钟全自动完成。更关键的是代码量。旧方案里MPI式的分片逻辑、内存释放逻辑、文件中间态清理逻辑加起来约600行新方案中Daft代管数据读取和分片Ray代管任务调度Lance代管存储和索引核心代码不到200行。这就是选对工具的价值。5.4 这套管线在专题课里怎么布置成作业课程里把这条管线拆成三个递进式作业方便学员逐步掌握作业一用Daft单独完成数据读取和初步过滤导出统计报告作业二在Ray上把作业一的单机流程改成多节点并行对比加速比作业三把处理结果写入Lance并构建索引完成一个语义检索Demo。作业的评分标准不只是跑通还要求画出每个环节的资源使用曲线解释为什么没有线性加速——Ray固有调度开销、数据倾斜、I/O瓶颈这些才是生产环境真正会遇到的问题。6. 专题课的大纲规划与学习环境建议最后聊聊课程本身的结构和配套资源。6.1 课程大纲从项目实战反推知识结构专题课一共规划了九章去掉前后介绍和总结实际技术内容七章第一章为什么新一代数据工具都在重做存储和计算第二章Daft基础——表达式、DataFrame、读多模态数据第三章Daft进阶——分布式执行计划、性能调优第四章Ray Core——远程函数、Actor、资源管理第五章Ray Data——数据集操作与流式处理第六章Lance列式存储格式——schema、fragment、文件组织第七章端到端项目——商品多模态检索系统搭建。每一章的作业都尽量贴近真实业务场景。Daft部分用的是电商日志分析Ray部分用的是批量模型推理服务Lance部分用的是向量检索系统。学员学完能直接把这套技能迁移到自己的业务里。6.2 学习环境本地开发与集群模拟的取舍我在课程里反复强调一个观点不要为了学分布式一上来就搭几十台机器的集群。Daft本地多线程模式已经完全够用Ray也支持在单机多核上模拟分布式调度。先在本地跑通逻辑再按需上云这个路径最高效平滑。推荐的学习环境分三档学习阶段环境配置目的入门8核CPU 16GB内存Daft单机模式掌握API和表达式进阶单机多核跑Ray本地集群理解分布式调度与Actor模型实战2台以上实体机或云上多实例验证多节点性能扩展代码统一用Python 3.10依赖管理用uv环境复现比conda稳定得多。6.3 给学员的三个实践建议第一个建议动手之前先读源码注释。这三个工具都有很好的文档和类型标注遇到不明白的API建议先跳进源代码看看实现逻辑。Daft的表达式系统源代码写得特别清晰读一遍胜过看三遍文档。第二个建议刻意练习读文件组织的能力。课程里花了整整一节课讲如何在Lance数据集的version/目录里找出当前版本指向的fragment再用代码验证查询是否真正跳过了无关文件。这是真正理解存储引擎性能特点的捷径。第三个建议做性能实验时一定要记录资源使用率。光看时间缩短不够CPU和内存的利用率曲线能暴露很多问题——比如数据倾斜、序列化开销、I/O等待。课程配套模板里提供了一个用Ray的Dashboard加psutil做的监控脚本学员可以改造成自己的性能观测工具。我在这轮更新课程的过程中最深的一个体会是技术的更新速度远超过文档的更新速度。Daft的分区策略、Ray的调度算法、Lance的文件格式都在持续演进。能做的就是把原理讲透把排查方法教给学员让他们在版本迭代时也能自己跟上。专题课的价值不在于背熟某个版本的API而在于理解这套体系解决问题的思路然后在自己的场景里举一反三。最后分享一个小tips学习这三件套时从Lance入手可能比从Daft或Ray入手更容易建立全局观。因为存储格式决定了数据如何读写理解了数据物理布局再回头看计算引擎的设计选择——为什么做列裁剪、为什么做谓词下推、为什么做分区——就全都说得通了。这也是我把章节顺序刻意安排在Daft → Ray → Lance → 项目实战的原因拿Lance当枢纽知识串成网才算真正吸收。