NATS integration for taskiq
Project description
Taskiq NATS
Taskiq-nats is a plugin for taskiq that adds NATS broker. This package has support for NATS JetStream.
Installation
To use this project you must have installed core taskiq library:
pip install taskiq taskiq-nats
Usage
Here's a minimal setup example with a broker and one task.
Default NATS broker.
import asyncio
from taskiq_nats import NatsBroker, JetStreamBroker
broker = NatsBroker(
[
"nats://nats1:4222",
"nats://nats2:4222",
],
queue="random_queue_name",
)
@broker.task
async def my_lovely_task():
print("I love taskiq")
async def main():
await broker.startup()
await my_lovely_task.kiq()
await broker.shutdown()
if __name__ == "__main__":
asyncio.run(main())
NATS broker based on JetStream
import asyncio
from taskiq_nats import (
PushBasedJetStreamBroker,
PullBasedJetStreamBroker
)
broker = PushBasedJetStreamBroker(
servers=[
"nats://nats1:4222",
"nats://nats2:4222",
],
queue="awesome_queue_name",
)
# Or you can use pull based variant
broker = PullBasedJetStreamBroker(
servers=[
"nats://nats1:4222",
"nats://nats2:4222",
],
durable="awesome_durable_consumer_name",
)
@broker.task
async def my_lovely_task():
print("I love taskiq")
async def main():
await broker.startup()
await my_lovely_task.kiq()
await broker.shutdown()
if __name__ == "__main__":
asyncio.run(main())
NatsBroker configuration
Here's the constructor parameters:
servers
- a single string or a list of strings with nats nodes addresses.subject
- name of the subect that will be used to exchange tasks betwee workers and clients.queue
- optional name of the queue. By default NatsBroker broadcasts task to all workers, but if you want to handle every task only once, you need to supply this argument.result_backend
- custom result backend.task_id_generator
- custom function to generate task ids.- Every other keyword argument will be sent to
nats.connect
function.
JetStreamBroker configuration
Common
servers
- a single string or a list of strings with nats nodes addresses.subject
- name of the subect that will be used to exchange tasks betwee workers and clients.stream_name
- name of the stream where subjects will be located.queue
- a single string or a list of strings with nats nodes addresses.result_backend
- custom result backend.task_id_generator
- custom function to generate task ids.stream_config
- a config for stream.consumer_config
- a config for consumer.
PushBasedJetStreamBroker
queue
- name of the queue. It's used to share messages between different consumers.
PullBasedJetStreamBroker
durable
- durable name of the consumer. It's used to share messages between different consumers.pull_consume_batch
- maximum number of message that can be fetched each time.pull_consume_timeout
- timeout for messages fetch. If there is no messages, we start fetching messages again.
NATS Result Backend
It's possible to use NATS JetStream to store tasks result.
import asyncio
from taskiq_nats import PullBasedJetStreamBroker
from taskiq_nats.result_backend import NATSObjectStoreResultBackend
result_backend = NATSObjectStoreResultBackend(
servers="localhost",
)
broker = PullBasedJetStreamBroker(
servers="localhost",
).with_result_backend(
result_backend=result_backend,
)
@broker.task
async def awesome_task() -> str:
return "Hello, NATS!"
async def main() -> None:
await broker.startup()
task = await awesome_task.kiq()
res = await task.wait_result()
print(res)
await broker.shutdown()
if __name__ == "__main__":
asyncio.run(main())
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
taskiq_nats-0.5.1.tar.gz
(7.7 kB
view details)
Built Distribution
File details
Details for the file taskiq_nats-0.5.1.tar.gz
.
File metadata
- Download URL: taskiq_nats-0.5.1.tar.gz
- Upload date:
- Size: 7.7 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: poetry/1.8.4 CPython/3.10.12 Linux/6.5.0-1025-azure
File hashes
Algorithm | Hash digest | |
---|---|---|
SHA256 | fe305fef7c959613bda893a0ffe92e0696194de9df1181c9b6e3a573bf5ac339 |
|
MD5 | cade2a9ccc52fcf8319e3327ea699f94 |
|
BLAKE2b-256 | 0f1ce46adc5031c92d2ce13ea13e426f664b258fd6e296c8fc6f1bb1009b698a |
File details
Details for the file taskiq_nats-0.5.1-py3-none-any.whl
.
File metadata
- Download URL: taskiq_nats-0.5.1-py3-none-any.whl
- Upload date:
- Size: 7.6 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: poetry/1.8.4 CPython/3.10.12 Linux/6.5.0-1025-azure
File hashes
Algorithm | Hash digest | |
---|---|---|
SHA256 | 87edcc082efe98435f59b344439b03436d4ef52eef52a6e34fd3e5bfe113e168 |
|
MD5 | 0adb7ffab8452df1788c006e5767e847 |
|
BLAKE2b-256 | c8a77ce378ba7653cd5382d5a8611e4fdb13e9b991d7a2f54bd55c081123290f |