Python dask 分布式计算:pandas 跑不动的数据,dask 怎么解决

发布时间:2026/9/10 23:55:57

Python dask 分布式计算:pandas 跑不动的数据,dask 怎么解决 Python dask 分布式计算pandas 跑不动的数据dask 怎么解决一张表 2000 万行pandas 直接 OOM 罢工试试 dask用法跟 pandas 几乎一样但数据再多也不怕。一、pandas 到 dask不是替代是升级做数据分析最崩溃的瞬间是什么对我来说是跑了 20 分钟的 pandas 脚本最后弹出一个MemoryError。pandas 的设计哲学是把数据全部加载到内存里操作。单机内存 16G、32G 的时候处理百万级数据没啥问题。但现在的业务数据动不动就上千万行、几十个字段宽表一张表几个 G 很常见pandas 是真的扛不住了。dask 的思路不一样它不把数据一次性全读进内存而是分块读、分块算用完一块就释放。具体来说dask 做了两件事惰性计算Lazy Evaluation你写操作的时候它不执行只是构建一个计算图。等你调用.compute()的时候它才真正执行。这样就可以在计算前做全局优化。任务分片Task Partitioning把大数据拆成很多小数据块partition每个块独立处理。一个块算完就释放内存再处理下一个。所以 dask 不是 pandas 的替代品而是当你数据大到 pandas 撑不住时的升级方案。日常小数据继续用 pandas大数据场景切 dask。为什么 dask 不是pandas 替代品而是pandas 超时后的升级方案这里有个性能悖论在处理 10 万行数据时pandas 比 dask 快 3-5 倍。原因是 dask 的惰性计算图构建、分片调度、序列化/反序列化都是有固定开销的——处理的数据量越小这些固定开销占的比重就越高。pandas 直接对内存中的 C 数组做向量化操作没有中间商。只有当真的一台机器装不下数据时比如 2000 万行 CSV 文件 16GBdask 的分批加载优势才开始抵消调度开销。所以不要用 dask 处理小数据——那等于抱着消防栓浇花。二、dask DataFrame跟 pandas 几乎一样的 APIdask 最友好的一点是它的 API 高度模仿 pandas。大部分时候你把import pandas as pd改成import dask.dataframe as dd代码基本不用动。import dask.dataframe as dd import pandas as pd # # 1. 读取 CSV —— 跟 pandas 语法一模一样 # # dask 不会立刻读文件只是记录了要读这个文件的元信息 # 用 blocksize 控制每个分片的大小建议 64MB~128MB df dd.read_csv( sales_2024.csv, blocksize64MB, # 每个分片 64MB控制内存峰值 dtype{user_id: str}, # 指定类型避免自动推断出错 parse_dates[order_date] ) # # 2. 基础操作 —— 写法跟 pandas 完全一致 # # 过滤、分组、聚合语法和 pandas 一毛一样 filtered df[df[amount] 100] # 筛选高客单价订单 monthly_sales ( df.groupby(df[order_date].dt.month) # 按月分组 [amount].sum() # 求和 ) # # 3. 触发计算 —— 这是 dask 和 pandas 最大的区别 # # .compute() 会把 Dask DataFrame 转成 pandas DataFrame # 调用之前所有操作都只是计划现在才真正执行 result monthly_sales.compute() # 此时才真正执行计算 print(result)需要注意一个细节dask 的惰性计算意味着每次调用.compute()都会重新从头算。如果你要对同一个数据做多次聚合最好先.persist()把中间结果缓存到内存# 错误做法每次 .compute() 都重算一遍 avg df[amount].mean().compute() # 第1次从头算 std df[amount].std().compute() # 第2次又从头算 !!! # 正确做法persist 缓存中间结果 df_cached df[[amount, order_date, region]].persist() avg df_cached[amount].mean().compute() # 从缓存算 std df_cached[amount].std().compute() # 从缓存算快很多三、dask 调优partition 是核心dask 的性能瓶颈通常出在 partition 数量的选择上。太少每个分片太大内存还是撑不住太多调度开销比计算开销还大。一个实用的经验公式每个 partition 大小控制在 100MB 左右。比如你的数据是 5G那就设 50 个分区。import dask.dataframe as dd df dd.read_csv(big_data.csv, blocksize64MB) # # 查看当前分片数量 —— 调优前先看一眼 # print(f当前分片数: {df.npartitions}) # # 调整分片数repartition 重新分配 # # npartitions 太小 每个分片太大 内存不够 # npartitions 太大 调度开销 计算开销 反而变慢 # 经验值每个 partition 约 100MB file_size_mb 5000 # 假设文件 5G target_partitions file_size_mb // 100 # 每个分片 100MB df_optimized df.repartition(npartitionstarget_partitions)另一个性能杀手是数据倾斜。比如分组聚合时某个 key 的数据量特别大它所在的分片就会成为瓶颈。可以用map_partitions在分片内部做预聚合来缓解# 分片内预聚合 —— 减少 shuffle 的数据量 def pre_agg_in_partition(partition): 在每个分片内先做一次聚合减少后续 shuffle 的数据量 return partition.groupby(user_id)[amount].sum().reset_index() # map_partitions 对每个分片独立执行函数 pre_aggregated df.map_partitions(pre_agg_in_partition) # 再做全局聚合数据量小了很多 final pre_aggregated.groupby(user_id)[amount].sum().compute()四、dask 分布式从单机到集群dask 单机模式可以解决内存不够的问题但要真正加速计算还得上分布式。from dask.distributed import Client, LocalCluster # # 方式一本地模拟分布式开发调试用 # cluster LocalCluster( n_workers4, # 4个worker进程 threads_per_worker2, # 每个worker 2线程 memory_limit4GB # 每个worker内存上限 ) client Client(cluster) # # 方式二连接已有集群生产环境用 # # client Client(scheduler-address:8786) # # 设置后所有 dask 操作自动用集群资源 # import dask.dataframe as dd df dd.read_csv(s3://my-bucket/big_data/*.csv) # 聚合操作会自动分布到各个 worker 上 result df.groupby(category)[sales].sum().compute() # 查看 dashboard默认端口 8787 # 可视化任务执行情况、内存使用、数据传输 print(fDashboard: {client.dashboard_link})分布式场景下有几个实用技巧用client.persist()代替.persist()数据会分散缓存在多个 worker 上避免全局排序.sort_values()需要把所有数据 shuffle 到一个 worker非常慢优先用有损操作比如approx_percentile替代精确百分位速度快很多五、总结 踩坑提醒dd.read_csv的blocksize不是每个分区精确大小如果你设blocksize64MBdask 会从 64MB 的偏移量开始切割文件。但 CSV 的行边界不一定在 64MB 处对齐——切割点可能正好落在某一行中间导致该行被截断。dask 会在每个分片的末尾往前找到下一个换行符所以实际分片大小可能比 64MB 略大。如果你对分片数量有严格要求读取后用.repartition()做二次调整。惰性求值 随机抽样会得到不随机的结果df.sample(frac0.01).compute()每次执行都会重新随机抽样但由于 dask 的惰性计算图是确定的如果你没设random_state两次compute()的抽样结果可能完全不同——这在需要对抽样结果做验证分析的场景下是个坑。建议每次.compute()前重新设一次random_state或者persist()后再抽样。分布式集群中memory_limit不是硬限制你设了memory_limit4GB期望 worker 到了 4GB 就停止接收新任务。但实际上 dask 的 memory_limit 是软限制——它只在 scheduler 调度新任务时检查当前内存使用不会主动 kill 已运行的任务。如果你的一个任务在运行中自己申请了 8GB比如一个大排序dask 不会阻止worker 会直接 OOM 被杀。解决方案是同时设置--memory-limit和--memory-spill4GB 时 spill 到磁盘或者在代码层面控制单任务数据量。dask 的核心价值就一句话用 pandas 的写法处理 pandas 吃不下的数据。它不要求你学一套全新的 API大部分代码从 pandas 迁移过来只是换个 import。但它通过惰性计算和分片机制让单机能处理的场景从百万行升级到十亿行。在实际工作中我现在的默认策略是100 万行以内pandas简单直接100 万 ~ 1000 万行dask 单机模式加个dd.read_csv就行1000 万行以上dask 分布式集群并行加速当然 dask 也有不擅长的场景比如频繁的 join 操作、需要全局排序的复杂计算。这些时候就该 Spark 上场了后面我们会聊到。为什么 dask 不擅长 JOIN 和全局排序JOIN 在 dask 里需要对两个 DataFrame 的 partition 做对齐——如果两个 DataFrame 的分区数不同或分区方式不同比如按不同列做了索引JOIN 时就需要 shuffle 全部数据等于把所有数据在集群里洗一遍牌。Spark 有专门的 Sort-Merge Join 和 Broadcast Join 优化器能根据数据量自动选择最佳策略dask 在这方面远没有 Spark 成熟。全局排序同理——sort_values()需要把所有数据集中到一个 partition如果数据是 100GB那这个 partition 就是 100GB——直接炸。在这些场景下dask 更像是能跑但不保证快而 Spark 提供了生产级的优化保障。本文由朱大喜原创欢迎点赞收藏有问题评论区交流~
延伸阅读

