Skip to content

Consistency Layer ​

Every getConfig call to Redis involves network I/O, serialization, and round-trip latency. At scale with thousands of service instances polling for configuration, this becomes a critical bottleneck -- the raw RedisConfigService achieves approximately 241,787 ops/s for reads. CoSky's Consistency Layer solves this by combining a local in-process cache (ConcurrentHashMap) with Redis PubSub-based invalidation. When a configuration changes, the Lua script that writes the new value also publishes a notification on the config's Redis key channel. All application instances subscribed to that channel receive the invalidation event and refresh their local cache. The result: getConfig jumps to 256,733,987 ops/s -- a roughly 1000x improvement -- while maintaining eventual consistency within milliseconds.

At a Glance ​

ComponentResponsibilityKey FileSource
RedisConsistencyConfigServiceDecorator around ConfigService that adds local caching with PubSub invalidationRedisConsistencyConfigService.ktRedisConsistencyConfigService.kt:33
ConfigEventListenerContainerInterface for subscribing to config change eventsConfigEventListenerContainer.ktConfigEventListenerContainer.kt:22
RedisConfigEventListenerContainerRedis PubSub implementation that listens to config key channelsRedisConfigEventListenerContainer.ktRedisConfigEventListenerContainer.kt:14
RedisConfigServiceStandard Redis-backed config service (the delegate)RedisConfigService.ktRedisConfigService.kt:41
ConfigChangedEventEvent model representing a config change with operation typeConfigChangedEvent.ktConfigChangedEvent.kt:20

The Performance Problem ​

In a typical microservice deployment, configuration is read far more frequently than it is written. A cluster of 100 service instances, each reading its configuration on startup and periodically refreshing it, generates thousands of read requests per second. When every getConfig call crosses the network to Redis, the throughput is limited by:

  • Network round-trip latency (even on localhost)
  • Redis command processing overhead
  • Reactive pipeline scheduling and context switching

The JMH benchmark results for RedisConfigService.getConfig confirm this constraint:

RedisConfigServiceBenchmark.getConfig   thrpt   241,787.679   ops/s

This is excellent for a single Redis-backed service, but it means every configuration read competes for the same Redis connection pool.

Source: jmh-cosky-config.json

The Solution: RedisConsistencyConfigService ​

RedisConsistencyConfigService is a decorator that wraps any ConfigService implementation (by default, RedisConfigService) and adds a two-layer caching strategy:

  1. Local cache -- A ConcurrentHashMap<NamespacedConfigId, Mono<Config>> stores the most recently fetched config as a cached Mono. Subsequent getConfig calls for the same configId return the cached value without hitting Redis.

  2. PubSub invalidation -- When any write operation (setConfig, removeConfig, rollback) modifies a configuration, the underlying Lua script publishes a message on the Redis channel corresponding to that config's key. The RedisConfigEventListenerContainer subscribes to these channels and triggers cache invalidation in the consistency service.

mermaid
flowchart TD
    subgraph app["Application Process"]
        style app fill:#161b22,stroke:#30363d,color:#e6edf3
        A["getConfig(namespace, configId)"]
        B["ConcurrentHashMap Cache"]
        C["Cache Hit?"]
        D["Return cached Mono"]
    end

    subgraph delegate["Redis Layer"]
        style delegate fill:#161b22,stroke:#30363d,color:#e6edf3
        E["RedisConfigService.getConfig()"]
        F["Redis HGETALL"]
        G["Cache result + return"]
    end

    subgraph invalidation["PubSub Invalidation"]
        style invalidation fill:#161b22,stroke:#30363d,color:#e6edf3
        H["Lua script: PUBLISH"]
        I["RedisConfigEventListenerContainer"]
        J["onConfigChanged()"]
        K["Refresh cache from Redis"]
    end

    A --> B --> C
    C -->|Yes| D
    C -->|No| E --> F --> G
    G --> B
    H --> I --> J --> K --> B

Architecture ​

The consistency layer sits between the application code and the raw Redis operations, intercepting reads and subscribing to change notifications.

