NewYour coding agent can read the release notes before it upgrades.Set up the MCP server →
PyPI · #1848 most downloaded on PyPI
NATS client for Python
Last release 19 days ago
16 Sep 2026
Ships fairly regularly
a new release about every 3 months
Nearly every release is documented
notes for 26 of 27 stable releases
Nothing withdrawn
no release was ever pulled
5 years old
34 releases · first in 2021
One column per quarter.
Client-side subject validation for publish, subscribe and request, with a skip_subject_validation connect option to opt out
skip_subject_validation connect option to opt out (#1016)user and password to be callables for credential refresh on reconnect (#891)StreamSource.consumer and AckPolicy.FLOW_CONTROL for pre-created sourcing consumers (#937)JetStreamManager.reset_consumer and ConsumerInvalidResetError (#941)PubAck fields (#938)BadSubjectError instead of a server error (#1016)Replace deprecated asyncio.iscoroutinefunction
Minor release of the nats-py client.
pip install nats-py
rtt method to Client for measuring round-trip time (#858)client_ip property to Client (#861)consumer_limits support to StreamConfig (#780)first_seq support to StreamConfig (#779)limit_marker_ttl support for KV watchers (#911)updates_only mode to object store watch (#666)StreamInfo (#772)auth_token alongside nkey/JWT in CONNECT (#900)email.parser path in _process_headers with a byte-level parser (#928)msg_ttl to the create and purge key-value operations (#834)$JS.API.STREAM.NAMES call since the stream name is known (#807)connect() to avoid a mutable default argument (#884)asyncio.iscoroutinefunction (#932)uv_build (#813)LICENSE file in the nats-py sdist (#840)KeyWatcher.stop() runs on a full queue (#899)PullSubscription.fetch hang due to an orphan lingering request (#934)close() against a None _io_reader (#839)Full Changelog: https://github.com/nats-io/nats.py/compare/v2.14.0...v2.15.0
This release adds ability to reconnect to a specific server.
This release adds ability to reconnect to a specific server.
Full Changelog: https://github.com/nats-io/nats.py/compare/v2.13.1...v2.14.0
This release adds ability to reconnect to a specific server.
Full Changelog: v2.13.1...v2.14.0
A patch release that fixes broken cluster fields introduced in 2.13.0
A patch release that fixes broken cluster fields introduced in 2.13.0 (#818)
Full Changelog: https://github.com/nats-io/nats.py/compare/v2.13.0...v2.13.1
A patch release that fixes broken cluster fields introduced in 2.13.0 (#818)
Full Changelog: v2.13.0...v2.13.1
Add token callback support #812
Add token callback support #812
The token parameter on connect now accepts a callable that is invoked on
each connection attempt, enabling dynamic token refresh on reconnect:
def get_token():
return fetch_token_from_auth_service()
nc = await nats.connect("nats://localhost:4222", token=get_token)
Add per-message TTL support for KV operations #783
KV create, delete, and purge now accept a msg_ttl parameter
(in seconds). Requires nats-server 2.11+.
kv = await js.create_key_value(bucket="SESSIONS")
await kv.create("sess-123", b"user-data", msg_ttl=3600)
await kv.delete("sess-123", msg_ttl=60)
Add consumer-configured inbox_prefix for JetStream pull_subscribe methods #781
A custom inbox_prefix can be passed to pull_subscribe and
pull_subscribe_bind to control the deliver subject prefix:
sub = await js.pull_subscribe("orders.>", "my-consumer", inbox_prefix=b"_CUSTOM_INBOX.")
msgs = await sub.fetch(10)
Add persist_mode to StreamConfig #773
Add raft_group, leader_since, and traffic_acc to ClusterInfo #766
StreamConfig omitempty fields for nats-server > 2.12 #788Add options to send custom WebSocket headers on connect
Add options to send custom WebSocket headers on connect
custom_headers = {
"Authorization": ["Bearer MySecretToken"],
"X-Client-ID": ["my-client-123"],
"Accept": ["application/json", "text/plain"]
}
nc = await nats.connect(
"ws://localhost:4222",
ws_connection_headers=custom_headers
)
> ⚠️ KV keys validation now happens by default, in order to opt-out from keys validation need to switch to use the validate_keys option:
⚠️ KV keys validation now happens by default, in order to opt-out from keys validation need to switch to use the
validate_keysoption:
Added validate_keys option to KV methods for controlling key validation #706
# Key validation enforced by default
await kv.put("valid-key", b"value")
# Opt out of validation for backwards compatibility
await kv.put("invalid.key.", b"value", validate_keys=False)
coroutine 'Queue.get' was never awaited warning #687RuntimeWarning: coroutine 'Queue.get' was never awaited in test #691subject_transforms #690Added KeysWithFilters method for key filtering in KV bucket (by @somratdutta in https://github.com/nats-io/nats.py/pull/602)
KeysWithFilters method for key filtering in KV bucket (by @somratdutta in https://github.com/nats-io/nats.py/pull/602) # Retrieve keys with filters
filtered_keys = await kv.keys(filters=['hello', 'greet'])
print(f'Filtered Keys: {filtered_keys}')
discard_new_per_subject to StreamConfig (by @caspervonb in https://github.com/nats-io/nats.py/pull/609)config = nats.js.api.StreamConfig(
name=stream_name,
discard=nats.js.api.DiscardPolicy.NEW,
discard_new_per_subject=True,
max_msgs_per_subject=100
)
await js.add_stream(config)
Added support for passing pathlib.Path derived types to user_credentials (by @johnweldon in https://github.com/nats-io/nats.py/pull/623)
Add an is_acked property to nats.aio.msg.Msg (by @charles-dyfis-net in https://github.com/nats-io/nats.py/pull/672)
JetStreamContext.publish by @rijenkii in https://github.com/nats-io/nats.py/pull/605REQUEST_TIMEOUT status code for a batch fetch with no_wait=True (by @diorcety in https://github.com/nats-io/nats.py/pull/618)deliver_subject in implicit subscription creation (by @m3nowak in https://github.com/nats-io/nats.py/pull/615)add_stream method (by @ff137 in https://github.com/nats-io/nats.py/pull/607)Improved server version semver handling (by @robinbowes in https://github.com/nats-io/nats.py/pull/679)
Bugfix release which includes:
Added micro module implementing services ADR-32
micro module implementing services ADR-32 (https://github.com/nats-io/nats.py/pull/566)Thank you to @charbonnierg for the community implementation that served as a kick-off point.
import asyncio
import contextlib
import signal
import nats
import nats.micro
async def echo(req) -> None:
"""Echo the request data back to the client."""
await req.respond(req.data)
async def main():
# Define an event to signal when to quit
quit_event = asyncio.Event()
# Attach signal handler to the event loop
loop = asyncio.get_event_loop()
for sig in (signal.Signals.SIGINT, signal.Signals.SIGTERM):
loop.add_signal_handler(sig, lambda *_: quit_event.set())
# Create an async exit stack
async with contextlib.AsyncExitStack() as stack:
# Connect to NATS
nc = await stack.enter_async_context(await nats.connect())
# Add the service
service = await stack.enter_async_context(
await nats.micro.add_service(nc, name="demo_service", version="0.0.1")
)
group = service.add_group(name="demo")
# Add an endpoint to the service
await group.add_endpoint(
name="echo",
handler=echo,
)
# Wait for the quit event
await quit_event.wait()
if __name__ == "__main__":
asyncio.run(main())
nc = await nats.connect()
js = nc.jetstream()
jsm = nc.jsm()
for i in range(300):
await jsm.add_stream(name=f"stream_{i}")
streams_page_1 = await jsm.streams_info(offset=0)
streams_page_1 = await jsm.streams_info(offset=256)
Added publish_async method to jetstream `python ack_future = await js.publish_async("foo", b'bar') await ack_future `
Added publish_async method to jetstream
ack_future = await js.publish_async("foo", b'bar')
await ack_future
Added the ability to file contents as user_credentials (https://github.com/nats-io/nats.py/pull/546)
from nats.aio.client import RawCredentials
...
await nats.connect(user_credentials=RawCredentials("<creds file contents as string>"))
Added heartbeat option to pull subscribers fetch API
Added heartbeat option to pull subscribers fetch API
await sub.fetch(1, timeout=1, heartbeat=0.1)
It can be useful to help distinguish API timeouts from not receiving messages:
try:
await sub.fetch(100, timeout=1, heartbeat=0.2)
except nats.js.errors.FetchTimeoutError:
# timeout due to not receiving messages
except asyncio.TimeoutError:
# unexpected timeout
Added subject_transform to add_consumer
await js.add_stream(
name="TRANSFORMS",
subjects=["test", "foo"],
subject_transform=nats.js.api.SubjectTransform(
src=">", dest="transformed.>"
),
)
Added subject_transform to sources as well:
transformed_source = nats.js.api.StreamSource(
name="TRANSFORMS",
# The source filters cannot overlap.
subject_transforms=[
nats.js.api.SubjectTransform(
src="transformed.>", dest="fromtest.transformed.>"
),
nats.js.api.SubjectTransform(
src="foo.>", dest="fromtest.foo.>"
),
],
)
await js.add_stream(
name="SOURCING",
sources=[transformed_source],
)
Added backoff option to add_consumer
await js.add_consumer(
"events",
durable_name="a",
max_deliver=3, # has to be greater than length as backoff array
backoff=[1, 2], # defined in seconds
ack_wait=999999, # ignored once using backoff
max_ack_pending=3,
filter_subject="events.>",
)
Added compression to add_consumer
await js.add_stream(
name="COMPRESSION",
subjects=["test", "foo"],
compression="s2",
)
Added metadata to add_stream
await js.add_stream(
name="META",
subjects=["test", "foo"],
metadata={'foo': 'bar'},
)
Added support for multiple filter consumers when using nats-server +v2.10 This is only supported when using the pull_subscribe_bind API:
pull_subscribe_bind API:await jsm.add_stream(name="multi", subjects=["a", "b", "c.>"])
cinfo = await jsm.add_consumer(
"multi",
name="myconsumer",
filter_subjects=["a", "b"],
)
psub = await js.pull_subscribe_bind("multi", "myconsumer")
msgs = await psub.fetch(2)
for msg in msgs:
await msg.ack()
subjects_filter option to js.stream_info() APIstream = await js.add_stream(name="foo", subjects=["foo.>"])
for i in range(0, 5):
await js.publish("foo.%d" % i, b'A')
si = await js.stream_info("foo", subjects_filter=">")
print(si.state.subjects)
# => {'foo.0': 1, 'foo.1': 1, 'foo.2': 1, 'foo.3': 1, 'foo.4': 1}
inactive_threshold cleanup timeout to be 5 minutes.
It can now be customized as well by passing inactive_threshold as argument in seconds:w = await kv.watchall(inactive_threshold=30.0)
pull_subscribe_bind first argument to be called consumer instead of durable
since it also supports ephemeral consumers. This should be backwards compatible.psub = await js.pull_subscribe_bind(consumer="myconsumer", stream="test")
Added support to ephemeral pull consumers
next_msgNotFoundError (#499 )Full Changelog: https://github.com/nats-io/nats.py/compare/v2.5.0...v2.6.0
fix: improve types on kv/object_store updates by @tekumara in https://github.com/nats-io/nats.py/pull/497
Full Changelog: https://github.com/nats-io/nats.py/compare/v2.4.0...v2.5.0
Fixed Python 3.7 compatibility: use uuid4 to gen unique task names by @Lancetnik in https://github.com/nats-io/nats.py/pull/457
connect() by @dsodx in https://github.com/nats-io/nats.py/pull/484Full Changelog: https://github.com/nats-io/nats.py/compare/v2.3.1...v2.4.0
Nothing published for this version
Added object store feature based on initial contribution by @domderen
pip install nats-py
Upload example:
import nats
import asyncio
async def main():
nc = await nats.connect("locahost", name="object.py")
js = nc.jetstream()
obs = await js.create_object_store("nats-py-object-store")
object_name = 'my-file.mp4'
with open(object_name) as f:
await obs.put(object_name, f)
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
Download example:
import nats
import asyncio
async def main():
nc = await nats.connect("localhost", name="object.py")
js = nc.jetstream()
obs = await js.object_store("nats-py-object-store")
files = await obs.list()
for f in files:
print(f.name, "-", f.size)
object_name = 'my-file.mp4'
with open("copy-"+object_name, 'w') as f:
await obs.get(object_name, f)
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
WebSocketTransport to detect the attempt to upgrade the connection to TLS by @allanbank (https://github.com/nats-io/nats.py/pull/443)next_msg and tasks cancellation (https://github.com/nats-io/nats.py/pull/446)republish value into dataclass by @orsinium (https://github.com/nats-io/nats.py/pull/405).travis.yml by @WillCodeCo (https://github.com/nats-io/nats.py/pull/420)Added support sending requests with reconnected client (by @charbonnierg https://github.com/nats-io/nats.py/pull/343)
Full Changelog: https://github.com/nats-io/nats.py/compare/v2.1.7...v2.2.0
Changed nc.request to have more unique inboxes to avoid accidental reuse of response tokens #335
nc.request to have more unique inboxes to avoid accidental reuse of response tokens #335Fixed SlowConsumerError that would appear when using async for msg in sub.messages after reaching default pending bytes limit
SlowConsumerError that would appear when using async for msg in sub.messages after reaching default pending bytes limitAdded pending_msgs_limit and pending_bytes_limit can now be set for push and pull consumers from JetStream. To disable limits -1 can be used instead:
pending_msgs_limit and pending_bytes_limit can now be set for push and pull consumers from JetStream. To disable limits -1 can be used instead:# Push Subscriber
await js.subscribe("push-example", pending_bytes_limit=-1, pending_msgs_limit=-1)
# Pull Subscriber
await js.pull_subscribe("pull-example", "durable", pending_msgs_limit=-1, pending_bytes_limit=-1)
sub.pending_bytes and sub.pending_msgs methods to confirm buffered bytes and messages from a SubscriptionFixed accounting bug when using sub.next_msg which would have caused SlowConsumer errors and dropping messages when reaching default limit
Fixed empty message being returned sometimes when calling sub.next_msg after future was cancelled
Fix to header compatibility across clients
num_replicas field to Consumer config (https://github.com/nats-io/nats.py/pull/326)Fixes to headers parser when handling empty keys
Added more methods to the JetStreamManager
JetStreamManager (https://github.com/nats-io/nats.py/pull/315)subscribe (https://github.com/nats-io/nats.py/pull/302)sub.Fetch to not block until all messages arrive, it now yields when there are some pending and handle new 4XX temporary status sent by the servercluster should be optional in Placement class (https://github.com/nats-io/nats.py/pull/310) by @olgeninext_msg() does not cancel internal _next_msg() task (https://github.com/nats-io/nats.py/pull/301)Added pending_size option on connect to set max internal pending size for flushing commands and a flush_timeout option to control the maximum time to
pending_size option on connect to set max internal pending size for flushing commands and a flush_timeout option to control the maximum time to wait for a flush.For example to set it to 2MB only which is the default:
nats.connect(pending_size=2*1024*1024)
And similar to the nats.go client, it can be disabled if set to be negative:
nats.connect(pending_size=-1)
When pending size is disabled, then any published message during a disconnection will result in a synchronous OutboundBufferLimitError thrown when publishing.
To control the flushing there is a flush_timeout option that can be passed on connect as well.
nats.connect(flush_timeout=10)
By default, there is no flush timeout (it is None) so the client can wait indefinitely during a flush similar to nats.go (eventually the ping interval ought to unblock the flushing). These couple of changes also improve the publishing performance of the client, thanks to @charbonnierg for contributing to this issue.
Changed nc.flush() to now throw FlushTimeoutError instead which is a subtype of TimeoutError that it used be using for backwards compatibility.
Changed EOF disconnections being reported as StaleConnectionError and they are instead a UnexpectedEOF error.
Several areas of the client got deprecated in this release:
Major upgrade to the APIs of the Python3 client. For this release the client has also been renamed to be nats-py from asyncio-nats-client, it can now be installed with:
pip install nats-py
# With NKEYS / JWT support
pip install nats-py[nkeys]
This version of the client is not completely compatible with previous versions of the client and it is designed to be used with Python 3.7.
Overall, the API of the client should resemble more the APIs of the Go client:
import nats
async def main():
nc = await nats.connect("demo.nats.io")
sub = await nc.subscribe("hello")
await nc.publish("hello")
msg = await sub.next_msg()
print(f"Received [{msg.subject}]: {msg.data}")
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
There is support for NATS Headers ⚡
import asyncio
import nats
from nats.errors import TimeoutError
async def main():
nc = await nats.connect("demo.nats.io")
async def help_request(msg):
print(f"Received a message on '{msg.subject} {msg.reply}': {msg.data.decode()}")
print("Headers", msg.header)
await msg.respond(b'OK')
sub = await nc.subscribe("hello", "workers", help_request)
try:
response = await nc.request("help", b'help me', timeout=0.5)
print("Received response: {message}".format(
message=response.data.decode()))
except TimeoutError:
print("Request timed out")
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
It also now includes JetStream support:
import asyncio
import nats
async def main():
nc = await nats.connect("demo.nats.io")
# Create JetStream context.
js = nc.jetstream()
# Persist messages on 'foo' subject.
await js.add_stream(name="sample-stream", subjects=["foo"])
for i in range(0, 10):
ack = await js.publish("foo", f"hello world: {i}".encode())
print(ack)
# Create pull based consumer on 'foo'.
psub = await js.pull_subscribe("foo", "psub")
# Fetch and ack messagess from consumer.
for i in range(0, 10):
msgs = await psub.fetch()
for msg in msgs:
print(msg)
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
As well as JetStream KV support:
import asyncio
import nats
async def main():
nc = await nats.connect()
js = nc.jetstream()
# Create a KV
kv = await js.create_key_value(bucket='MY_KV')
# Set and retrieve a value
await kv.put('hello', b'world')
entry = await kv.get('hello')
print(f'KeyValue.Entry: key={entry.key}, value={entry.value}')
# KeyValue.Entry: key=hello, value=world
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
The following site has been created to host the API of the Python3 client: https://nats-io.github.io/nats.py/ The contents of the doc site can be found in the following branch from this same repo: https://github.com/nats-io/nats.py/tree/docs/source
Changed the return type of subscribe instead of returning a sid.
Changed suffix of most errors to follow PEP-8 style and now use the Error suffix. For example, ErrSlowConsumer is now SlowConsumerError. Old style errors are subclasses of the new ones so exceptions under try...catch blocks would be still caught.
Several areas of the client got deprecated in this release:
is_async parameter for subscribeClient.timed_requestloop parameter to most functionsauto_unsubscribe instead preferring sub.unsubscribeSpecial thanks to @brianshannan @charliestrawn @orsinium @charbonnierg for their contributions in this release!
Nothing published for this version
Nothing published for this version
Nothing published for this version
Nothing published for this version
Several areas of the client got deprecated in this release:
Major upgrade to the APIs of the Python3 client! For this release the client has also been renamed to be nats-py from asyncio-nats-client, it can now be installed with:
pip install nats-py
# With NKEYS / JWT support
pip install nats-py[nkeys]
This version of the client is not completely compatible with previous versions of the client and it is designed to be used with Python 3.7.
Overall, the API of the client should resemble more the APIs of the Go client:
import nats
async def main():
nc = await nats.connect("demo.nats.io")
sub = await nc.subscribe("hello")
await nc.publish("hello")
msg = await sub.next_msg()
print(f"Received [{msg.subject}]: {msg.data}")
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
There is support for NATS Headers ⚡
import asyncio
import nats
from nats.errors import TimeoutError
async def main():
nc = await nats.connect("demo.nats.io")
async def help_request(msg):
print(f"Received a message on '{msg.subject} {msg.reply}': {msg.data.decode()}")
print("Headers", msg.header)
await msg.respond(b'OK')
sub = await nc.subscribe("hello", "workers", help_request)
try:
response = await nc.request("help", b'help me', timeout=0.5)
print("Received response: {message}".format(
message=response.data.decode()))
except TimeoutError:
print("Request timed out")
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
It also now includes JetStream support:
import asyncio
import nats
async def main():
nc = await nats.connect("demo.nats.io")
# Create JetStream context.
js = nc.jetstream()
# Persist messages on 'foo' subject.
await js.add_stream(name="sample-stream", subjects=["foo"])
for i in range(0, 10):
ack = await js.publish("foo", f"hello world: {i}".encode())
print(ack)
# Create pull based consumer on 'foo'.
psub = await js.pull_subscribe("foo", "psub")
# Fetch and ack messagess from consumer.
for i in range(0, 10):
msgs = await psub.fetch()
for msg in msgs:
print(msg)
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
As well as JetStream KV support:
import asyncio
import nats
async def main():
nc = await nats.connect()
js = nc.jetstream()
# Create a KV
kv = await js.create_key_value(bucket='MY_KV')
# Set and retrieve a value
await kv.put('hello', b'world')
entry = await kv.get('hello')
print(f'KeyValue.Entry: key={entry.key}, value={entry.value}')
# KeyValue.Entry: key=hello, value=world
await nc.close()
if __name__ == '__main__':
asyncio.run(main())
The following site has been created to host the API of the Python3 client: https://nats-io.github.io/nats.py/ The contents of the doc site can be found in the following branch from this same repo: https://github.com/nats-io/nats.py/tree/docs/source
Changed the return type of subscribe instead of returning a sid.
Changed suffix of most errors to follow PEP-8 style and now use the Error suffix. For example, ErrSlowConsumer is now SlowConsumerError. Old style errors are subclasses of the new ones so exceptions under try...catch blocks would be still caught.
Several areas of the client got deprecated in this release:
Deprecated is_async parameter for subscribe
Deprecated Client.timed_request
Deprecated passing loop parameter to most functions
Nothing published for this version
Nothing published for this version
Your coding agent can read these notes before it upgrades. Set up the MCP server →