Overview

Cache invalidation is the hardest problem in distributed systems because the database is the source of truth and the cache is a copy — and copies get stale. The classic approach is TTL: give cached values a time-to-live, accept that reads within that window may be stale, and size the TTL to balance freshness against DB load. For product catalogues or exchange rates, this is fine. For account balances in a regulated bank, it is not.

The TTL failure scenario is elementary: a customer initiates a SAR 10,000 transfer at 14:32:01. The core banking system posts the debit at 14:32:03. The account balance cache entry was written at 14:31:45 with a 60-second TTL. At 14:32:20, the customer opens their mobile banking dashboard. The cache serves the pre-debit balance. The customer sees SAR 42,300 instead of SAR 32,300. SAMA’s Cyber Security Framework requires that balance display reflects the current authorised state — not “the state as of the last cache write.”

Application-level invalidation — where the application deletes the cache key immediately after writing to the database — closes the TTL gap for writes the application makes. It does not cover out-of-band changes: a DBA running an emergency correction, a batch job reconciling SAMA settlement figures, a core banking module updating balances through a database procedure the API layer never calls. Those writes are invisible to the application. CDC is not.

Change Data Capture reads the database’s own transaction log (the DB2 journal on z/OS, the redo log on Oracle, the binlog on MySQL). Every committed change appears in the log within milliseconds of commit. Debezium tails the log, transforms each change into a structured event, and publishes it to Kafka. A dedicated cache invalidation consumer reads these events and deletes or updates the corresponding Redis key. The result: the cache is invalidated whenever the database changes — regardless of which system caused the change.

CDC invalidation makes the cache eventually consistent, not strongly consistent

There is a window between the database write and the Redis DEL completing — typically 100–500 ms. During that window, a cache read returns the stale value. Design your UX accordingly: if a user submits a payment and immediately refreshes their balance, they may see the pre-payment balance for up to half a second. This is acceptable in most designs. If your SLA requires zero staleness, you need synchronous invalidation at write time, not CDC — and you accept the coupling and failure modes that come with it.

CDC-Based Invalidation Architecture

The architecture has four components in sequence. Each has a distinct failure domain, which is why CDC is more resilient than application-level invalidation — the components can fail and recover independently without losing change events.

Component 1: Database transaction log. DB2 for z/OS uses the DB2 journal (also called the log). Every committed INSERT, UPDATE, and DELETE is recorded here before the transaction completes. Debezium reads the journal using DB2’s replication API; it does not poll the table, so it adds negligible load to the database server.

Component 2: Debezium connector. Debezium transforms raw log entries into structured change events — JSON or Avro records containing the before and after values of each changed row, the operation type (c/r/u/d), and the transaction timestamp. The connector stores its read position (offset) in a Kafka topic, so a connector restart resumes from exactly where it stopped.

Component 3: Kafka topic. One topic per source table (or one topic per logical domain if you prefer). Topics are durable: a consumer outage does not lose events; it just falls behind. The invalidation consumer is a dedicated consumer group that reads these topics. It does not share a consumer group with business consumers reading the same topics for other purposes — partition assignment and offset management must be independent.

Component 4: Cache invalidation consumer. A Spring Boot application (or a lightweight Kafka Streams job) that reads change events, extracts the cache key corresponding to the changed row, and calls Redis DEL or UNLINK. The consumer is stateless: it maps event fields to cache keys deterministically. Restart it and it re-reads from its committed offset — already-processed events result in idempotent deletes (deleting a key that is already absent is a no-op in Redis).

Debezium Connector for DB2

IBM DB2 for z/OS requires supplemental logging enabled at the table level before Debezium can read before-images of updated rows. Without before-images, update events carry only the after state — which is sufficient for cache invalidation (you can still derive the cache key from the after row), but not for change-data scenarios that need to know what changed. Enable supplemental logging during the DB2 maintenance window, not as an emergency change on a live table.

