dagio: Asynchronous I/O - with DAGs!
dagio is an embarassingly simple Python package for running directed acyclic
graphs of asynchronous I/O operations. It is built using and to be used with
Python's built-in asyncio
module, and provides a veeeery thin layer of functionality on top of it.
:sweat_smile:
- Git repository: https://github.com/brendanhasz/dagio
- Bug reports: https://github.com/brendanhasz/dagio/issues
Getting Started
Suppose you have a set of potentially long-running I/O tasks (e.g. hit a web service, query a database, read a large file from disk, etc), where some of the tasks depend on other tasks having finished. That is, you've got a directed acyclic graph (DAG) of tasks, where non-interdependent tasks can be run asynchronously.
For example, if you've got a task G which depends on E and F, but E
depends on D, and F depends on both C and D, etc:
A
|
B C
\ /|
D |
/ \|
E F
\ /
G
Coding that up using raw asyncio might look something like this:
import asyncio
class MyDag:
async def task_a(self):
# does task a stuff...
async def task_b(self):
# does task b stuff...
async def task_c(self):
# does task c stuff...
async def task_d(self):
# does task d stuff...
async def task_e(self):
# does task e stuff...
async def task_f(self):
# does task f stuff...
async def task_g(self):
# does task g stuff...
async def run():
obj = MyDag()
task_a = asyncio.create_task(obj.task_a())
task_c = asyncio.create_task(obj.task_c())
await task_a
await obj.task_b()
await task_c
await obj.task_d()
task_e = asyncio.create_task(obj.task_e())
task_f = asyncio.create_task(obj.task_f())
await task_e
await task_f
await obj.task_g()
asyncio.run(run())
Which is... fine, I guess :roll_eyes: But, you have to be careful about what
task you start before what other task, and which tasks can safely be run
asynchronously vs those which can't. And then you have to type out all that
logic and ordering manually! With the confusing asyncio API! So: a lot of
thought has to go into it, especially for complex DAGs.
And thinking is hard! Less thinking! :fist:
With dagio, you just use the depends decorator to specify what methods any
other given method depends on, and it'll figure everything out for you, and run
them in the correct order, asynchronously where possible:
import asyncio
from dagio import depends
class MyDag:
async def task_a(self):
# does task a stuff...
@depends("task_a")
async def task_b(self):
# does task b stuff...
async def task_c(self):
# does task c stuff...
@depends("task_b", "task_c")
async def task_d(self):
# does task d stuff...
@depends("task_d")
async def task_e(self):
# does task e stuff...
@depends("task_c", "task_d")
async def task_f(self):
# does task f stuff...
@depends("task_e", "task_f")
async def task_g(self):
# does task g stuff...
async def run():
obj = MyDag()
await obj.task_g()
asyncio.run(run())
Note that:
- Each task in your DAG has to be a method of the same class
- Task methods must be
asyncmethods - Calling a task method decorated with
dependsruns that task and all its dependencies - Task methods should not take arguments nor return values. You can handle
inter-task communication using object attributes (e.g.
self._task_a_output = ...). If you need a lock, you can set up anasyncio.Lockin your class's__init__.
You can also run a non-async method asynchronously in a thread pool using the run_async decorator:
import asyncio
from dagio import depends, run_async
class MyDag:
@run_async
def task_a(self):
# a sync method which does task a stuff...
@run_async
def task_b(self):
# a sync method which does task b stuff...
@depends("task_a", "task_b")
async def task_c(self):
# does task c stuff...
async def run():
obj = MyDag()
await obj.task_c() #runs a and b concurrently, then c
That's it. That's all this package does.
Installation
pip install dagio
Support
Post bug reports, feature requests, and tutorial requests in GitHub issues.
Contributing
Pull requests are totally welcome! Any contribution would be appreciated, from things as minor as fixing typos to things as major as adding new functionality. :smile:
Why the name, dagio?
It's for making DAGs of IO operations. DAG IO. Technically it's asynchronous
DAG-based I/O, and the name adagio would have been siiiick, but it was
already taken! :sob:
Metadata
Release files for dagio 0.0.2
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| dagio-0.0.2.tar.gz | 5.8 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| dagio-0.0.2-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 11.0 kB
Release files / dagio-0.0.2.tar.gz
| Download URL | dagio-0.0.2.tar.gz |
|---|---|
| Size | 5.8 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
f01076dcf478c5d7523a9a48f5cb4b54b97270d12434ba99562a4d4df84ddde4
|
|
BLAKE2b-256 checksum How to use checksums |
4977506d18203ea0c2e1dc553e44e777f0d55481c6a5d7a1c54d22d88a743ccd
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/3.4.2 importlib_metadata/4.8.1 pkginfo/1.7.1 requests/2.26.0 requests-toolbelt/0.9.1 tqdm/4.62.3 CPython/3.9.7
|
Release files / dagio-0.0.2-py3-none-any.whl
| Download URL | dagio-0.0.2-py3-none-any.whl |
|---|---|
| Size | 5.2 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
a8e5226d44b7e099eb5d3485f8462e63b0d5ba2eafe6cc221de739dfe0216ce5
|
|
BLAKE2b-256 checksum How to use checksums |
08b514a2953f9559f3b3fb82f614ee5d686a89387bf5fcb7dba5c57c0db61067
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/3.4.2 importlib_metadata/4.8.1 pkginfo/1.7.1 requests/2.26.0 requests-toolbelt/0.9.1 tqdm/4.62.3 CPython/3.9.7
|