You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Dask处理20GB大CSV的分组与复杂函数应用问题

处理20GB大型CSV文件的Dask技术疑问

数据结构

SID             DATETIME            LAT        LON         HEADING  SPEED NAME 
0  0ff58f68b3a1f  2023-06-06 16:47:38  43.589027  3.813129    297.0    7.0  AAAAAA
1  798cd24a678e1  2023-06-06 09:53:34  43.588792  3.812747    226.0   25.0  BBBBB 
3  0ff58f68b3a1f  2023-06-06 16:47:32  43.589069  3.813142    129.0    4.0  AAAAAA 
4  0ff58f68b3a1f  2023-06-06 16:47:33  43.589062  3.813133    217.0    5.0  AAAAAA

处理需求

  • 按SID分组;
  • 每组内按DATETIME排序;
  • 应用复杂函数getTripStatistics(完成距离计算、时长统计、地理点REST请求补充数据),每组输出一行聚合信息;
  • 导出结果数据集。

编写的Dask代码

import dask.dataframe as dd

names = ['SID','DATETIME','LAT','LON','HEADING','SPEED','NAME']
ddf = dd.read_csv(FILE, header=0, names=names)

# 按datetime排序
ddf1 = ddf.set_index('DATETIME').reset_index()

# 方案1
ddf2 = ddf1.groupby('SID').apply(getTripStatistics, meta=(None, 'object'))

# 方案2
ddf2 = ddf1.map_partitions(lambda df: df.groupby('SID').apply(getTripStatistics), meta=(None, 'object'))

ddf2.compute().to_csv("./tmp/my_one_file.csv", index=False)

测试补充与疑问

测试补充

  • 1个分区时,两种方案结果数量一致(长度683),耗时约12.68秒;
  • 重分区为4个后,方案2结果长度变为811(耗时约14.05秒),方案1仍为683(耗时约13.53秒)。

疑问

  1. read_csv后的分区数量是否取决于数据集大小?是否会导致内存溢出?
  2. 按SID分组后,是否每个组会单独成为一个分区?是否需要手动调用repartition?如何实现“分组后每组对应一个分区”?
  3. 测试小数据集时方案1和方案2结果一致,是否因为数据仅在一个分区中?
  4. getTripStatistics函数是性能瓶颈,针对Dask中复杂多步骤函数的应用,有哪些优化技巧?
  5. 重分区为4个后,方案2结果长度增加,是忽略了部分分区数据吗?

问题解答

1. read_csv分区与内存溢出问题

  • Dask的read_csv默认按文件大小划分分区,默认单分区约64MB(可通过blocksize参数调整),数据集越大,分区数自然越多。
  • 只要单分区大小不超过机器可用内存,就不会触发内存溢出。如果你的单分区设置过大(比如远大于内存),才会出问题。建议根据机器内存调整blocksize,比如16G内存可设为100-200MB/分区。

2. 分组后的分区逻辑与“每组一个分区”实现

  • 默认情况下,Dask分组后不会自动让每个组对应一个分区。分组操作会先做shuffle,把同一SID的数据集中到同一分区,但一个分区可能包含多个组,且单个组不会跨分区。
  • 不建议手动调用repartition实现每组一个分区,因为这会生成大量小分区,反而降低性能。如果确实需要,可在分组后结合repartition(npartitions=总组数),但前提是你明确总组数,且组数不多。

3. 小数据集下两方案结果一致的原因

  • 是的。当小数据集仅占一个分区时,map_partitions里的分组是对整个数据集执行的,和全局groupby.apply效果完全一致,所以结果相同。当数据多分区时,map_partitions会在每个分区内独立分组,同一SID的数据如果分散在多个分区,就会被多次处理,这就是方案2结果变长的核心原因。

4. 复杂函数getTripStatistics的优化技巧

  • 批量处理REST请求:把单条数据的REST请求改成批量请求(比如攒一批地理点再发请求),减少IO等待时间。
  • 向量化计算:替换函数内的逐行循环(比如距离计算)为Pandas/NumPy的向量化操作,大幅提升计算速度。
  • 缓存复用:给REST请求加本地缓存(比如functools.lru_cache或Redis),避免重复请求相同地理点。
  • 提前过滤数据:分组前先过滤无用数据,减少每个分组的处理量。
  • 用Dask Delayed封装:如果函数里有无法向量化的步骤,用dask.delayed封装,让Dask更好地并行调度任务。
  • 调整并行度:根据CPU核心数设置dask.config.set(scheduler='processes', num_workers=核心数),充分利用多核资源。

5. 多分区下方案2结果变长的原因

  • 不是忽略数据,而是重复处理了跨分区的同一SID数据。方案2中map_partitions在每个分区内独立分组,同一个SID的数据如果分布在多个分区,每个分区里的该SID子数据集都会单独传入getTripStatistics生成一行结果,最终合并后同一SID会有多行输出,导致总长度增加。而方案1的全局groupby.apply会先聚合同一SID的所有数据,再处理生成一行结果,所以结果长度正确(683)。方案2不符合你的需求,应该用方案1。

内容的提问来源于stack exchange,提问作者user2894156

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 12:25:45