debezium-db2-connector.jsonjson
{
  "name": "saib-account-cdc",
  "config": {
    "connector.class": "io.debezium.connector.db2.Db2Connector",
    "database.hostname": "db2-zos-host.saib.internal",
    "database.port": "50000",
    "database.user": "${file:/opt/kafka/secrets:db2.username}",
    "database.password": "${file:/opt/kafka/secrets:db2.password}",
    "database.dbname": "SAIBPROD",
    "database.server.name": "saib",

    // Snapshot mode: initial captures the full table on first start;
    // schema_only skips the snapshot and starts from the current log position.
    // Use schema_only when the cache can tolerate a cold-start period.
    "snapshot.mode": "initial",

    "table.include.list": "SAIB.ACCOUNT_BAL,SAIB.LIMIT_MATRIX,SAIB.CUSTOMER_REF",

    // Column filtering: exclude audit columns that don't affect cache keys.
    // Reduces event payload size significantly for wide tables.
    "column.exclude.list": "SAIB.ACCOUNT_BAL.LAST_MODIFIED_BY,SAIB.ACCOUNT_BAL.RECORD_VERSION",

    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "key.converter.schemas.enable": "false",
    "value.converter.schemas.enable": "false",

    // Offset store: MUST be a durable Kafka topic, not an in-memory store.
    // The connector reads from this offset on restart.
    "offset.storage.topic": "debezium-offsets",
    "offset.flush.interval.ms": "5000",

    // Schema history: required for DB2 connector to track DDL changes.
    "schema.history.internal.kafka.topic": "debezium-schema-history",
    "schema.history.internal.kafka.bootstrap.servers": "kafka-bootstrap:9092",

    // Heartbeat: emit a heartbeat event every 30 s on idle tables.
    // Prevents the offset from falling behind on low-change tables.
    "heartbeat.interval.ms": "30000",

    "topic.prefix": "db",
    "poll.interval.ms": "100",
    "max.batch.size": "2048"
  }
}

The snapshot.mode choice has a significant operational consequence. initial causes Debezium to read every row of the included tables before streaming live changes — for a table with 50 million account balance rows, this is hours of snapshot time during which the connector holds a consistent read, potentially blocking table statistics updates. schema_only starts streaming immediately from the current log position. Use schema_only when you can tolerate a cold-start cache (the cache will be empty and will self-populate on demand for the first few hours). Use initial when the cache must be warm immediately after the connector starts.

Kafka Consumer for Invalidation

The invalidation consumer must be idempotent and must batch Redis operations. At-least-once delivery means the same change event may arrive twice after a consumer restart. Deleting a Redis key that is already absent is a no-op — the idempotency requirement is already satisfied by Redis semantics. Do not build deduplication logic into the consumer; it adds complexity without benefit.

Batch processing matters at scale. A payment-heavy morning session at a large Saudi bank produces tens of thousands of balance-change events per minute. Processing them one Redis DEL per Kafka record wastes round-trip latency. Collect a batch of records, build a DEL call with multiple keys, and send one Redis command per batch.

CacheInvalidationConsumer.javajava
@Component
public class CacheInvalidationConsumer {

    private static final Logger log = LoggerFactory.getLogger(CacheInvalidationConsumer.class);
    private static final int MAX_KEYS_PER_DEL = 500;

    private final RedisTemplate<String, String> redis;
    private final CacheKeyMapper keyMapper;

    public CacheInvalidationConsumer(RedisTemplate<String, String> redis,
                                      CacheKeyMapper keyMapper) {
        this.redis = redis;
        this.keyMapper = keyMapper;
    }

    // @KafkaListener with batch mode — processes up to 500 records per poll
    @KafkaListener(
        topics = { "db.saib.ACCOUNT_BAL", "db.saib.LIMIT_MATRIX", "db.saib.CUSTOMER_REF" },
        groupId = "cache-invalidation",       // dedicated group, not shared with business consumers
        containerFactory = "batchKafkaListenerFactory"
    )
    public void onChangeEvents(List<ConsumerRecord<String, String>> records) {
        List<String> keysToDelete = new ArrayList<>(records.size() * 2);

        for (ConsumerRecord<String, String> record : records) {
            try {
                ChangeEvent event = parseEvent(record.value());
                if (event == null || event.getOp().equals("r")) {
                    continue; // skip snapshot reads — no cache exists yet
                }

                // keyMapper.deriveKeys() is deterministic: same row → same cache keys
                List<String> keys = keyMapper.deriveKeys(record.topic(), event);
                keysToDelete.addAll(keys);
            } catch (Exception e) {
                // Log and continue — never throw from a Kafka listener;
                // the batch offset will still be committed.
                log.error("Failed to process change event on {}: {}",
                          record.topic(), record.offset(), e);
            }
        }

        if (!keysToDelete.isEmpty()) {
            // Partition into batches of MAX_KEYS_PER_DEL to avoid oversized Redis commands
            Lists.partition(keysToDelete, MAX_KEYS_PER_DEL).forEach(batch -> {
                redis.delete(batch);     // issues one DEL key1 key2 ... keyN command
                log.debug("Invalidated {} cache keys", batch.size());
            });
        }
    }

