如何使用Dask对URL列表中的文件执行处理函数实现并行计算加速
并行化实现建议
你的任务属于典型的IO密集型场景,大部分时间消耗在网络请求、数据库写入上,不需要用复杂的Dask,用Python标准库自带的多线程工具就能得到很好的加速效果,实现成本也低。
方案1:用concurrent.futures.ThreadPoolExecutor实现(最推荐)
IO密集型任务多线程开销远低于多进程,也不需要处理进程间数据传递的问题,非常适配你的场景,修改成本极低,只需要调整主逻辑代码即可:
from concurrent.futures import ThreadPoolExecutor criteria = 'temperature' threshold = 22 filenames =[url1.html, url2.html, url3.html] # max_workers可以根据机器配置、站点服务器压力调整,IO密集型任务一般可以设为CPU核心数的4~10倍,也可以直接不填,默认值为CPU核心数*5 with ThreadPoolExecutor(max_workers=20) as executor: for file in filenames: executor.submit(get_and_filter_data, file, criteria, threshold)
使用该方案需要注意几个细节:
- 数据库连接部分建议用连接池实现,避免每个任务都新建连接导致连接数过多被数据库拒绝,可直接用SQLAlchemy的内置连接池,或者MySQL官方的连接池工具。
- 给网络请求加超时参数,避免个别站点响应过慢卡住整个任务。
- 若担心对站点服务器并发请求量过高被封禁,可以适当调低
max_workers的数值。
方案2:先优化现有串行代码,再做并行(收益更高)
你当前的串行代码本身还有可优化的空间,改完后即使不并行也能提升不少运行速度:
- 替换
filter_criteria函数里的iterrows遍历逻辑,改用pandas向量化操作,速度可提升数倍:
def filter_criteria(df,threshold,criteria): if criteria in df.columns: return df[pd.to_numeric(df[criteria], errors='coerce') >= threshold].index else: return []
- 省略XML下载到本地临时文件的步骤,直接把内容读到内存里解析,省掉磁盘IO开销:
import requests def xml_to_dataframe(ood_xml): resp = requests.get(ood_xml, timeout=10) prstree = ETree.fromstring(resp.content) root = prstree.getroot() all_obs = [] for obs in root.iter('observations'): ood_observation = [] for n, param in enumerate(list_parameters): x=obs.find(variable_to_check).text ood_observation.append(x) all_obs.append(ood_observation) return(pd.DataFrame(all_obs, columns=list_parameters))
其他可选方案
如果后续你需要扩展到分布式多机器跑任务,可以再考虑用Dask或者Celery,当前单机场景下ThreadPoolExecutor完全够用,实现成本最低,提速效果也最明显,原300秒的串行IO任务,开20个线程一般可以压缩到20秒以内。
内容的提问来源于stack exchange,提问作者Guimeteo
相关产品推荐
相关产品推荐

