KubeMQ
Client SDKsPythonHow-to guidesRPC

Cached Query

Use server-side query response caching for repeated KubeMQ queries with the Python SDK.

Overview

Query response caching lets the broker answer repeat requests without re-running your handler — useful when a query is expensive to compute (a database lookup, an aggregation, a downstream call) but the same input is asked for repeatedly in a short window. Only the first request pays the processing cost; every other caller gets the same answer straight from the broker.

Set cache_key and cache_ttl_in_seconds on the QueryMessage. The first query with a given key is a miss: it reaches the responder, and the broker stores the response under that key for the TTL. A subsequent query with the same key is a hit — the broker returns the stored response directly without calling the responder. cache_hit on the response tells you which happened.

Gotchas: the cache is keyed by the string you choose, not by the query body — if the underlying data changes mid-TTL, callers can get a stale answer until it expires. Keys are scoped per channel, so the same key on another channel is a separate entry. Caching only helps when requests genuinely repeat with the same key.

Prerequisites

  • KubeMQ server running on localhost:50000
  • Python SDK installed (pip install kubemq)

Code

cached_query.py
"""Example: Cached query — use server-side caching to avoid redundant processing."""

from __future__ import annotations

import asyncio

from kubemq import (
    AsyncCancellationToken,
    AsyncCQClient,
    QueriesSubscription,
    QueryMessage,
    QueryReceived,
    QueryResponse,
)


async def main() -> None:
    async with AsyncCQClient(
        address="localhost:50000",
        client_id="python-queries-cached-query-client",
    ) as client:
        token = AsyncCancellationToken()

        async def query_handler() -> None:
            async for query in client.subscribe_to_queries(
                subscription=QueriesSubscription(
                    channel="python-queries.cached-query",
                    on_receive_query_callback=lambda e: None,
                    on_error_callback=lambda e: print(f"Error: {e}"),
                ),
                cancellation_token=token,
            ):
                print(f"Responder received query: {query.body.decode('utf-8')}")
                await client.send_response(
                    QueryResponse(
                        query_received=query,
                        is_executed=True,
                        body=b"cached response data",
                    )
                )

        handler_task = asyncio.create_task(query_handler())
        await asyncio.sleep(1)

        # First query — responder processes and response is cached
        response1 = await client.send_query(
            QueryMessage(
                channel="python-queries.cached-query",
                body=b"fetch data",
                timeout_in_seconds=10,
                cache_key="my-cache-key",
                cache_ttl_in_seconds=30,
            )
        )
        print(f"First query  — cache_hit: {response1.cache_hit}, body: {response1.body}")

        # Second query — served from cache, responder NOT called
        response2 = await client.send_query(
            QueryMessage(
                channel="python-queries.cached-query",
                body=b"fetch data",
                timeout_in_seconds=10,
                cache_key="my-cache-key",
                cache_ttl_in_seconds=30,
            )
        )
        print(f"Second query — cache_hit: {response2.cache_hit}, body: {response2.body}")

        token.cancel()
        handler_task.cancel()
        try:
            await handler_task
        except asyncio.CancelledError:
            pass


if __name__ == "__main__":
    asyncio.run(main())

How It Works

cache_key and cache_ttl_in_seconds are set on QueryMessage. On the first send_query, KubeMQ forwards the query to the responder and stores the response keyed by cache_key for 30 seconds. The second query with the same cache_key is served entirely from the broker cache — the responder is not called. Check response.cache_hit to confirm: False on the first call, True on the second. Cache entries expire after cache_ttl_in_seconds and the next query after expiry hits the responder again.

Was this page helpful?

On this page