    private ChangeEvent parseEvent(String json) {
        if (json == null) return null; // tombstone — deletion event, key field only
        return objectMapper.readValue(json, ChangeEvent.class);
    }
}
Debezium connector offset store must be persistent

The connector stores its log read position in a Kafka topic (debezium-offsets). If that topic is deleted, recreated with fewer partitions, or compacted aggressively, the connector loses its position. On restart, it will either error or fall back to a new snapshot, depending on configuration. Both outcomes are bad: an error stops invalidation; a new snapshot in initial mode floods your application’s cache with events for rows that haven’t changed. Configure the offset topic with cleanup.policy=compact, retention.ms=-1 (indefinite), and replication factor 3. Treat it as infrastructure, not an application topic.

Redis Invalidation Operations

The choice of Redis command matters for performance and for correctness on large keys.

OperationCommandBehaviourWhen to use
Simple key deleteDEL keySynchronous; blocks Redis thread until key is removedSmall keys, default choice
Async deleteUNLINK keyMarks key as deleted immediately; reclaims memory in background threadLarge cached objects (> 1 MB) to avoid blocking the Redis event loop
Tag-based invalidationSMEMBERS tag + DEL keysReads a set of related cache keys, then deletes all of them atomically via LuaOne DB row backs multiple cache keys (e.g. account balance + account summary + account card list)
Pub/Sub notificationPUBLISH channel eventPushes invalidation notification to subscribed application caches (L1 invalidation)Multi-tier caching where applications keep an in-process cache in front of Redis

Tag-based invalidation is the correct pattern when one database row backs multiple cache keys. An account balance row change should invalidate: the balance key (acct:balance:{accountId}), the account summary key (acct:summary:{accountId}), and potentially the customer dashboard key (cust:dashboard:{customerId}) if it includes the balance. Tracking these relationships in a Redis Set (the tag) keeps the consumer code simple: it knows only the account ID; the tag stores which keys to delete.

tag-invalidation.lualua
-- Atomically read all keys in a tag set and delete them, then delete the tag.
-- KEYS[1] = tag set key, e.g. "tag:acct:12345678"
-- Runs inside a single Redis command slot; no other command interleaves.

local tag_key = KEYS[1]
local members = redis.call('SMEMBERS', tag_key)

if #members == 0 then
    return 0  -- tag is empty or already deleted; no-op
end

-- Build the DEL argument list: tag_key + all member keys
local to_delete = { tag_key }
for _, member in ipairs(members) do
    table.insert(to_delete, member)
end

-- DEL accepts multiple keys; returns the count of deleted keys
return redis.call('DEL', unpack(to_delete))

Load the Lua script once at application startup using Redis SCRIPT LOAD and cache the SHA. Call it via EVALSHA sha 1 tag:acct:{accountId} from the invalidation consumer. The script executes atomically in a single Redis roundtrip regardless of how many keys the tag contains.

Cache-Aside + CDC Hybrid

CDC invalidation does not replace cache-aside; it replaces the TTL in a cache-aside pattern. The read path is unchanged: cache miss → DB read → cache write. The write path changes: instead of the application invalidating the cache after its own write, the database change log triggers the invalidation — regardless of which component made the write.

  1. Read path (unchanged)

    Application reads acct:balance:{accountId} from Redis. Cache hit: return cached value immediately. Cache miss: read from DB2, write result to Redis with a safety-net TTL (see below), return to caller. The safety-net TTL is typically 5–30 minutes — long enough that it only fires if Debezium lag is severe. It is not the primary freshness control; CDC is.

  2. Write path (application writes)

    Application writes the debit to DB2. Transaction commits. The application does not invalidate the cache. Within 100–500 ms, Debezium reads the committed change from the DB2 journal, the invalidation consumer deletes the Redis key, and the next cache read triggers a DB refresh. The coupling between the application write and the cache invalidation is broken — the invalidation always happens, but the application does not drive it.

  3. Out-of-band write path (DBA correction, batch job)

    DBA runs UPDATE ACCOUNT_BAL SET AVAILABLE_BALANCE = ... WHERE ACCOUNT_ID = .... Transaction commits. DB2 journal records the change. Debezium tails the journal. Kafka receives the event. Invalidation consumer deletes the Redis key. Application next read is fresh. The application never knew the DBA update happened — and it does not need to.