更多相关文章

2026/9/8 14:45:44

OptiScaler终极问题解决指南:从新手到专家的完整教程

OptiScaler终极问题解决指南:从新手到专家的完整教程 【免费下载链接】OptiScaler OptiScaler bridges upscaling/frame gen across GPUs. Supports DLSS2/XeSS/FSR2 inputs, replaces native upscalers, enables FSR-FG/XeFG on non-FG titles. Supports Nukem mod…

2026/9/9 19:21:37

告别简陋界面:用foobox-cn打造专业级音乐播放器体验

告别简陋界面:用foobox-cn打造专业级音乐播放器体验 【免费下载链接】foobox-cn DUI 配置 for foobar2000 项目地址: https://gitcode.com/GitHub_Trending/fo/foobox-cn 还在为foobar2000那单调的默认界面而烦恼吗?你是不是也渴望拥有一个既美观…

2026/9/9 16:52:56

CANN/asc-devkit TensorDesc GetIndex方法文档

GetIndex 【免费下载链接】asc-devkit 本项目是CANN 推出的昇腾AI处理器专用的算子程序开发语言,原生支持C和C标准规范,主要由类库和语言扩展层构成,提供多层级API,满足多维场景算子开发诉求。 项目地址: https://gitcode.com/c…

2026/9/10 23:54:43

