I am trying to set up a script that does the following:
- Making several different queries from database via http
- Parsing query results in pandas dataframes
- Joining dataframes
- Other intensive computations on merged data
I have looked into Dask and these examples/tutorials:
https://distributed.dask.org/en/stable/asynchronous.html
https://distributed.dask.org/en/stable/client.html
but I am struggling to wrap my head around how to integrate my current asyncio/aiohttp setup with Dask (in order to speed up the joins and following computations).
Code is as follows for now:
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
What I have written so far works but:
- Joining dataframes is CPU bound and very slow. I also need to aggregate the (long) list of dataframes I get and do some calculations.
- I am looking into Dask to parallelize the joins (like in "Custom computation: Tree summation" here) but the examples are not really helpful
- I am very new to asyncio and Dask so I'm not sure I am doing anything right.
Question: how do I adapt the code above to work with dask.distributed.Client(asynchronous=True)?