如何用Dask read_csv加载大CSV文件中指定分段的数据?
好问题!这种分段式的大型CSV文件确实很常见,Dask完全能帮你精准加载需要的区块,不用把整个文件硬塞进内存。我给你分享两种实用的方案,你可以根据自己的需求选择:
方案一:先扫描行号,再用Dask read_csv精准加载
这个方案适合文件特别大的场景,利用Dask原生的并行读取能力,效率很高。核心思路是先快速扫一遍文件,记录A、B两段的行号范围,再告诉Dask只读取对应范围的内容:
第一步:编写行号扫描函数
这个函数只会逐行读取文件的标记行,不会加载所有数据到内存,速度非常快:
def find_segment_boundaries(file_path): with open(file_path, 'r') as f: line_num = 0 a_start = None a_end = None b_start = None b_end = None for line in f: line = line.strip() line_num += 1 # 跳过前两行注释 if line_num <= 2: continue # 标记A段起始行 if a_start is None: a_start = line_num # 找到A段结束标记 if line == '0 /END OF A DATA': a_end = line_num b_start = line_num + 1 continue # 找到B段结束标记后停止扫描 if b_start is not None and line == '0 /END OF B DATA': b_end = line_num break # 计算各段的有效行数 a_rows = a_end - a_start b_rows = b_end - b_start return a_start, a_rows, b_start, b_rows
第二步:用Dask加载指定分段
调用上面的函数拿到行号后,就可以用read_csv的skiprows和nrows参数精准加载:
import dask.dataframe as dd # 获取分段行号信息 a_start, a_rows, b_start, b_rows = find_segment_boundaries('datafile.csv') # 加载A段:跳过前2行注释,读取a_rows行数据 df_a = dd.read_csv( 'datafile.csv', skiprows=2, nrows=a_rows, header=None, quotechar="'" ) # 加载B段:跳过B段之前的所有行,读取b_rows行数据 df_b = dd.read_csv( 'datafile.csv', skiprows=b_start - 1, nrows=b_rows, header=None, quotechar="'" ) # 现在可以对分段数据做统计,比如计算各列总和 print("A段列总和:", df_a.sum().compute()) print("B段列总和:", df_b.sum().compute())
方案二:用dask.delayed自定义读取逻辑
如果不想预先扫描行号,也可以用dask.delayed包装一个自定义读取函数,直接根据结束标记停止读取,灵活性更强:
自定义分段读取函数
这里用csv.reader来处理可能的引号和转义,避免解析错误:
from dask import delayed import pandas as pd import csv def read_a_segment(file_path): data = [] with open(file_path, 'r') as f: reader = csv.reader(f, quotechar="'") # 跳过前两行注释 next(reader) next(reader) for row in reader: # 匹配A段结束标记 if row[0] == '0' and len(row) >= 2 and row[1] == '/END OF A DATA': break data.append(row) # 转为DataFrame并指定数据类型(避免默认的object类型) df = pd.DataFrame(data, columns=[f'col{i}' for i in range(1, 5)]) df = df.astype({f'col{i}': int for i in range(1, 5)}) return df def read_b_segment(file_path): data = [] with open(file_path, 'r') as f: reader = csv.reader(f, quotechar="'") # 跳过直到A段结束 for row in reader: if row[0] == '0' and len(row) >= 2 and row[1] == '/END OF A DATA': break # 读取B段数据直到结束标记 for row in reader: if row[0] == '0' and len(row) >= 2 and row[1] == '/END OF B DATA': break data.append(row) df = pd.DataFrame(data, columns=[f'col{i}' for i in range(1, 8)]) df = df.astype({f'col{i}': int for i in range(1, 8)}) return df
转为Dask DataFrame并统计
# 将自定义函数转为延迟对象,再生成Dask DataFrame df_a = dd.from_delayed(delayed(read_a_segment)('datafile.csv')) df_b = dd.from_delayed(delayed(read_b_segment)('datafile.csv')) # 执行统计操作 print("A段列均值:", df_a.mean().compute()) print("B段列均值:", df_b.mean().compute())
关键注意事项
- 列数差异处理:你的A段是4列,B段是7列,必须分开加载,不能同时读取整个文件,否则Dask会因为列数不统一报错。
- 解析准确性:如果文件里有带引号的字段,一定要用
csv.reader或者read_csv的quotechar参数,避免用split(',')导致解析错误。 - 性能选择:如果分段数据量极大,优先选方案一,因为Dask的
read_csv支持并行分块读取;如果文件结构多变,方案二的灵活性更高。
内容的提问来源于stack exchange,提问作者Tims
相关产品推荐
相关产品推荐

