Deep Dive

Streaming subscriptions & the embedded cache

The subscription model end to end — keyed and patch-mode filters, the subscribe-ack barrier, the five event types, and local mirrors that stay in sync.

Why it matters

Polling a cache is an admission of defeat: you ask “has anything changed?” on a timer and pay latency plus wasted round trips for the privilege. Aeron Cache inverts it. Clients subscribe, and the cluster pushes every change as it is committed. On top of that primitive sits a local mirror — the embedded cache — that turns a networked store into something that reads like an in-process map, because for your reads it is one. This is what powers real-time pricing fan-out, feature-flag distribution, and presence stores where the value coming to you beats you going to fetch it.

The subscription model

A subscription is a standing request to receive updates for one or more caches. The cluster service tracks subscribers in a CacheSubscriptionService, and AbstractCacheClusterService keeps three of them side by side:

this.subscriptionService        = new CacheSubscriptionServiceImpl<>(/* full value */);
this.patchSubscriptionService   = new CacheSubscriptionServiceImpl<>(/* patch deltas */);
this.countersSubscriptionService = new CacheSubscriptionServiceImpl<>(/* counters */);

Because subscriptions live inside the deterministic state machine (see RAFT consensus), delivery is a side effect of the same committed log that performs the mutation. When an entry is added, removed, patched, or a cache is cleared or deleted, the handler encodes the event once into the shared egress buffer and the relevant subscription service fans it out. Session lifecycle is handled too — when a client disconnects, onSessionClose drops its subscriptions from all three services, so a dead client leaves no residue.

The five update event types

Every streamed change is one of exactly five event types, defined once in the SBE schema (see Aeron + SBE transport) and mirrored across every transport and client:

ADD_ITEM      an entry was created or its value updated (also carries post-patch values)
REMOVE_ITEM   an entry was removed — by request or by TTL expiry
CLEAR_CACHE   all entries in a cache were cleared
DELETE_CACHE  the cache itself was deleted
PATCH_ITEM    a merge-patch delta (patch-mode subscribers only)

Note that REMOVE_ITEM covers both explicit removals and TTL expiry. When a cluster timer fires, the removal flows to subscribers through the same path as a client-initiated remove:

protected <VT extends Reusable> void handlePostRemoveTimerCacheEntry(I cacheId, K key,
        RemoveCacheEntryResult<I, K> removeCacheEntryResult,
        CacheResponseEncoder<I, K, VT> encoder,
        CacheSubscriptionService<I, K, VT> subscriptionService) {
    var length = encoder.encodeRemoveCacheEntryResult(cacheId, key, removeCacheEntryResult, egressBuffer);
    subscriptionService.handleTimerEntryRemoved(removeCacheEntryResult, egressBuffer, length);
}

So a subscriber sees an expiring key disappear with no special-case code — it is just another REMOVE_ITEM.

Keyed and patch-mode subscriptions

Two filters let a client ask for exactly what it needs, and no more. Both are expressed in the GatewaySubscribe message, which carries a repeating group of cache ids, each with its own mode and optional key:

<sbe:message name="GatewaySubscribe" id="2">
    <field name="sendSnapshot" id="1" type="BooleanType"/>
    <field name="counters"     id="2" type="BooleanType"/>
    <group name="cacheIds" id="10" dimensionType="groupSizeEncoding">
        <field name="mode"    id="12" type="SubscriptionMode"/>  <!-- FULL | PATCH -->
        <data  name="cacheId" id="11" type="varStringEncoding"/>
        <data  name="key"     id="13" type="varStringEncoding"/> <!-- optional: keyed -->
    </group>
    <data name="correlationId" id="3" type="varStringEncoding"/>
</sbe:message>
  • Keyed subscriptions supply a key, so the client only receives updates for the keys it cares about — ideal when a cache holds thousands of entries but a given consumer tracks a handful of instruments or flags.
  • Patch-mode subscriptions (SubscriptionMode.PATCH) receive only PATCH_ITEM deltas instead of full values — the streaming half of the patch-native story in JSON Merge Patch.

The sendSnapshot flag lets a new subscriber ask for the current contents up front, so it starts from a complete picture and then stays current from the live stream.

The subscribe-ack barrier

There is a subtle race in any push system: between sending subscribe and the subscription actually registering, updates can slip through unseen. Aeron Cache closes it with an explicit acknowledgement. The schema defines a GatewaySubscribeAck, emitted once per subscribe request, that confirms the cluster has registered the subscription:

<sbe:message name="GatewaySubscribeAck" id="15"
    description="Confirms a GatewaySubscribe is live: the cluster has registered
                 the subscription, so updates for the requested caches will now be delivered.">

This ack is a barrier. A client can await it before trusting the stream — and the embedded client libraries expose exactly that, so callers can block until a subscription is guaranteed live before reading or acting. No “did I miss the first tick?” ambiguity.

Near cache vs. client-side embedded cache

Two patterns put data close to reads; they differ in where the mirror lives.

  • Near cache (read-ahead), server-side. The http-server-near-javalin module fronts the cluster: it subscribes and locally caches items for all keys (not just ones you have GET’d), so GETs are served from the near cache’s local copy. Creating a cache through the near endpoints provisions a normal Aeron Cache remotely and wires up the subscription and local copy automatically. This accelerates read-heavy HTTP services without changing the client.

  • Embedded cache, client-side. The embedded client libraries keep the mirror in your process. EmbeddedAeronCache (string), EmbeddedObjectCache (JSON, deep-merges patches), and EmbeddedCounterCache (int64) each subscribe to the stream and apply incoming events to a local map automatically. Reads become local map lookups with no network round trip at all.

How the mirror stays in sync

The embedded cache’s correctness reduces to one rule: apply every event to the local map in order. ADD_ITEM writes, REMOVE_ITEM deletes, CLEAR_CACHE empties, DELETE_CACHE drops the cache, and PATCH_ITEM deep-merges (object caches only). Because the cluster emits these events in committed log order and the transport delivers them reliably and in order (see Aeron + SBE transport), the local mirror converges to exactly the cluster’s state. Start from a sendSnapshot, cross the subscribe-ack barrier, then let the ordered event stream carry you.

And it works over any transport: HTTP+WebSocket and bidirectional WebSocket in all four client languages, plus the native Aeron SBE gateway for Java and Rust. The subscription semantics — keyed, patch-mode, ack barrier, five event types — are identical regardless of the wire underneath. Pick the transport your latency budget and language demand; the streaming model does not change (see embedded clients).

Takeaway: subscribe, cross the ack barrier, and apply ordered events to a local map. That is the whole secret to reads with no round trip — and it holds on every transport Aeron Cache speaks.