Never invalidate and re-populate in the same CDC event handler

A common mistake: the invalidation consumer deletes the Redis key, then immediately reads the DB and writes the new value back to Redis. This creates a read-your-own-write race condition. Between the delete and the re-populate, another thread reads the key (cache miss), goes to the DB, and writes its copy to Redis. Now two writes are racing to populate the same key, and the one that lands second wins — which may be the stalest one if the DB read happens before the transaction that triggered the invalidation is fully visible. Let the application re-populate the cache on the next organic cache miss. The consumer’s only job is to delete.

Invalidation Latency Budget

The total time between a database commit and the Redis key being deleted is the sum of delays in each pipeline stage. Understanding this budget determines whether your architecture meets the SAMA balance-display SLA.

StageTypical latencyTuning lever
DB2 journal flush to disk0–5 msDB2 group commit interval; not tunable for this use case
Debezium polling interval50–100 mspoll.interval.ms — default 100 ms; reduce to 50 ms for latency-sensitive tables
Kafka producer batch timeout0–50 mslinger.ms on the Debezium producer; set 0–10 ms for low-latency
Kafka consumer poll interval0–100 msmax.poll.interval.ms and consumer loop design
Redis DEL round trip1–5 msCo-locate consumer and Redis in the same OpenShift cluster
Total typical100–500 msAchievable without heroics on a well-tuned pipeline
Total under Debezium lag (peak load)500 ms–5 sScale Debezium connector tasks; monitor consumer group lag as a SLI

SAMA does not publish a numeric staleness SLA for balance display. The practical constraint comes from mobile banking UX: if a customer initiates a payment and sees the unchanged balance for more than 2 seconds after the confirmation screen, they file a support ticket. The 100–500 ms window is imperceptible; the 5-second window under peak load is not. Alert when Debezium-to-Redis lag (measured as the difference between the DB2 commit timestamp in the change event and the Redis operation timestamp) exceeds 2 seconds.

Account Balance Cache Pattern

Account balance is the canonical use case for CDC invalidation in retail banking. It is read-heavy (every dashboard load, every payment initiation pre-check), changes frequently (every debit and credit), and has strict accuracy requirements (SAMA requires that the displayed balance be the authorised available balance, not an approximation).

The choice between dual-write and CDC for balance cache invalidation comes down to failure safety. In a dual-write approach, the application writes to DB2 first, then publishes an invalidation event to Kafka. If the application crashes after the DB write but before the Kafka publish, the balance is updated in DB2 but the cache key is never invalidated — it serves stale data until TTL expires. In CDC, the DB2 journal is the source; if the application crashes after the DB write, the commit is in the journal, Debezium will read it on the next poll, and invalidation will happen. CDC is safer under failure precisely because it reads from the database’s own commit record, not from the application’s self-reported write.

The balance cache key should include the account type and currency suffix to avoid key collisions in multi-currency accounts: acct:balance:{accountId}:{currency}. The tag set for an account should include all currency variants: tag:acct:{accountId} → { acct:balance:{accountId}:SAR, acct:balance:{accountId}:USD, acct:summary:{accountId} }. A single Debezium event for the account row invalidates all of them in one Lua script execution.

Pitfalls

Cache key derivation must be deterministic

The invalidation consumer derives the cache key from the Debezium change event fields. The application code derives the cache key from the query parameters. If these two derivation functions are not identical — same algorithm, same encoding, same field selection — the invalidation consumer will delete a key that does not exist while the actual stale key survives. This is a class of bug that is nearly impossible to detect in testing: the system appears to work (the deleted key is gone), but the stale key is never cleaned up. Keep the key derivation logic in a single shared library, tested with a table of known inputs and expected outputs, used by both the consumer and the application.

Debezium snapshot competing with live invalidation during restart

When Debezium restarts with snapshot.mode=initial and the offset is lost (or the connector is reconfigured), it takes a full table snapshot. During the snapshot, it emits read events ("op": "r") for every row — not update events. These should be filtered in the consumer (as shown in the code example above). If you accidentally process snapshot read events as invalidations, you will delete every cache key for every account — causing a thundering herd of DB2 reads as the application re-populates the cache from scratch under load. Filter "op": "r" events before passing them to the invalidation logic.

Redis eviction racing with CDC invalidation

