ZCore LogoZCore
Core concepts

Caching & Real-Time Streaming

Deep dive into ZCore's distributed Redis cache with resilient in-memory LRU fallback and the cluster-wide PubSub streaming engine.

ZCore provides unified infrastructure for high-speed distributed data access and real-time event streaming, designed to operate seamlessly with Redis in clustered environments while gracefully degrading to in-memory fallbacks when standalone.


1. The BaseCache Dual-Layer Resilience Architecture

BaseCache acts as a transparent, fault-tolerant caching gateway. When configured with a Redis connection URL (init_cache(redis_url)), all operations are directed to distributed Redis keyspaces.

Yes Success Connection Error / Timeout No / Unconfigured cache.get / cache.set Is Redis Connected? Execute Redis Command Return Data Log Error & Fallback Execute TTLLRUCache In-Memory

Graceful Network Degradation: If a network partition or Redis outage occurs mid-request, ZCore catches the connection error, logs the diagnostic warning, and immediately serves/stores the record in the local, thread-safe in-memory TTLLRUCache. When Redis recovers, operations automatically route back to distributed caching without restarting the application.

Dynamic Framework Boundaries

Both BaseCache and TTLLRUCache dynamically bind to centralized framework configuration settings with sensible fallbacks:

  • Default Lifespan (CACHE_DEFAULT_TTL): Fallback TTL (default: 3600 seconds) applied automatically when ttl is omitted in cache.set().
  • Local In-Memory Capacity (CACHE_LOCAL_MAXSIZE): Maximum key volume (default: 1000 items) retained in the local fallback store before LRU eviction triggers.
  • Garbage Collection Sweep (CACHE_EVICTION_INTERVAL): Interval in seconds (default: 60 seconds) between non-blocking background sweeps purging expired records.

Automatic Pydantic V2 Deserialization

You can supply target_type=TaskResponse to cache.get(). ZCore automatically deserializes the raw JSON string and executes TaskResponse.model_validate(), returning strongly-typed model instances:

# Fetches from Redis/Local cache and validates schema automatically
task: TaskResponse | None = await cache.get("task_123", target_type=TaskResponse)

Background Garbage Collection (weakref.WeakSet)

TTLLRUCache instances are registered in a global weakref.WeakSet. This allows the background asynchronous eviction loop (_start_eviction_loop) to sweep expired keys across all active caches at intervals defined by settings.CACHE_EVICTION_INTERVAL without creating circular memory references or preventing clean de-allocation.


2. Real-Time Streaming Engine (StreamManager)

For real-time features like WebSockets, Server-Sent Events (SSE), and live notifications, ZCore provides the StreamManager.

Cluster-Wide Message Propagation

StreamManager connects to Redis PubSub using pattern subscription (stream:user:*). When any worker node in your cluster publishes an event to a user, all nodes receive the frame and route it to the active local queues for that user.

# Stream Manager Subscription Lifecycle (e.g., in an SSE or WebSocket route)
from fastapi.responses import StreamingResponse
from zcore.web import StreamManager

stream_manager = StreamManager()

@router.get("/tasks/stream/{user_id}")
async def stream_user_events(user_id: str):  # Accepts polymorphic user IDs (UUID, int, or str)
    async def event_generator():
        # Automatically subscribes and handles cleanup upon connection closure
        async with stream_manager.subscription(user_id) as queue:
            while True:
                data = await queue.get()
                yield f"data: {data}\n\n"
                
    return StreamingResponse(event_generator(), media_type="text/event-stream")

Slow-Consumer Overflow Protection

Each connected client listener is allocated a bounded asyncio.Queue whose capacity is resolved dynamically from settings.STREAM_QUEUE_MAXSIZE (default: 100 items):

  • If a publisher pushes events faster than a client can consume them and the queue reaches capacity, StreamManager catches asyncio.QueueFull.
  • It immediately unregisters and drops the overflowing queue, safeguarding node memory against buffer bloat and memory exhaustion.

Standalone Mode: If Redis is uninitialized, StreamManager routes events directly through local memory queues, allowing real-time features to function in single-node development environments out of the box.

On this page