Simplify nested asyncio operations for string modification function

Viewed 228

I've an async code that looks like this:

  • There's a third-party function that performs some operations on the string and returns a modified string, for the purpose of this question, it's something like non_async_func.

  • I've an async def async_func_single function that wraps around the non_async_func that performs a single operation.

  • Then another async def async_func_batch function that nested-wraps around async_func_single to perform the function for a batch of data.

The code kind of works but I would like to know more about why/how, my questions are

  • Is it necessary to create the async_func_single and have async_func_batch wrap around it?

  • Can I directly just feed in a batch of data in async_func_batch to call non_async_func?

  • I have a per_chunk function that feeds in the data in batches, is there any asyncio operations/functions that can avoid the use of pre-batching the data I want to send to async_func_batch?

import nest_asyncio
nest_asyncio.apply()

import asyncio
from itertools import zip_longest

from loremipsum import get_sentences

def per_chunk(iterable, n=1, fillvalue=None):
  args = [iter(iterable)] * n
  return zip_longest(*args, fillvalue=fillvalue)

def non_async_func(text):
  return text[::-1]

async def async_func_single(text):
  # Perform some string operation.
  return non_async_func(text)

async def async_func_batch(batch):
  tasks = [async_func_single(text) for text in batch]
  return await asyncio.gather(*tasks)

# Create some random inputs
thousand_texts = get_sentences(1000)

# Loop through 20 sentence at a time.
for batch in per_chunk(thousand_texts, n=20):  
  loop = asyncio.get_event_loop()
  results = loop.run_until_complete(async_func_batch(batch))
  for i, o in zip(thousand_texts, results):
    print(i, o)
2 Answers

Note that marking your functions as "async def", rather than "def" doesn't make them automatically asynchronous - you can have "async def" functions that are synchronous. The difference between asynchronous functions and synchronous ones is that asynchronous functions define places (using "await") where it waits on either another asynchronous function or waits on an asynchronous IO operation.

Also note that asyncio is not magic - it is basically a scheduler that schedules asynchronous functions to be run based on whether the function/operation that is being "awaited" has completed. And, as the scheduler and the asynchronous functions all run on a single thread, then at any given moment, only a single asynchronous function can be running.

So, going back to your code, the only thing your "async_func_single" function is doing is calling an synchronous function, therefore, despite being marked as "async def", it is still a synchronous function. And the same logic applies to the "async_func_batch" function - the "async_func_single" tasks passed to "asyncio.gather" are all synchronous, so the "asyncio.gather" is just running each task synchronously (so it is not offering up any benefits over a simple for loop waiting on each task), so again the "async_func_batch" is a synchronous function. Because you are just calling synchronous functions, then asyncio is not offering any benefits to your program.

If you want multiple synchronous functions that all run at the same time, you don't use asynchronous functions. You need to run them in parallel processes/threads:

import sys
import itertools
import concurrent.futures

from loremipsum import get_sentences

executor = concurrent.futures.ProcessPoolExecutor(workers=sys.cpu_count())

def per_chunk(iterable, n=1):
    while True:
        chunk = tuple(itertools.islice(iterable, n))
        if chunk:
            yield chunk
        else:
            break

def non_async_func(text):
    return text[::-1]

def process_batches(batches):
    futures = [
        executor.submit(non_async_func, batch)
        for batch in batches
    ]
    concurrent.futures.wait(futures)    

thousand_texts = get_sentences(1000)
process_batches(per_chunk(thousand_texts, n=20))

If you still want to use an asynchronous function to process the batches, then asyncio provides asynchronous wrappers around the concurrent futures:

async def process_batches(batches):
    event_loop = asyncio.get_running_loop()
    futures = [
        event_loop.run_in_executor(executor, non_async_func, batch)
        for batch in batches
    ]
    await asyncio.wait(futures)

thousand_texts = get_sentences(1000)
asyncio.run(process_batches(per_chunk(thousand_texts, n=20)))

but it gives no advantages unless you have other asynchronous functions that can be run while it is waiting.

I have tried to answer your questions below.

The code kind of works but I would like to know more about why/how, my questions are

  • Is it necessary to create the async_func_single and have async_func_batch wrap around it?

    No, this is absolutely not necessary.

  • Can I directly just feed in a batch of data in async_func_batch to
    call non_async_func?

    You could do something like the example 1 below, where you feed all the data directly.

  • I have a per_chunk function that feeds in the data in batches, is
    there any asyncio operations/functions that can avoid the use of
    pre-batching the data I want to send to async_func_batch?

    It's possible to use Asyncio Queues with a max size and then process data until the queue is empty and fill it up again. Check out example 2.

Example 1

import asyncio
from concurrent.futures import ThreadPoolExecutor
from loremipsum import get_sentences

def non_async_func(text):
    return text[::-1]

async def async_func_batch(batch):
    with ThreadPoolExecutor(max_workers=20) as executor:
        futures = [loop.run_in_executor(executor, non_async_func, text) for text in batch]
    return(await asyncio.gather(*futures))

# Create some random inputs
thousand_texts = get_sentences(1000)

# Loop through 20 sentence at a time.
loop = asyncio.get_event_loop()
results = loop.run_until_complete(async_func_batch(thousand_texts))
for i, o in zip(thousand_texts, results):
    print(i, o)

Example 2 Queues can be infinite in size. If you do not specify maxsize it will Queue up all elements before processing. If you remove maxsize, then you need to move join outside of the for-loop and remove if taskQueue.full():.

from loremipsum import get_sentences
import asyncio

async def async_func(text, taskQueue, resultsQueue):
    await resultsQueue.put(text[::-1]) # add the result to the resultsQueue
    taskQueue.task_done() # Tell the taskQueue that the task is finished
    taskQueue.get_nowait() # Don't wait for it (unblocking)

async def main():
    taskQueue = asyncio.Queue(maxsize=20)
    resultsQueue = asyncio.Queue()
    thousand_texts = get_sentences(1000)
    results = []
    for text in thousand_texts:
        await taskQueue.put(asyncio.create_task(async_func(text, taskQueue, resultsQueue)))
        if taskQueue.full(): # If maxsize is reached
            await taskQueue.join() # Will block until finished
    while not resultsQueue.empty():
        results.append(await resultsQueue.get())
    for i, o in zip(thousand_texts, results):
        print(i, o)

if __name__ == "__main__":
    asyncio.run(main())
Related