基于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秒)。
疑问
read_csv后的分区数量是否取决于数据集大小?是否会导致内存溢出?- 按SID分组后,是否每个组会单独成为一个分区?是否需要手动调用
repartition?如何实现“分组后每组对应一个分区”? - 测试小数据集时方案1和方案2结果一致,是否因为数据仅在一个分区中?
getTripStatistics函数是性能瓶颈,针对Dask中复杂多步骤函数的应用,有哪些优化技巧?- 重分区为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
相关产品推荐
相关产品推荐

