High-performance Python client for LAX NATS JetStream with smart routing
Project description
LAX NATS JetStream Python SDK
A high-performance Python client SDK for LAX NATS JetStream broker with intelligent routing between direct NATS (for ultra-low latency) and gRPC (for advanced features).
Perfect for FastAPI and other async Python applications.
Features
- 🚀 Smart Routing: Automatically routes to direct NATS (<1ms) or gRPC based on requirements
- 🔄 Async/Await: Built for modern async Python applications
- ⚡ FastAPI Ready: First-class FastAPI integration
- 🛡️ Circuit Breaker: Protects against cascading failures
- 📊 Prometheus Metrics: Built-in observability
- 🔁 Retry Logic: Configurable retry with exponential backoff
- 🏊 Connection Pooling: gRPC connection pool for high throughput
- 📨 Subscriptions: Stream processing with consumer groups and acknowledgments (v0.2.0+)
- 🎯 Load Balancing: Consumer groups for distributed message processing
Installation
pip install lax-nats-client
# For FastAPI integration
pip install lax-nats-client[fastapi]
Quick Start
import asyncio
from lax_client import LaxClient, PublishOptions, Tier
async def main():
# Create client
async with LaxClient() as client:
# Publish to memory tier (uses direct NATS, <1ms latency)
await client.publish(
"events.user.login",
{"user_id": "123", "ip": "192.168.1.1"}
)
# Publish to persistent tier (uses gRPC, guaranteed delivery)
await client.publish(
"orders.created",
{"order_id": "456", "amount": 99.99},
options=PublishOptions(
tier=Tier.PERSISTENT,
require_ack=True
)
)
asyncio.run(main())
FastAPI Integration
from fastapi import FastAPI, Depends
from contextlib import asynccontextmanager
from lax_client import LaxClient, ClientOptions
# Global client
lax_client = None
@asynccontextmanager
async def lifespan(app: FastAPI):
global lax_client
# Startup
lax_client = LaxClient(ClientOptions(
nats_urls=["nats://localhost:4222"],
broker_addr="localhost:50051",
))
await lax_client.connect()
yield
# Shutdown
await lax_client.close()
app = FastAPI(lifespan=lifespan)
def get_client() -> LaxClient:
return lax_client
@app.post("/events")
async def publish_event(
data: dict,
client: LaxClient = Depends(get_client)
):
# Automatically uses direct NATS for memory tier
message_id = await client.publish("events.api", data)
return {"message_id": message_id}
@app.post("/orders")
async def create_order(
order: dict,
client: LaxClient = Depends(get_client)
):
# Uses gRPC for persistent tier
message_id = await client.publish(
"orders.created",
order,
options=PublishOptions(tier=Tier.PERSISTENT, require_ack=True)
)
return {"order_id": order["id"], "message_id": message_id}
Configuration
from lax_client import ClientOptions
options = ClientOptions(
# NATS settings
nats_urls=["nats://nats1:4222", "nats://nats2:4222"],
nats_connect_timeout=5.0,
# gRPC settings
broker_addr="broker:50051",
connection_pool_size=20,
# Performance tuning
max_concurrent_publishes=1000,
publish_timeout=5.0,
retry_attempts=3,
# Circuit breaker
circuit_breaker_threshold=10,
circuit_breaker_timeout=30.0,
# Metrics
enable_metrics=True,
metrics_prefix="myapp_lax",
# Logging
log_level="INFO",
)
client = LaxClient(options)
Publishing Options
from lax_client import PublishOptions, Tier
# Memory tier (default) - Ultra fast, no persistence
await client.publish("topic", data)
# With specific tier
await client.publish(
"topic",
data,
options=PublishOptions(tier=Tier.PERSISTENT)
)
# Require acknowledgment (forces gRPC path)
await client.publish(
"topic",
data,
options=PublishOptions(require_ack=True)
)
# With headers and timeout
await client.publish(
"topic",
data,
options=PublishOptions(
tier=Tier.REPLICATED,
headers={"trace-id": "abc123"},
timeout=2.0
)
)
# Batch publish for efficiency
messages = [
{"subject": "events.click", "data": {"button": "submit"}},
{"subject": "events.view", "data": {"page": "/home"}},
]
results = await client.batch_publish(messages)
print(f"Sent {results['success_count']} messages")
Subscriptions (v0.2.0+)
Basic Subscription
from lax_client import SubscribeOptions
# Iterator style
async for msg in client.subscribe("orders.*"):
print(f"Received: {msg.get_data_as_json()}")
await msg.ack() # Acknowledge message
# Callback style
async def handler(msg):
data = msg.get_data_as_json()
print(f"Order: {data}")
await msg.ack()
await client.subscribe("orders.*", callback=handler)
Consumer Groups
Create consumer groups for load-balanced message processing:
from lax_client import ConsumerGroupConfig
# Create/join a consumer group
group = await client.create_consumer_group(
name="order-processors",
stream="PERSISTENT_STREAM",
subject_filter="orders.*",
config=ConsumerGroupConfig(
deliver_policy="new", # Start with new messages
max_deliver=3, # Max redelivery attempts
ack_wait_seconds=30 # ACK timeout
)
)
# Subscribe as part of the group
async for msg in client.subscribe(
"orders.*",
consumer_group="order-processors",
options=SubscribeOptions(
max_inflight=10, # Max unacked messages
manual_ack=True, # Require explicit ACK
prefetch=5 # Prefetch for performance
)
):
try:
# Process the message
await process_order(msg.get_data_as_json())
await msg.ack()
except Exception as e:
# Requeue for retry
await msg.nack(delay=5.0)
Message Object
Received messages provide these methods:
msg.id # Message ID
msg.subject # Subject/topic
msg.data # Raw bytes
msg.headers # Headers dict
msg.sequence # Stream sequence
msg.timestamp # Message timestamp
# Helper methods
msg.get_data_as_string() # Decode as string
msg.get_data_as_json() # Parse as JSON
await msg.ack() # Acknowledge
await msg.nack(delay=5) # Requeue with delay
Smart Routing Logic
The SDK automatically chooses the optimal path:
| Condition | Route | Latency | Use Case |
|---|---|---|---|
tier=memory & require_ack=False |
Direct NATS | <1ms | Events, metrics, logs |
tier=persistent/replicated |
gRPC to broker | 10-50ms | Orders, payments |
require_ack=True |
gRPC to broker | 10-50ms | Critical data |
Performance
Based on benchmarks:
| Method | Throughput | P99 Latency |
|---|---|---|
| Direct NATS (memory) | 100k+ msg/sec | <5ms |
| gRPC (persistent) | 5k msg/sec | 50ms |
Error Handling
from lax_client.exceptions import (
PublishError,
CircuitBreakerOpen,
ConnectionError
)
try:
await client.publish("topic", data)
except CircuitBreakerOpen:
# Too many failures, circuit breaker is open
logger.error("Circuit breaker open, service degraded")
except PublishError as e:
# Publish failed after retries
logger.error(f"Publish failed: {e}")
except ConnectionError as e:
# Connection issues
logger.error(f"Connection error: {e}")
Metrics
When enabled, Prometheus metrics are exposed:
# In your FastAPI app
from prometheus_client import generate_latest
@app.get("/metrics")
async def metrics():
return Response(generate_latest(), media_type="text/plain")
Available metrics:
lax_client_publish_total{tier,status}lax_client_publish_latency_seconds{tier,method}lax_client_direct_nats_totallax_client_grpc_totallax_client_circuit_breaker_state
Examples
See the examples directory for:
Development
# Install development dependencies
pip install -e ".[dev]"
# Run tests
pytest
# Format code
black lax_client/
# Type checking
mypy lax_client/
License
MIT License - see LICENSE file for details.
Project details
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file lax_nats_client-0.2.0.tar.gz.
File metadata
- Download URL: lax_nats_client-0.2.0.tar.gz
- Upload date:
- Size: 29.5 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.11.5
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
2c86884e7271ff8afdb7ee849f12359fa0de73c39946fe12f1ee46614302968d
|
|
| MD5 |
865b1c801835f372d615c09a789e9ff5
|
|
| BLAKE2b-256 |
6724ee1499bcab43d44cf30df34f07cbf610b10038acb8219fb2155634abcab2
|
File details
Details for the file lax_nats_client-0.2.0-py3-none-any.whl.
File metadata
- Download URL: lax_nats_client-0.2.0-py3-none-any.whl
- Upload date:
- Size: 25.4 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.11.5
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
af22ed0e71c318e69c03c54e9f61d95cd6ed12d81ffc8a6bca0f3f4f04116e08
|
|
| MD5 |
202cac0959005aed4844566b733f210b
|
|
| BLAKE2b-256 |
91ae2c4da7c412d35472425d4776b6c400c97a8695bcdff5c4bfb296016c9716
|