If Redis is configured with maxmemory-policy allkeys-lru and memory pressure causes eviction, Redis may evict a key before the CDC invalidation consumer gets to it — which is harmless (the key is already gone). More dangerous: Redis evicts a key, the application re-populates it from the database, and then the CDC consumer delivers a delayed invalidation event (from Debezium lag) and deletes the freshly populated key. The result is unnecessary cache churn. Use maxmemory-policy volatile-lru with TTLs on all keys so eviction only applies to keys that have a TTL; your most-accessed hot keys can have long TTLs and will survive memory pressure better than allkeys-lru.

Steps to Replace TTL-Based Invalidation with CDC

  1. Enable supplemental logging on target tables

    Work with the DB2 DBA to enable supplemental logging on ACCOUNT_BAL, LIMIT_MATRIX, and other tables that back cached entities. This is a one-time change per table. Verify that before-images appear in change events after enabling. Schedule during a low-traffic window.

  2. Deploy Debezium connector with schema_only snapshot

    Start with snapshot.mode=schema_only so the connector begins streaming from the current log position without a full table scan. The cache will be cold for these tables until organic reads re-populate it. This avoids the snapshot thundering herd during the initial rollout.

  3. Build the CacheKeyMapper with shared library

    Extract the cache key derivation logic from the application into a shared library. Implement it in the invalidation consumer using the same library. Write unit tests for every table with a table of (input fields) → (expected cache keys) mappings. Include edge cases: null fields, multi-currency accounts, composite keys.

  4. Deploy the invalidation consumer and validate against shadow TTL

    Run the consumer alongside the existing TTL-based invalidation. Log every invalidation event with the account ID and timestamp. Compare with the DB2 audit log for balance updates. Verify that every out-of-band DB update triggers an invalidation within your SLA window. Run this validation for at least one full business day, including the end-of-day batch window.

  5. Shorten the safety-net TTL

    Once CDC is validated, shorten the cache TTL from your current value (e.g. 60 seconds) to the safety-net value (e.g. 5 minutes). The TTL no longer drives freshness — it is a backstop for Debezium outages. Document the acceptable Debezium lag SLA and the safety-net TTL value in the runbook.

  6. Remove application-level cache invalidation calls

    Remove the explicit cache.evict() calls from application write paths. These are now redundant — CDC will invalidate within 500 ms of any write. Removing them simplifies the application code and eliminates the dual-write race condition.

  7. Alert on Debezium lag

    Add a Kafka consumer group lag alert on the cache-invalidation consumer group for each CDC topic. Alert at 1,000 unconsumed records (approximately 5 seconds of peak throughput). Page at 10,000 records (approximately 50 seconds). A lagging invalidation consumer is a cache accuracy incident, not just a performance issue.

Invalidation strategyStaleness windowCovers out-of-band DB writesImplementation complexityFailure safety
TTL-only0 to TTL (up to minutes)Yes (eventually)Minimal — set TTL at write timeHigh — no runtime dependency; degrades gracefully
Application-levelNear-zero for app writes; TTL for out-of-bandNo — blind to DBA changes and batch jobsMedium — each write path must call evict()Low — crash between write and evict() leaves stale key
CDC-driven100–500 ms typicalYes — reads the DB journal, not the applicationHigh — Debezium + Kafka + consumer infrastructureHigh — journal offset survives outages; catch-up on restart
Dual-writeNear-zero for app writes; uncovered for out-of-bandNo — only as good as application coverageMedium-high — write to cache and DB in same operationLow — partial write on crash leaves cache inconsistent with DB

Production Checklist

  1. Supplemental logging enabled on all target DB2 tables; verified before-images in change events.
  2. Debezium offset topic: cleanup.policy=compact, retention.ms=-1, replication-factor=3.
  3. Schema history topic: same retention settings as offset topic.
  4. CacheKeyMapper in a shared library, unit-tested with known-input/known-output table per entity.
  5. Snapshot read events ("op": "r") filtered before invalidation logic; tested with a connector restart.
  6. Tag-based invalidation Lua script loaded via SCRIPT LOAD; SHA cached at consumer startup.
  7. Consumer group: cache-invalidation, separate from all business consumer groups.
  8. Redis maxmemory-policy volatile-lru; all cache keys written with TTL (safety-net).
  9. Safety-net TTL documented in runbook; reduced to 5 minutes after CDC validation.
  10. Consumer group lag alert: warn at 1,000 records, page at 10,000 records per CDC topic.
  11. Debezium-to-Redis latency metric (commit timestamp delta); alert at >2 s staleness window.
  12. Application-level cache evict() calls removed after CDC validation.
  13. Debezium connector health dashboard; alert on connector status != RUNNING.