mermaid
flowchart LR
    subgraph client["Service Instance"]
        style client fill:#161b22,stroke:#30363d,color:#e6edf3
        App["Application Code"]
        RCS["RedisConsistencyConfigService"]
        Cache["Local Cache<br>(ConcurrentHashMap)"]
    end

    subgraph redis["Redis"]
        style redis fill:#161b22,stroke:#30363d,color:#e6edf3
        RDS["RedisConfigService<br>(delegate)"]
        Lua["Lua Scripts<br>config_set.lua"]
        PubSub["PubSub Channels"]
    end

    App -->|"getConfig()"| RCS
    RCS -->|cache miss| RDS
    RDS -->|HGETALL| Redis["Redis Storage"]
    RDS -->|"return Config"| RCS
    RCS -->|"store"| Cache
    Cache -->|"return"| RCS
    RCS -->|"return Config"| App

    Lua -->|"PUBLISH"| PubSub
    PubSub -->|"invalidate"| RCS

Cache Hit Flow ​

When a configuration has been previously loaded, the consistency service returns the cached Mono<Config> directly without touching Redis.

mermaid
sequenceDiagram
    autonumber
    participant App as Application
    participant RCS as RedisConsistencyConfigService
    participant Cache as ConcurrentHashMap
    participant RDS as RedisConfigService

    App->>RCS: getConfig(namespace, configId)
    RCS->>Cache: get(NamespacedConfigId)
    Cache-->>RCS: Mono~Config~ (cached)
    Note over RCS: Replays the cached Mono<br>to all subscribers
    RCS-->>App: Mono~Config~
    Note over App: No Redis call made

The key detail is that the cache stores a Mono<Config> that has been .cache()-d by Reactor. This means the initial subscription triggers the Redis call, but all subsequent subscribers receive the same result without re-executing the reactive pipeline. When computeIfAbsent inserts a new entry, it also subscribes to the PubSub channel for that configId to listen for future changes.

Cache Invalidation Flow ​

When any process modifies a configuration, the Lua script publishes a change event on the Redis PubSub channel. Every instance listening on that channel receives the invalidation and lazily refreshes its cache.

mermaid
sequenceDiagram
    autonumber
    participant Producer as Service Instance A
    participant RDS as RedisConfigService
    participant Redis as Redis (Lua)
    participant PS as Redis PubSub
    participant Listener as RedisConfigEventListenerContainer
    participant RCS as RedisConsistencyConfigService
    participant Cache as ConcurrentHashMap

    Producer->>RDS: setConfig(namespace, configId, newData)
    RDS->>Redis: EVAL config_set.lua
    Note over Redis: 1. Check hash<br>2. Write new data<br>3. Archive old version<br>4. PUBLISH change
    Redis->>PS: PUBLISH {ns}:cfg:{configId} "set"

    PS-->>Listener: Message received
    Listener->>Listener: Parse ConfigChangedEvent
    Listener-->>RCS: onConfigChanged(event)
    RCS->>Cache: Overwrite entry with a new cold<br>Mono cached with 1-minute TTL
    Note over RCS,Cache: No Redis call happens here -- the fetch is deferred<br>until the next getConfig() subscriber arrives
    Note over Cache: Next read re-executes the fetch<br>and returns fresh data

The invalidation path is a single overwrite: onConfigChanged replaces the map entry with delegate.getConfig(...).cache(CONFIG_CACHE_TTL) -- a cold Mono that is never subscribed at invalidation time (hence the @Suppress("ReactiveStreamsUnusedPublisher") in the source). The event type (SET, REMOVE, ROLLBACK) is not inspected; all three follow the identical overwrite path, and a removed config simply yields a cached empty Mono.

Event Listener Container ​

The ConfigEventListenerContainer interface abstracts the subscription to configuration change events. Its Redis implementation, RedisConfigEventListenerContainer, wraps a ReactiveRedisMessageListenerContainer and subscribes to Redis PubSub channels using the config's Redis key as the channel name.

mermaid
classDiagram
    class EventListenerContainer~Topic, Event~ {
        <<interface>>
        +receive(topic: Topic) Flux~Event~
    }

    class ConfigEventListenerContainer {
        <<interface>>
    }

    class RedisEventListenerContainer~Topic, Event~ {
        <<abstract>>
        #delegate: ReactiveRedisMessageListenerContainer
        +receive(topic: Topic) Flux~Event~
        #receiveEvent(topic: Topic) Flux~Event~
    }

    class RedisConfigEventListenerContainer {
        +receiveEvent(topic: NamespacedConfigId) Flux~ConfigChangedEvent~
        -asEvent(message) ConfigChangedEvent
    }

    EventListenerContainer <|.. ConfigEventListenerContainer
    EventListenerContainer <|.. RedisEventListenerContainer
    RedisEventListenerContainer <|-- RedisConfigEventListenerContainer

