先说一个我自己的经历。去年手头有个日志清洗管道,每天几十GB数据要先做过滤、join维度表、再groupby按小时聚合,最后还要转成Tensor喂模型。第一版全用Pandas,跑一趟大概40分钟,数据量翻倍后直接卡死。后来换成cuDF,同一条链路跑下来不到4分钟,代码改动比想象中小得多。但正因为经历过“从能跑变成跑得更快”的过程,我才意识到:cuDF不是把import pandas改成import cudf那么简单,它的架构组织、内存模型、并行调度方式,和Pandas完全是两套思维。
这篇文章我想从源码层面做一个全景式拆解,围绕cuDF的架构分层、GPU数据加速原理和工程落地方法展开。适合刚开始碰RAPIDS的数据工程、准备优化数据管道的算法工程师,以及正在做数据处理中间件选型的架构师。文中不会停在“怎么调用API”这种层面,而是尽量讲清楚每条数据在GPU世界里经历了什么,以及为什么某些坑只有看源码才能想明白。
1. 从Pandas到cuDF的认知升级:GPU数据框到底改变了什么
1.1 数据框任务的本质是内存带宽问题
很多人把Pandas慢归结于“Python本身慢”,其实不全是。df[df['a'] > 0]这种过滤,Python层生成的只是整棵执行树的顶点,真正花时间的是对一列数据逐个元素做比较,再把结果汇聚成新的列。这类操作是典型的数据密集型任务,吃的是内存带宽与访存效率。
DDR4时代CPU内存带宽大约几十GB/s,而NVIDIA数据中心显卡上的HBM2e/HBM3显存带宽可以到1.5~3.2TB/s。CUDA GPU有几千个计算核心,相比CPU十几个核心,天然适合“每条线程处理一行”的SIMT模式。但这里有个容易误解的点:GPU快,不只是因为核心多,更因为用它执行数据并行任务时,可以同时大量访问不同内存区域,带宽利用率远高于CPU。cuDF就是把这个底层优势包装成DataFrame接口,让JOIN、GROUPBY、SORT这些“看起来笨重”的操作在GPU上重新被并行化。
1.2 cuDF的能力边界:不是“Pandas on GPU”这么简单
初看cuDF,会觉得它就是把Pandas搬到了GPU上,可以加列、过滤、合并、聚合,还能resample。但真正深入用会发现,两者兼容不等于替代。
第一,数据模型有差异。Pandas里有object类型、Categorical、时区感知datetime,cuDF虽然提供了对应的cudf.DataFrame和cudf.Series,但对于object列、嵌套结构列的处理很有限。String列在cuDF内部是用offset数组加chars数组实现的,算是支持,但复杂的字符串正则、逐字符串逻辑仍然比CPU场景更容易踩到性能坑。
第二,内存模型变化了。Pandas运行时数据主要在主机内存,cuDF数据主要在显存。显存比主机内存小得多,所以“把所有数据放进去再计算”这个前提本身就需要重新审视。早期版本显存不足会直接OOM,现在cuDF有了spilling和统一内存(UVM)作为兜底,但这不意味着可以无限塞数据——它只是把“爆显存”变成了“慢速换页”。
第三,API兼容度不是100%。大多数常用API都能用,但部分函数和参数行为不同。比如apply、transform这类需要调用Python回调的接口,在cuDF里性能极差,因为它需要把每一块数据从GPU拉回CPU执行函数,再传回GPU,来回拷贝的开销足以把加速红利全部吃掉。工程上要习惯把逻辑改成向量化写法。
1.3 源码视角的整体调用链
如果只看Python层,cuDF和Pandas长得像,但顺着调用链往下挖,会发现它是一个非常深的栈。以读CSV为例,实际路径大致是:
cudf.read_csv (Python API) -> cudf.io.csv (Cython 封装) -> libcudf::io::read_csv (C++ 核心) -> CUDA kernel:CSV tokenizer / converter这个链路意味着每一层都有可能成为瓶颈。Python层负责参数校验和对象转换,Cython层负责把Python对象翻译成C++的column_view,libcudf层负责任务拆分和kernel调度,而真正的数据搬运与计算发生在CUDA kernel里。你写的每一行cudf.to_pandas()或pandas_to_cudf(),本质都是在打破这个链路,把数据从显存搬回主机内存或反向搬回。
看源码后最直观的收获是:不要在三行代码里反复横跳。数据在GPU上,就应该让它尽量全程留在GPU上,直到真正需要传给另一个框架时再做一次跨设备传输。
2. cuDF分层架构:Python API、cython胶水层与libcudf引擎的真实分工
2.1 Python层的数据结构与延迟行为
cuDF的Python代码主要在python/cudf/cudf目录下,cudf.DataFrame、cudf.Series、cudf.Index是日常接触最多的三个对象。但注意,这些类并不是继承自Pandas,而是抽象出了一套自己的列式存储协议。
在源码里可以看到,cudf.DataFrame内部维护了一个ColumnAccessor对象,它管理着多个Column,每个Column对应一列数据。ColumnAccessor做的事情不仅仅是存字段,还包括列标签映射、切片索引、类型一致性检查等。这些Python层的额外逻辑会带来解释器开销,所以在“每行执行一次”的循环里操作cuDF效率远不如直接整列操作。
Python层还有不少防御性代码:类型检查、索引对齐、元数据拷贝。比如你手写一个循环,对每个分组调用一次groupby[col].mean(),每次调用的固定成本可能在几毫秒到几十毫秒不等,几千个分组叠加起来就非常难看。源码里和apply相关的路径尤其明显,它会通过_apply这类机制构造一个复杂的执行包装,最终绕不开CPU-GPU数据传输。这也是很多从Pandas迁移过来的人觉得“为什么GPU反而更慢”的第一个原因——用Python循环碎片化调用了GPU。
2.2 Cython层的作用与核心接口
再往下是Cython胶水层,代码在python/cudf/cudf/_lib目录下。为什么需要这一层?因为libcudf本身是纯C++库,Python解释器无法直接调用C++对象,必须通过Cython生成的CPython扩展模块来做桥接。
这一层做的事情很机械但很关键。它把Python的Column对象转换成C++的column_wrapper或column_view,把Python的DataFrame对象改造成C++的table对象;调用完后,再把C++返回的table重新拼装成Python DataFrame。异常处理也在这里做:libcudf抛出的std::logic_error、std::runtime_error会被捕获,翻译成Python异常后重新抛出。
很多人容易忽略Cython层的另一个贡献:GIL释放。libcudf核心计算不会长时间持有Python全局解释器锁,所以当CUDA kernel在GPU上跑的时候,Python解释器理论上还能去执行别的线程释放。如果直接用pybind11做绑定,同样也能释放GIL,但RAPIDS早期选Cython有历史原因,也方便直接复用很多DBMS风格的数据结构定义。我们写Python扩展时不妨记住:粘合层不只是转换指针,还需要考虑并发语义和异常边界,否则一个小错误会一直拖垮上层性能。
2.3 libcudf C++层的列式存储与算法注册
libcudf是cuDF真正的计算内核,所有面向用户的计算最终都被翻译成libcudf中的函数调用。它的核心数据结构是cudf::column和cudf::table,在cpp/src/column/column.cpp和cpp/include/cudf/table/table.hpp中能读到完整定义。
cudf::column最少包含几个部分:类型信息(data_type)、数据缓冲区(rmm::device_buffer),以及一个可选的null mask缓冲区。为了表示String等嵌套类型,它还会持有子列(children)。我们看一个简化版本:
struct column { data_type type; // 列类型,如 int32、float64、string std::unique_ptr<rmm::device_buffer> data; // 列数据主体(设备内存) std::unique_ptr<rmm::device_buffer> mask; // null mask std::vector<std::unique_ptr<column>> children; // 子列,如string的offset列 size_type size; // 行数 size_type offset; // 视图偏移 };offset字段的存在很有意思,意味着我们可以用O(1)时间构建一个子表视图,而不需要复制底层数据。这种视图思想贯穿整个libcudf:column_view只是一组指针和长度描述,真正的数据落在一个共享的buffer里。所以读源码时会发现很多操作直接接受column_view,而返回新的column,中间并不总是发生显式拷贝。
在libcudf中,算法层实现了众多数据操作,包括copying、sorting、joining、groupby、strings、io等。从这里开始,性能分析不能再套用Pandas的“逐行执行”模型,而是要理解每个函数底层启动的CUDA kernel,以及每个kernel对显存访问的pattern。
2.4 CUDA kernels与stream管理:底层加速的真正来源
真正让cuDF跑快的,是GPU上那些并行内核。以哈希连接为例,libcudf内部会有一个build表过程,把要连接的key写入GPU全局内存,再启动hash_probe类内核,每个线程处理一个或多个probe行的key计算与匹配。排序操作则会调用GPU上的并行归并排序内核,可能在O((N/P) log N)级别的并行复杂度上完成。
这些内核并不是各自孤立执行的,它们需要被放到CUDA stream中。libcudf许多接口允许传入可选stream参数,默认使用默认流。生产环境里如果既有cuDF计算又有深度学习推理,最好显式管理stream,避免kernel之间不必要地串行等待。RMM内存分配器在这一层扮演了重要角色:cudaMalloc一次大约有几十微秒到毫秒级开销,频繁调用会明显拖慢小任务,而RMM会维护内存池,把分配过的块的释放和复用留在池内,显著降低分配代价。
这也是为什么源码评测里不能只看算法复杂度。一个设计合理的内核,如果每次都在默认流里等待上一个内核结束,性能会折损;反之,一个复杂度稍微高一点的内核,如果用对了stream和pool,可能整体反而更快。工程上,我们应该学会用nsys profile这类工具看kernel时间线,而不是凭“复杂度低”来判断快慢。
3. GPU数据加速原理:从列式内存布局到内核启动的完整链路
3.1 为什么列式格式是GPU加速的基础
数据框在GPU里几乎都是列式存储的,这不是偶然。行式存储时,一条记录的多个字段连续存放,但如果GPU线程要处理“某一列的所有值”,访问的地址会跨很大范围,每次读取都可能只命中一个字段,搬运了大量无用数据。
列式存储则相反,一列中的值在显存里是连续排布的。当32个线程一起读取一列里相邻的32个int32时,GPU可以把这128字节合并成一个或少数几个内存事务,极大利用内存带宽。这个“合并访问”是GPU性能的铁律。libcudf里几乎所有算子都围绕“对列进行操作”这一假设设计,连字符串类型也都采用“offsets列 + chars列”的列式布局,而不是每行一个独立缓冲区。
如果数据原本是行式CSV,读入时就需要做转置工作。读CSV的内核本身就会把每行的多个字段拆分到不同列缓冲区里,这个过程也会消耗时间和显存,但比起后续计算节省的带宽,通常还是值得的。你可以在源码里看到io/csv目录下的tokenizer kernel,核心任务之一就是做这种行列转换。
3.2 内存资源的分层与UVM/spilling机制
做CUDA开发的人都清楚,显存不是唯一的内存层级。每个流多处理器有共享内存和寄存器,访问速度远高于全局显存,但容量很小。libcudf算法设计时会在共享内存里放哈希表桶的局部缓存、排序的共享块等,目的就是减少global memory的访问。
对于用户层,最敏感的是“数据放得下吗”这个问题。cuDF支持两种扩展内存的方式。一种是UVM(managed memory),分配由CUDA统一管理,GPU访问缺页时会把数据从主机内存搬到显存。另一种是显存不足时的spilling,把暂时用不到的device buffer转存到host memory,用的时候再搬回来。开启spilling后,逻辑显存容量可以超过物理显存,但换页开销非常大,可能把原本的几十倍加速拉回个位数。
工程上建议先估算峰值内存。比如一个1亿行的DataFrame,每列int64约8字节,10列就是8GB;如果后面做hash join,还得额外加上哈希表的内存开销。用脚本模拟一遍数据规模变化,比硬着头皮算要靠谱得多。大多数时候,我会先把CUDF_SPILL_DEVICE_MEM_TRY=ON打开,再配合nvidia-smi监控实际显存占用,而不是让程序一直裸奔在大显显存的边缘。
3.3 CUDA stream与并发调度
当一次数据处理有多个算子时,比如“过滤列a,再排序列b”,朴素做法是让这两个kernel顺序执行。默认流模式下,前一个kernel结束后后一个kernel才开始。如果两条计算路径之间没有数据依赖,我们就可以用多个stream让它们并发执行,进一步压满GPU。
libcudf部分API已经能接受stream参数,但用户直接操作stream的机会不算多。实际生产里,特别是想要和PyTorch训练流水线交错执行时,可以把cuDF处理放在一个stream,把深度学习前向推理放在另一个stream,使用事件或cudaStreamSynchronize做同步点。这里要注意:CUDA默认流对没有显式指定流的操作会施加隐式同步,可能让并发失效。正确做法是对每个流都显式调用库接口,或者确保所有操作都使用同一个非默认流池。
源码中还可以看到RMM的执行策略(exec policy),它把分配、释放与stream绑定。这不仅是性能优化,也是正确性要求:如果分配器使用的stream和kernel使用的stream不一致,数据可能在你以为已经完成时还没真正到达目标缓冲区。踩过这个坑之后,我现在几乎每个生产脚本都会统一管理stream句柄,不依赖默认流。
3.4 算子融合与表达式模板
GPU数据加速还有一个常用手段:减少中间结果写入显存的次数。假设你要执行df['new'] = df['a'] * 2 + df['b'],若拆成两步做:先算a*2生成临时列,再把临时列与b相加生成结果,那就写了两次全局内存。更聪明的做法是把表达式树传给计算引擎,由引擎把整个表达式展开成单一kernel,每个线程直接读取a、b原始列,计算出最终结果后只写一次。
libcudf在表达式处理上走了模板实例化路线:核心表达式类在cpp/include/cudf/expressions/中,支持列引用、字面量、二元运算、一元运算等节点。编译期会把这些节点组合成具体的内核模板实例,而不是运行时JIT。这样做的好处是没有启动解释器解析AST的开销,代码路径清晰可读。如果你在源码里看到类似expression_parser这样的模块,它负责把外部表达式字符串转换成内部表达式树,再进入模板实例化阶段。
不过算子融合是有边界的。它主要覆盖“逐行计算”类操作,比如算术、比较、逻辑运算。对于join、groupby、sort这类需要重排数据的算子,很难融合到同一个kernel里,只能通过中间数据在显存中传递。这也是为什么一个复杂管道里,每个阶段都要尽量在GPU上完成,而不是中途把结果搬回CPU。
4. 源码关键模块评测:列、表、hash join与groupby的实现路径
4.1 Column与Table:所有操作的最小单位
libcudf中所有操作都围绕column和table展开。table可以简单理解成std::vector<std::unique_ptr<column>>的集合,但在源码里它还有自己的表视图、行计数、列数等元信息。日常接触最多的是table_view,它是vector<column_view>,不拥有数据。这意味着多数算法只读取输入,不强制要求调用方让出数据所有权,对接口设计非常友好。
要理解cuDF的列为什么可以做到“切片不拷贝”,关键就在column的offset和size字段。一个从第100行开始、长度为50的子列视图,可以直接用原buffer加偏移表达,不需要复制50个元素。几乎所有算法在启动kernel时都会基于原始buffer偏移计算出真实地址,这也是为什么源码里大量使用指针加偏移的写法。
当数据类型是嵌套的,比如String,column内部会包含一个offsets子列,每个元素在chars子列中有一个起始偏移和结束偏移。读取字符串值时,需要先查offset,再去做字符搬运。这也是字符串操作比数值操作慢的原因之一:间址访问与变长数据很难做到全带宽合并访问。需要对字符串做大量高频操作时,建议先把字符串映射成category编码或者字典ID,能避开很多性能陷阱。
4.2 hash join:build、probe与哈希表的并发插入
join是数据处理里最常用的重操作。cuDF的hash join实现大致分两个阶段。第一阶段是build,从两张表中选择一张作为build表,把连接键的所有值插入GPU端哈希表。第二阶段是probe,遍历另一张表的每一行,用同样的哈希函数计算键的hash值,到哈希表里查找匹配行,命中则输出匹配结果。
源码里可以关注三个细节。第一,build阶段会优先选择小表,因为哈希表占用显存,选小表能减少内存和插入耗时。第二,并发插入用到了CAS原子指令,多个线程可能同时命中同一个bucket,必须保证插入不丢数据。如果冲突严重,哈希表退化严重,性能会大幅下滑。第三,为了处理重复键,库内会为每个桶维护冲突链或开放寻址的额外探测路径。这个设计直接决定了带重复键的大表join内存占用会有多高。
实际使用中,如果两张表都非常大,即使有spilling,依然可能在build表阶段把显存撑爆。建议先通过过滤和投影把参与join的列数压到最少,必要时拆分成多轮,比如按日期分区后逐区join,再合并结果。
4.3 groupby:hash聚合与两阶段合并
Pandas里的groupby操作习惯性按索引排序,但cuDF的groupby并不强制全局排序。源码中groupby采用hash聚合路径:先把分组键hash到不同分区,每个分区内部执行部分聚合,最后再合并各分区的中间结果。这样避免了全局排序的O(N log N)开销,更贴合GPU的并行模型。
这个实现里,分组数多少和每组数据多少会影响策略选择。如果分组数量很少,每个GPU线程可以直接在寄存器或共享内存里维护自己负责的key对应累计值。如果分组数很多,就需要全局哈希表,这时内存占用会明显上升。源码中有针对不同聚合类型和分组规模的分支,这也是为什么groupby在数据量变化时性能曲线不是单调递增的原因之一。
需要注意,cuDF groupby返回的组顺序默认与Pandas并不总是一致的,因为hash聚合天然不保证键的有序性。如果你的下游依赖分组有序,要在groupby之后显式调用sort,而不是默认假设。
4.4 表达式计算器:从表达式树到具体kernel
表达式计算器是现代cuDF里一个比较核心的模块。用户调用df['c'] = df['a'] + df['b']时,Python层会生成一个表达式对象,而不是立即执行循环。这个表达式对象进入libcudf后,会被展开成表达式树,再被编译期模板实例化为具体内核。
源码里cudf这个类会持有表达式树的根节点,比如column_reference、operation、literal。遍历时,每个节点都会对应一段代码生成逻辑,最终组装成一个完整的CUDA kernel调用。这样做的好处是省去了多次kernel启动。对于一个大DataFrame,一次融合kernel执行的时间可能比拆成三个kernel快几倍,因为每次kernel启动和中间数据写回都不是免费的。
不过表达式计算器并不是为了优化“万能场景”设计的,它更擅长算术、比较和简单的函数调用。碰到用户自定义UDF时,它往往无法融合,只能退回逐块执行。所以源码评测后得出的建议是:能用原生表达式表达的逻辑,尽量用表达式;需要UDF时,优先考虑用向量化或者分组重写,而不是在每一行调用Python函数。
5. 工程落地指南:从环境搭建到生产调优的完整实践
5.1 硬件平台与CUDA环境的匹配
cuDF对GPU架构有要求,至少需要支持CUDA的NVIDIA GPU,推荐主流产品线。安装时最省心的是用conda创建RAPIDS环境,官方会先把CUDA、RMM、libcudf这些依赖都整理好:
conda create -n rapids -c rapidsai -c conda-forge \ cudf=24.10 python=3.11如果你想在已有PyTorch环境里集成cuDF,一定要检查CUDA运行时版本是否一致。PyTorch自带CUDA运行库,cuDF也可能链接不同的CUDA版本,混装容易导致符号冲突。我吃过这个亏,明明nvidia-smi显示的驱动版本很高,但import cudf后再import torch就报错,最后发现是两个包分别带了不同版本的libcublas,解决办法是统一用conda从同一channel安装,或者用pip安装时对齐nvidia-*依赖。
物理机上还要确认GPU驱动版本与nvidia-container toolkit版本匹配,尤其是在容器里用Kubernetes调度GPU时,需要确保nvidia-smi能在容器内正常输出。不要小看这些环境问题,生产环境里一半以上的“cuDF很慢”其实都源于驱动或库版本不匹配,导致CUDA无法使用最新的内存管理特性。
5.2 显存规划与spilling配置
显存是硬约束。在设计数据管道时,先算一次理论峰值:假设有10列int64,1亿行就是大约8GB,加上结果集和中间hash表,轻松超过12GB。所以在开始写cuDF代码之前,我会先做个“数据体检”:
- 原始数据多大?哪些列需要被使用?
- 会做join吗?join键基数高不高?构建哈希表的内存大约是多少?
- 哪些中间结果可以立刻释放?
- 是否可以分批处理或行组过滤?
开启spilling通常有环境变量:
export CUDF_SPILL_DEVICE_MEM_TRY=ON export CUDF_SPILL_HOST_MEM_TRY=ON这里要分清楚spilling和UVM。UVM把host memory映射到GPU进程空间,靠缺页机制按需搬移,用户不感知;spilling则是cuDF显式把某些device buffer换出到host memory。前者适合小范围超限,后者更适合可预测的大规模数据绕过显存限制。但两者都不能当主内存用,一旦频繁搬移,性能会打折到让你怀疑人生。一个靠谱的策略是:先设定好显存预算,把数据切块,而不是完全依赖spilling。
5.3 与PyTorch、分布式框架的集成实践
数据工程最终经常要把处理结果交给深度学习框架。cuDF与PyTorch之间的传输,理论上可以直接走DLPack或者__cuda_array_interface__,不用经过CPU副本。例如把cudf Series转为torch Tensor,可以借助库内部的to_dlpack或torch.utils.dlpack能力。
import torch from torch.utils.dlpack import from_dlpack # cudf_col 是 cudf.Series torch_tensor = from_dlpack(cudf_col.to_dlpack())这个转换基本是零拷贝的,因为两边都指向同一块显存。但要注意生命周期管理:如果cudf对象被释放,显存可能被RMM池回收,而torch Tensor还持有旧指针。最好在转换后保留对原始cudf对象的引用,等torch这边的计算完成再释放。
分布式场景常用Dask-cuDF,它会将一个大数据集划分成多个分区,每个分区是一个cuDF DataFrame,由不同GPU处理。Dask负责调度分区间的shuffle、join和聚合。使用方式类似Dask DataFrame,只是在dask_cudf里读出来的分区天然跑在GPU上。如果你的集群里既有CPU节点又有GPU节点,要注意别让Dask把GPU任务调度到CPU节点上。
5.4 性能调优与常见坑:碎片化、同步开销与GPU OOM
我踩过最深的坑是“中间过程反复拷贝”。最开始写迁移脚本时,每个阶段处理完都会to_pandas()回主机内存看一眼,或者用pandas做一个小校验,结果一条GPU管道被拆成了几十个“GPU->CPU->GPU”片段,最终比纯Pandas还慢。看源码才意识到,to_pandas()会触发cudaMemcpy和同步,这会打断GPU流水线。
另一个高频问题是“小数据反而慢”。cuDF每个算子背后都有kernel启动和RMM内存池分配开销,启动一次CUDA kernel就有几十微秒级的耗时。如果你拿一万行数据反复调用几百次,当然会输给Pandas。正确姿势是利用批处理,把多个列的表达式融合成一次kernel调用。
排查GPU崩溃时,我经常先看系统日志里的gpu crash dump triggered关键字,这个通常意味着某个kernel发生了非法内存访问。可能是数组越界、空mask处理不对,或者不同版本之间ABI不兼容。首选工具是compute-sanitizer,它能报告具体kernel和触发位置,比盲改代码高效太多。还有,显存不足不一定是“崩溃”,可能只是抛出MemoryError或rmm::bad_alloc,此时优先看nvidia-smi的显存状态,再考虑重新设计join中间表。
6. 实测对比与反直觉结论:什么时候该用cuDF
6.1 一组直观的耗时对比
拿我手头一个实际环境,GPU是RTX 4090,测试数据50万到5000万行,结果大致如下(硬件不同会有差异,只作趋势参考):
| 数据规模/任务 | Pandas耗时 | cuDF耗时 | 加速倍数 |
|---|---|---|---|
| 100万行过滤+分组求和 | 420ms | 75ms | 约5.6倍 |
| 1000万行join(含hash表构建) | 6.8s | 370ms | 约18倍 |
| 5000万行多列groupby | 38s | 1.9s | 约20倍 |
| 1万行小表反复调用100次 | 200ms | 1.2s | 反而慢6倍 |
结论其实很清晰:数据规模越大,cuDF相对Pandas的优势越明显;数据量小、算子碎片化时,GPU反而会输。你可以把它理解成“集装箱货车”和“小轿车”:装满货物时货车性价比极高,但送个文件跑空车就非常浪费。
6.2 哪些场景不适合cuDF
虽然cuDF很香,但很多场景并不适合一上来就全量迁移。
- 数据量长期小于十万行,且算子调用很频繁。这种情况下,kernel启动和Python对象转换的开销可能主导耗时,纯CPU方案更稳。
- 需要大量Python UDF。比如对每一行调用一个复杂的外部函数,改写起来成本极高,这属于CUDA并行化的敌人,建议先把函数向量化,再回到cuDF。
- 超大字符串处理。比如上亿条短文本做正则替换,字符串在GPU上变长访问的劣势会被放大。可以先转成ID或长度特征再做计算。
- 可用显存非常小,比如8GB以下,且数据量动不动几十GB。如果spilling无法满足性能需求,不如用Dask在CPU上分片,至少调度逻辑更均匀。
6.3 落地几条经验
第一条,永远先做数据规模和显存预算的预估,再写代码。第二条,把数据管道拆成可以单独测试的算子,分别验证正确性,因为cuDF某些API的行为边界和Pandas不完全一致,比如索引默认不排序、groupby结果顺序不保证。第三条,尽量保持数据在GPU上连续流转,不要在一次处理过程中反复导入导出。
我自己的最后一条建议是:如果你刚开始接触cuDF,先用cudf.pandas这个模式在原来Pandas代码上跑通,它会自动把常见的Pandas调用路由到GPU路径,虽然不能覆盖所有API,但能让你快速看到提速潜力。跑通后,再针对热点阶段手工改成原生cuDF接口,收益会远超一上来就全局重写。
我把这套方法用在一个日增几十GB日志的项目上,最终把原来的40分钟管道压到4分钟以内,而且因为中间结果不再写磁盘做序列化,单机内存占用也降下来了。所以说到底,不要盲目追求“把所有代码改成GPU”,而是先把数据路径和瓶颈找出来,再让cuDF在它擅长的地方承担计算任务。