如何将基于asyncio/aiohttp的数据库查询脚本适配到dask.distributed.Client(asynchronous=True)
如何将asyncio/aiohttp的数据获取流程与Dask异步Client集成以加速DataFrame合并和计算?
我正尝试搭建一个具备以下功能的脚本:
- 通过HTTP从数据库执行多个不同查询
- 将查询结果解析为Pandas DataFrame
- 合并DataFrame
- 对合并后的数据执行其他密集型计算
我已查阅Dask相关文档,但仍无法理解如何将当前的asyncio/aiohttp架构与Dask集成(以加速DataFrame合并及后续计算)。现有代码可正常运行,但DataFrame合并属于CPU密集型操作,速度极慢;同时我还需要对获取到的长DataFrame列表进行聚合及其他计算。我曾尝试参考Dask的“Custom computation: Tree summation”示例来并行化合并操作,但帮助有限。由于我对asyncio和Dask均不熟悉,不确定当前实现是否正确。
现咨询:如何将上述代码适配到dask.distributed.Client(asynchronous=True)?
现有代码如下:
import asyncio, aiohttp from io import StringIO import pandas as pd async def fetch_html(session, field1, field2): query = 'some InfluxDB Flux query {f1} {f2}' response = await session.get(query.format(f1=field1, f2=field2)) return response async def get_data(session, field1, field2): response = await fetch_html(session, field1, field2) content = await response.content.read() df = pd.read_csv(StringIO(content), parse_dates=['time'], infer_datetime_format=True)\ .set_index(['time','other_index']) return df async def get_data_field1(session, field1): tasks = [] for field2 in ['val1', 'val2']: tasks.append(get_data(session, field1, field2)) df_list = await asyncio.gather(*tasks) df_joined = df_list[0].merge(df_list[1], right_index=True, left_index=True, suffixes=('','')).droplevel('other_index') return df_joined def get_all_data(field_list): async with aiohttp.ClientSession() as session: tasks = [] for field1 in field_list: tasks.append(get_data_field1(session, field1)) df_list = await asyncio.gather(*tasks) return df_list if __name__=='__main__': df_list = asyncio.run(get_all_data(['foo','bar','baz'])) #very long list # code that aggregates df_list
解决方案
核心思路是:用asyncio/aiohttp处理IO密集型的HTTP查询(这部分它们擅长),把CPU密集型的DataFrame合并、计算任务交给Dask异步Client,让Dask集群的workers并行处理这些任务。下面是具体的适配步骤和修改后的代码:
1. 准备工作:安装依赖并理解核心分工
首先确保你安装了Dask分布式模块:
pip install dask distributed
我们要明确分工:
- IO密集操作(HTTP请求、数据读取):留在asyncio/aiohttp中处理,发挥其异步IO优势
- CPU密集操作(DataFrame合并、聚合计算):封装成普通函数,交给Dask集群并行执行
2. 修改后的完整代码
import asyncio import aiohttp from io import StringIO import pandas as pd from dask.distributed import Client # -------------------------- # CPU密集型任务:封装成普通函数(交给Dask执行) # -------------------------- def merge_two_dfs(df_pair): """合并两个DataFrame,封装为普通函数方便Dask调用""" df1, df2 = df_pair return df1.merge(df2, right_index=True, left_index=True, suffixes=('', '')).droplevel('other_index') def aggregate_all_dfs(df_list): """聚合所有合并后的DataFrame,替换为你的实际业务逻辑""" combined_df = pd.concat(df_list) return combined_df.groupby(level='time').agg({'value': ['mean', 'max']}) # -------------------------- # IO密集型任务:保留asyncio/aiohttp实现 # -------------------------- async def fetch_html(session, field1, field2): query = 'some InfluxDB Flux query {f1} {f2}' response = await session.get(query.format(f1=field1, f2=field2)) return response async def get_data(session, field1, field2): response = await fetch_html(session, field1, field2) content = await response.content.read() df = pd.read_csv(StringIO(content), parse_dates=['time'], infer_datetime_format=True)\ .set_index(['time','other_index']) return df # -------------------------- # 整合Dask异步Client:提交CPU任务 # -------------------------- async def process_single_field1(session, client, field1): # 1. 异步获取当前field1对应的两个DataFrame(IO操作) df1, df2 = await asyncio.gather( get_data(session, field1, 'val1'), get_data(session, field1, 'val2') ) # 2. 把合并任务提交给Dask集群,返回Future对象 merge_future = client.submit(merge_two_dfs, (df1, df2)) return merge_future async def get_all_data(field_list): # 初始化Dask异步客户端,与asyncio事件循环兼容 async with Client(asynchronous=True, n_workers=4) as client: async with aiohttp.ClientSession() as session: # 提交所有field1的处理任务,得到一批Future futures = [] for field1 in field_list: future = await process_single_field1(session, client, field1) futures.append(future) # 收集所有合并后的DataFrame结果 merged_dfs = await client.gather(futures) # 提交最终聚合任务给Dask aggregate_future = client.submit(aggregate_all_dfs, merged_dfs) final_result = await aggregate_future return final_result if __name__=='__main__': # 运行异步主函数 final_result = asyncio.run(get_all_data(['foo','bar','baz'])) # 替换为你的长列表 print(final_result)
3. 关键改动说明
- 异步Dask客户端:用
Client(asynchronous=True)创建客户端,它能直接在asyncio事件循环中工作,无需额外线程切换 - 任务拆分:把CPU密集的合并、聚合逻辑抽成普通函数,通过
client.submit()提交给Dask集群,返回的Future可以在async代码中等待结果 - 结果收集:用
client.gather()批量获取Dask任务的结果,比asyncio.gather更适合分布式场景,能自动处理集群任务的调度 - 集群配置:初始化Client时可以指定
n_workers(比如和CPU核心数一致),或者连接远程Dask集群进一步提升算力
4. 额外优化建议
- 改用Dask DataFrame:如果数据量极大,单个Pandas DataFrame内存压力大,可以把解析后的Pandas DataFrame转成Dask DataFrame,用Dask内置的
merge、concat方法,实现全流程分布式处理 - 批量提交任务:如果field_list特别长,可以分批次提交任务,避免一次性压满集群资源
- 监控任务状态:可以通过Dask的Dashboard(默认
http://localhost:8787)查看任务执行情况,排查瓶颈
内容的提问来源于stack exchange,提问作者VErYSEYmPhY
相关产品推荐
相关产品推荐