K3S 基础命令集

K3S 基础命令集Pod查看 Pod查看所有Pod查看详情查看Pod 日志查看 IP查看yaml编辑配置删除Pod强制删除进入容器PVC查看所有PVC查看详情查看PVC YAML创建 PVC编辑 PVC删除 PVCService查看 Service查看 Service 详情查看 Service YAML创建 Service编辑 Service删除 Servicedepoly查…

2026/9/10 23:54:43

2027届论文降AI率平台哪个靠谱?六款实测横评

高校对AI生成内容的检测力度逐年加强,2027届毕业生面临的降AI率压力比往届更大。论文写作过程中适度使用AI辅助已成常态,但如何让最终稿顺利通过检测,成为不少学生头疼的问题。这篇测评选取市面上讨论度较高的六款降AI率工具,逐一…

2026/9/10 23:54:43

2027届论文降AI率工具实测,六款平台谁更靠谱

高校对论文AI生成内容的检测逐年收紧,降AI率从加分项变成了毕业答辩前的硬性门槛。不少2027届学生已经开始为学位论文发愁,市面上号称能降AI率的平台数量庞大,实际效果却参差不齐。本文选取六款有代表性的工具,用同一批论文样本做…

2026/9/10 23:54:43

生物医学SCI三区期刊投稿策略与高录用率期刊推荐

1. 生物医学领域投稿策略解析 在科研论文发表的道路上,选择合适的期刊往往是决定成败的关键一步。作为一名在生物医学领域发表过十余篇SCI论文的研究者,我深知投稿过程中的焦虑与期待。今天我要分享的这个期刊选择策略,可能会颠覆你对SCI期刊…

2026/9/10 23:49:43

Tracy Profiler 快速上手实战指南:4步定位游戏掉帧元凶

Tracy Profiler 快速上手实战指南:4步定位游戏掉帧元凶 【免费下载链接】tracy Frame profiler 项目地址: https://gitcode.com/GitHub_Trending/tr/tracy 你正在跑一个游戏,画面突然卡了一瞬,重跑一遍问题又消失了,日志里…

2026/9/10 16:39:38

超人会飞不算本事:系统稳定依赖清晰规则与边界设计

开头先不绕弯子。“#斯坦李吐槽dc 所以超人是无缘无故会飞的嘛哈哈哈哈哈哈哈锤哥真是技术人才啊!#雷神 #复联”这类调侃式短标题,第一波冲击力在于它把两个宇宙的角色塞进同一个吐槽箱里,但细想一下就能发现,它真正碰到的根本不是…

2026/9/10 11:16:38

超人VS蜘蛛侠:拆解超级IP的影响力与传播方法论

把“蜘蛛侠 vs 超人”放在 CSDN 上聊,可能很多人第一反应是走错片场了。但如果把这两个角色看成“两个持续运营了 80 多年的文化产品”,你会发现,这场比较本质上是两个不同 IP 策略的长期结果对比:超人赢在定义了整个超级英雄题材…

2026/9/9 16:31:09

基于CNN的调制信号识别:MATLAB实现时频图分类实战

简介:本资源是一套面向通信工程与信号处理方向学习者、研究者的深度学习实践方案,聚焦调制信号自动检测与识别这一典型无线通信任务,解决传统方法依赖人工特征、低信噪比下性能下降等痛点。压缩包共12个文件(10.73MB)&…

2026/9/10 0:00:55

目录对比去重实战:用哈希算法精准清理重复文件

我电脑里现在还有一块换了三次机的“数据墓地”硬盘,里面存着2016年以前所有旧笔记本的完整备份。平时不觉得有什么,直到前阵子想把它整理归档,发现同一个安装包、同一批照片、同一份论文草稿,在几个不同的备份目录里反复出现。更…

2026/9/10 0:00:55

Leaflet离线地图完整Demo合集:内网部署与坐标纠偏实战

简介:这是一份面向Web GIS开发者的LeafLet离线地图示例合集,帮助开发者快速掌握离线地图从搭建到交互的完整流程。压缩包共723个文件,大小14.06MB,以319个js脚本、175个html页面和29个css样式文件为主体,配合png/svg图…

2026/9/10 0:00:55

MATLAB读取Rinex 3.02观测文件:多系统GNSS数据解析实战

简介:基于MATLAB开发的Rinex3.02版观测文件(o文件)读取代码包,面向卫星定位导航方向的学习者与研究人员,用于解决新版观测文件的数据解析、历元提取与时间转换问题。压缩包共4个文件,包含两个m脚本、一个19…

2026/9/10 12:32:02

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/10 15:19:50

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/10 15:49:53

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码