StreamManager
API reference for the real-time event streaming engine using Redis PubSub, tunable queue limits, and bounded memory queues.
StreamManager coordinates active listener queues and Redis PubSub operations. It facilitates cluster-wide event routing via subscription channels formatted as stream:user:<user_id>, falling back to local memory queues if Redis is unconfigured or offline.
Class Definition
from zcore.web import StreamManager
class StreamManager:
def __init__(self) -> None:
...Global Setup Function
init_stream_redis
Initializes the shared Redis connection client for the streaming subsystem.
from typing import Any
from zcore.web import init_stream_redis
def init_stream_redis(client: Any) -> None: ...Prop
Type
Properties
Prop
Type
Methods
subscribe
Subscribes a user, returning a bounded async listener queue. Queue capacity is dynamically resolved from settings.STREAM_QUEUE_MAXSIZE (Default: 100). Initializes background Redis PubSub listeners if this is the first active subscription for a user on this node.
async def subscribe(self, user_id: Any) -> asyncio.Queue[Any]: ...Prop
Type
unsubscribe
Unsubscribes a specific listener queue for a user. Shuts down background Redis tasks if no active listeners remain in the local registry.
async def unsubscribe(self, user_id: Any, queue: asyncio.Queue[Any]) -> None: ...Prop
Type
subscription
Asynchronous context manager safely wrapping active user event streams. Guarantees automatic cleanup and queue unregistration upon block exit or connection drops.
@asynccontextmanager
async def subscription(
self,
user_id: Any
) -> AsyncGenerator[asyncio.Queue[Any], None]: ...Prop
Type
publish
Publishes an event payload to a target user's stream. Broadcasts cluster-wide across Redis PubSub channel stream:user:{user_id}, falling back to local memory delivery if unconfigured.
async def publish(self, user_id: Any, data: dict[str, Any]) -> None: ...Prop
Type
start_listening
Explicitly initiates background Redis PubSub pattern-matching subscribers (stream:user:*).
async def start_listening(self) -> None: ...Usage Example (Server-Sent Events)
from fastapi import APIRouter
from fastapi.responses import StreamingResponse
from zcore.web import StreamManager
router = APIRouter()
stream_manager = StreamManager()
@router.get("/notifications/stream/{user_id}")
async def stream_notifications(user_id: str):
async def event_generator():
# Automatically registers and safely tears down queue on client disconnect
async with stream_manager.subscription(user_id) as queue:
while True:
payload = await queue.get()
yield f"data: {payload}\n\n"
return StreamingResponse(event_generator(), media_type="text/event-stream")Slow-Consumer Overflow Protection:
If a client consumes messages too slowly and their bounded queue reaches maximum capacity (settings.STREAM_QUEUE_MAXSIZE, default: 100 items), StreamManager catches asyncio.QueueFull, immediately unregisters the stalled queue, and drops it to protect host node memory from buffer bloat.