taskiq + ydb
Plugin for taskiq that adds a new result backend, broker and schedule source based on YDB.
Installation
This project can be installed using pip/poetry/uv (choose your preferred package manager):
pip install taskiq-ydb
Quick start
Basic task processing
- Define your broker with asyncpg:
# broker_example.py
import asyncio
from ydb.aio.driver import DriverConfig
from taskiq_ydb import YdbBroker, YdbResultBackend
driver_config = DriverConfig(
endpoint='grpc://localhost:2136',
database='/local',
)
broker = YdbBroker(
driver_config=driver_config,
).with_result_backend(
YdbResultBackend(driver_config=driver_config),
)
@broker.task('solve_all_problems')
async def best_task_ever() -> None:
"""Solve all problems in the world."""
await asyncio.sleep(2)
print('All problems are solved!')
async def main() -> None:
await broker.startup()
task = await best_task_ever.kiq()
print(await task.wait_result())
await broker.shutdown()
if __name__ == '__main__':
asyncio.run(main())
- Start a worker to process tasks (by default taskiq runs two instances of worker):
taskiq worker broker_example:broker
- Run
broker_example.pyfile to send a task to the worker:
python broker_example.py
Your experience with other drivers will be pretty similar. Just change the import statement and that's it.
Task scheduling
- Define your broker and schedule source:
# scheduler_example.py
import asyncio
from taskiq import TaskiqScheduler
from ydb.aio.driver import DriverConfig
from taskiq_ydb import YdbBroker, YdbScheduleSource
driver_config = DriverConfig(
endpoint='grpc://localhost:2136',
database='/local',
)
broker = YdbBroker(driver_config=driver_config)
scheduler = TaskiqScheduler(
broker=broker,
sources=[
YdbScheduleSource(
driver_config=driver_config,
broker=broker,
),
],
)
@broker.task(
task_name='solve_all_problems',
schedule=[
{
'cron': '*/1 * * * *', # type: str, either cron or time should be specified.
'cron_offset': None, # type: str | timedelta | None, can be omitted.
'time': None, # type: datetime | None, either cron or time should be specified.
'args': [], # type list[Any] | None, can be omitted.
'kwargs': {}, # type: dict[str, Any] | None, can be omitted.
'labels': {}, # type: dict[str, Any] | None, can be omitted.
},
],
)
async def best_task_ever() -> None:
"""Solve all problems in the world."""
await asyncio.sleep(2)
print('All problems are solved!')
- Start worker processes:
taskiq worker scheduler_example:broker
- Run scheduler process:
taskiq scheduler scheduler_example:scheduler
Metadata
Release files for taskiq-ydb 0.4.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| taskiq_ydb-0.4.0.tar.gz | 7.8 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| taskiq_ydb-0.4.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 18.0 kB
Release files / taskiq_ydb-0.4.0.tar.gz
| Download URL | taskiq_ydb-0.4.0.tar.gz |
|---|---|
| Size | 7.8 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
354c9a9769f2a6c164d0bab434caac78508b687db5d583523b53131b210e219e
|
|
BLAKE2b-256 checksum How to use checksums |
724382db49934230bb9b27f3b3c8ddbd82f93ebb6b2aae5bb17a2a7e4a6af31b
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
uv/0.9.5
|
Release files / taskiq_ydb-0.4.0-py3-none-any.whl
| Download URL | taskiq_ydb-0.4.0-py3-none-any.whl |
|---|---|
| Size | 10.2 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
3ccf18b377aa96f4241d4c49f9144f642f14757c724181f126e67874c39607e1
|
|
BLAKE2b-256 checksum How to use checksums |
7c2c5e970742b74757309a1173642204041492751a61fd191025b187dc148078
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
uv/0.9.5
|