Dask distributed asynchronous with aiohttp

Viewed 48

I am trying to set up a script that does the following:

  1. Making several different queries from database via http
  2. Parsing query results in pandas dataframes
  3. Joining dataframes
  4. 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)?

0 Answers
Related