The channel naming convention uses the same key pattern as the configuration hash itself: {namespace}:cfg:{configId}. This means the PubSub channel is naturally namespaced and config-specific, so each service instance only subscribes to the configurations it actually uses.

Source: RedisConfigEventListenerContainer.kt:18

Cache State Machine ​

The local cache for each configuration entry goes through a well-defined set of states as requests and invalidation events occur.

mermaid
stateDiagram-v2
    [*] --> CachedNoTTL: First getConfig() call<br>entry = getConfig().cache()

    CachedNoTTL --> CachedNoTTL: Subsequent getConfig() (cache hit,<br>never expires on its own)
    CachedNoTTL --> CachedTTL: Any PubSub event (set / remove /<br>rollback) overwrites the entry

    CachedTTL --> CachedTTL: getConfig() within TTL (cache hit)
    CachedTTL --> Refetching: Cache TTL reached (1 minute)
    Refetching --> CachedTTL: Next getConfig() re-executes the fetch<br>and caches for another minute
    CachedTTL --> CachedTTL: Next PubSub event (overwrite)

    CachedNoTTL --> [*]: Listener terminated<br>(entry removed from map)
    CachedTTL --> [*]: Listener terminated<br>(entry removed from map)

The CONFIG_CACHE_TTL is set to Duration.ofMinutes(1), but it applies only to entries written by the PubSub refresh path (onConfigChanged uses .cache(CONFIG_CACHE_TTL)). The initial cache entry, created on the first getConfig() call, uses plain .cache() with no TTL and never expires on its own. If the PubSub listener terminates, its doFinally hook removes the entry from the map, so the next getConfig() re-creates the entry and re-subscribes.

Source: RedisConsistencyConfigService.kt:40

Performance Benchmarks ​

The JMH benchmarks (run on MacBook Pro M1 with local Redis, 50 threads) demonstrate the dramatic performance difference between the standard and consistency-backed config services.

OperationImplementationThroughput (ops/s)ImprovementSource
getConfigRedisConfigService241,787baselinejmh-cosky-config.json
getConfigRedisConsistencyConfigService256,733,987~1062xjmh-cosky-config.json
setConfigRedisConfigService140,461N/A (write)jmh-cosky-config.json
mermaid
xychart-beta
    title "getConfig Throughput Comparison"
    x-axis ["RedisConfigService", "RedisConsistencyConfigService"]
    y-axis "Operations / second" 0 --> 260000000
    bar [241787, 256733987]

The consistency service achieves 256,733,987 ops/s because after the first read, all subsequent getConfig calls resolve from the in-process ConcurrentHashMap without any network I/O. The only Redis traffic is the initial load and the PubSub invalidation messages (which are infrequent since configuration changes are rare relative to reads).

Source: jmh-cosky-config.json

Trade-offs ​

The consistency layer introduces a set of trade-offs that are important to understand:

Eventual Consistency ​

When a configuration is updated, there is a brief window (typically sub-millisecond on a local network, slightly longer across data centers) during which some instances may serve stale data. The propagation path is:

Lua script PUBLISH -> Redis PubSub channel -> subscriber callback -> cache refresh

For most microservice configuration use cases (feature flags, connection pool sizes, timeout values), this sub-second eventual consistency is acceptable. It would not be appropriate for configuration that requires strict linearizable reads.

Memory Overhead ​

Each cached configuration entry occupies memory in the JVM heap for the ConcurrentHashMap entry and the cached Mono<Config>. For services with hundreds of configurations, this is negligible. For services with tens of thousands of configurations, monitor heap usage.

PubSub Reliability ​

Redis PubSub is a fire-and-forget protocol -- if an instance is temporarily disconnected (network blip, GC pause), it will miss invalidation messages. The CONFIG_CACHE_TTL of 1 minute bounds this staleness only for entries that have already been refreshed by a previous invalidation event: those entries are re-fetched from Redis at most 60 seconds after their last refresh. The initial cache entry, however, is created with plain .cache() (no TTL) and never expires on its own -- if the very first invalidation event for a config is missed, that entry keeps serving stale data until a later event arrives, the listener terminates (which evicts the entry), or the process restarts.

Source: RedisConsistencyConfigService.kt:64-76

Delegation Pattern ​

RedisConsistencyConfigService uses Kotlin's by delegate syntax to forward all non-getConfig methods directly to the underlying RedisConfigService. This means write operations (setConfig, removeConfig, rollback) always go to Redis and trigger PubSub events. The consistency optimization only applies to reads.

Source: RedisConsistencyConfigService.kt:37

References ​

Released under the Apache License 2.0.