Millet Porridge

English version of https://corvo.myseu.cn

0%

asyncio Combined with Threads and Processes (Translation)

Note: this is a translation; the original is Combining Coroutines with Threads and Processes

Many existing Python libraries are not yet ready to work with asyncio. They may block, or depend on concurrency features the module doesn’t provide. It is still possible to use these libraries in asyncio-based programs. The method is to use the executors provided by concurrent.futures to run this code in separate threads or processes.

Threads

The run_in_executor method in the event loop takes an executor object, a callable object, and some arguments to pass. It returns a Future object, which can wait for the function to finish and pass back the return value. If we don’t pass an executor object, a ThreadPoolExecutor will be created; the example below explicitly creates the executor to limit the maximum number of concurrent threads.

A ThreadPoolExecutor starts its own threads and then calls each function passed in within a thread. This example shows how to combine run_in_executor() with wait() so the event loop still has the ability to yield while these blocking functions run, and when the functions finish, the event loop will reactivate the caller.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
# asyncio_executor_thread.py
import asyncio
import concurrent.futures
import logging
import sys
import time


def blocks(n):
log = logging.getLogger('blocks({})'.format(n))
log.info('running')
time.sleep(0.1)
log.info('done')
return n ** 2


async def run_blocking_tasks(executor):
log = logging.getLogger('run_blocking_tasks')
log.info('starting')

log.info('creating executor tasks')
loop = asyncio.get_event_loop()
blocking_tasks = [
loop.run_in_executor(executor, blocks, i)
for i in range(6)
]
log.info('waiting for executor tasks')
completed, pending = await asyncio.wait(blocking_tasks)
results = [t.result() for t in completed]
log.info('results: {!r}'.format(results))

log.info('exiting')


if __name__ == '__main__':
# Configure logging to show the name of the thread
# where the log message originates.
logging.basicConfig(
level=logging.INFO,
format='%(threadName)10s %(name)18s: %(message)s',
stream=sys.stderr,
)

# Create a limited thread pool.
executor = concurrent.futures.ThreadPoolExecutor(
max_workers=3,
)

event_loop = asyncio.get_event_loop()
try:
event_loop.run_until_complete(
run_blocking_tasks(executor)
)
finally:
event_loop.close()

asyncio_executor_thread.py uses logging to conveniently show which function or thread generated the record. Because each blocking function calls a separate logger, the output also shows some threads being reused to complete the work.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
$ python asyncio_exector_thread.py
MainThread run_blocking_tasks: starting
MainThread run_blocking_tasks: creating executor tasks
ThreadPoolExecutor-0_0 blocks(0): running
ThreadPoolExecutor-0_1 blocks(1): running
ThreadPoolExecutor-0_2 blocks(2): running
MainThread run_blocking_tasks: waiting for executor tasks
ThreadPoolExecutor-0_0 blocks(0): done
ThreadPoolExecutor-0_0 blocks(3): running
ThreadPoolExecutor-0_2 blocks(2): done
ThreadPoolExecutor-0_2 blocks(4): running
ThreadPoolExecutor-0_1 blocks(1): done
ThreadPoolExecutor-0_1 blocks(5): running
ThreadPoolExecutor-0_0 blocks(3): done
ThreadPoolExecutor-0_2 blocks(4): done
ThreadPoolExecutor-0_1 blocks(5): done
MainThread run_blocking_tasks: results: [4, 9, 0, 16, 25, 1]
MainThread run_blocking_tasks: exiting

Processes

A ProcessPoolExecutor works similarly to the above, except it creates worker processes instead of threads. Using separate processes requires more system resources, but for compute-intensive operations it can fully utilize each CPU core.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
# asyncio_executor_process.py
# changes from asyncio_executor_thread.py

if __name__ == '__main__':
# Configure logging to show the id of the process
# where the log message originates.
logging.basicConfig(
level=logging.INFO,
format='PID %(process)5s %(name)18s: %(message)s',
stream=sys.stderr,
)

# Create a limited process pool.
executor = concurrent.futures.ProcessPoolExecutor(
max_workers=3,
)

event_loop = asyncio.get_event_loop()
try:
event_loop.run_until_complete(
run_blocking_tasks(executor)
)
finally:
event_loop.close()

The only difference in the code concerning processes is creating a different type of executor. This example also changes the log format, printing the process id instead of the thread name, to show that the tasks run in separate processes.

1
2
3
4
5
6
7
8
9
10
11
12
13
$ python asyncio_exector_process.py
PID 23876 run_blocking_tasks: starting
PID 23876 run_blocking_tasks: creating executor tasks
PID 23876 run_blocking_tasks: waiting for executor tasks
PID 23877 blocks(0): running
PID 23878 blocks(1): running
PID 23879 blocks(2): running
PID 23877 blocks(0): done
PID 23878 blocks(1): done
PID 23879 blocks(2): done
PID 23878 blocks(3): running
PID 23877 blocks(4): running
PID 23879 blocks(5): running

My Understanding and Extension

The first part of the article about threads gives the feeling of a thread pool. First, you create some threads — the executor is a kind of pool; here only a maximum is given — and then submit tasks. The difference is that here await is used to glue the thread results to the context of the current event loop thread, avoiding writing callbacks.

Maybe this really is a transitional approach — forcibly extending support to libraries that temporarily don’t support asyncio.

The official documentation‘s introduction of run_in_executor, compared with the original blog above, may be more basic; for normal usage it’s still more suitable to follow the example above.