如何用Dask并行处理DataFrame中无返回值的逐行任务?
加速DataFrame逐行本地文件写入任务的可行方案
针对你需要加速DataFrame逐行写入本地文件的需求(示例为文本,实际是图片),这里整理几种可行方案,包括你已经修正后可用的Dask实现,以及其他替代思路:
已验证可行的Dask.apply实现
你当前的代码已经能正常运行,关键是正确指定了meta参数——因为你的函数无返回值,Dask需要明确返回类型,所以(None, 'object')是合适的设置:
import dask.dataframe as dd import pandas as pd doc = pd.DataFrame({'file_name': ['Bob', 'Jane', 'Alice','Allan'], 'text': ['text1','text2', 'text3','text4']}) def func(row): with open(row['file_name']+'.txt', 'w') as f: f.write(row['text']) ddf = dd.from_pandas(doc, npartitions=2) k = ddf.apply(func, axis=1, meta=(None,'object')) k.compute()
其他高效替代方案
1. Dask map_partitions(更低调度开销)
逐行apply会有较多调度成本,改用按分区处理能提升效率:把每个分区的DataFrame传给处理函数,在函数内完成逐行写入:
def process_partition(df): for _, row in df.iterrows(): with open(row['file_name']+'.txt', 'w') as f: f.write(row['text']) return df # 也可返回None,需对应调整meta参数 ddf = dd.from_pandas(doc, npartitions=2) # 这里meta传入原DataFrame的结构,匹配返回值 result = ddf.map_partitions(process_partition, meta=doc).compute()
2. 标准库concurrent.futures并行
如果不想依赖Dask,直接用Python标准库就能实现并行:
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor def func(row): with open(row['file_name']+'.txt', 'w') as f: f.write(row['text']) # 把DataFrame转成字典列表,方便逐个传参 rows = doc.to_dict('records') # IO密集型任务(比如文件写入)用ThreadPoolExecutor更高效 with ThreadPoolExecutor() as executor: executor.map(func, rows) # 如果任务包含CPU密集型计算,改用ProcessPoolExecutor # with ProcessPoolExecutor() as executor: # executor.map(func, rows)
3. swifter自动并行
swifter库会自动判断任务适合串行还是并行,底层会自动调用Pandas或Dask实现最优处理:
import swifter doc.swifter.apply(func, axis=1)
内容的提问来源于stack exchange,提问作者Avatrin
相关产品推荐
相关产品推荐

