Overview

Stateless stream processing — filtering, transforming, routing — answers questions about a single event in isolation. It cannot answer questions that span multiple events over time: Has this customer made five transactions in the last ten minutes? Is this session of transactions consistent with the customer’s historical behaviour profile? Has this account received three or more incoming transfers today that individually fall below the SAMA reporting threshold but collectively exceed it?

Answering those questions requires state: information accumulated from past events that is available when the current event is processed. The stream processor must maintain that state reliably, recover it correctly after a failure, and evict it when it is no longer needed to prevent unbounded memory growth. Stateful processing is therefore a more demanding operational commitment than stateless processing — it requires a state backend, a recovery mechanism, and an operational discipline around state growth and schema migration.

The two relevant engines in the Kafka ecosystem are Apache Flink (a dedicated stream processing cluster with rich stateful primitives and exactly-once via Chandy-Lamport checkpointing) and Kafka Streams (an embedded library that runs inside your JVM process with state managed via RocksDB and backed by Kafka changelog topics). The right choice depends on the complexity of the stateful computation, the size of the state, and the team’s operational maturity with distributed systems.

Window Types

Windows are the mechanism by which a stream processor groups unbounded event streams into bounded computations. Picking the wrong window type for a use case is a silent mistake — the computation succeeds, but the result answers a different question than intended.

Window typeTime alignmentMemory per windowPrimary banking use caseLate data handling
TumblingFixed boundaries; non-overlappingOne window per keyDaily transaction totals, hourly SAMA reporting aggregatesSingle window per period; missed events update next window
HoppingFixed size, fixed advance; overlappingMultiple windows per key (size/advance ratio)Rolling 5-minute velocity checks with 1-minute advanceEvent contributes to all overlapping windows
SessionGap-based; no fixed sizeOne window per active session per keyLogin session fraud grouping, ATM usage sessionLate events can merge or extend sessions
GlobalNo boundary — entire stream per keyGrows with all events for key (unbounded)Lifetime totals — use with explicit state TTL onlyAlways included; no late data concept

For SAMA velocity limit controls — maximum N transactions in a rolling window — hopping windows provide the most accurate real-time check because the window advances continuously. For daily reporting aggregations aligned to the business day boundary, tumbling windows are the correct and simplest choice. For fraud patterns based on transaction clustering in time (a burst of activity followed by silence), session windows model the behaviour naturally. For customer lifetime statistics that never expire, use a KeyedProcessFunction with explicit state — not a global window, which has no eviction mechanism.

Flink’s state backend determines where operator state is stored during job execution and how it is persisted to durable storage for fault tolerance.

HashMapStateBackend (formerly HeapStateBackend) stores all state as Java objects on the JVM heap. Access is fast — in-process, no serialisation overhead — but state size is bounded by the TaskManager heap and state is lost on a process restart unless checkpointed. Suitable for small per-key state (counters, flags) where the total state across all active keys fits comfortably in 60–70% of the TaskManager heap. In a payment context with millions of active customer keys, this ceiling is reached quickly.

EmbeddedRocksDBStateBackend stores all state in a RocksDB instance on local disk, serialised as byte arrays. State size is bounded only by available disk, not by JVM heap, enabling orders-of-magnitude larger state per TaskManager. The trade-off is serialisation overhead on every state read and write, and RocksDB’s compaction behaviour under high write load (see pitfalls). For most banking stateful pipelines — velocity ring buffers, customer session state, running aggregations over tens of millions of customers — RocksDB is the correct backend.

State TTL is configured per state descriptor and is independent of the state backend. It specifies how long state is retained after the last write or read (configurable). Expired state is not immediately removed — it is lazily evicted on access or during RocksDB background compaction. For payment velocity state, set TTL to 2× the longest velocity window (e.g. TTL=60 minutes for a 30-minute velocity window) to ensure state is available for late-arriving events within the allowed lateness budget.

Kafka Streams State Stores

Kafka Streams manages state through state stores — key-value stores embedded in the Streams application process. Two implementations are available: an in-memory store (fast, lost on restart, bounded by heap) and a persistent RocksDB-backed store (durable across restarts, larger capacity). For production stateful applications, the RocksDB store with a Kafka changelog topic for fault tolerance is the standard choice.

A state store changelog is a compacted Kafka topic that records every state mutation. On recovery from a crash, the Streams application replays the changelog to restore state — the same Chandy-Lamport-style recovery that makes Kafka Streams stateful operations exactly-once compatible with the Kafka transactional API. Configure standby replicas (num.standby.replicas=1) to keep a warm copy of each state store on another instance; standby replicas reduce recovery time from the time to replay the changelog to the time to catch up on missed events since the last checkpoint.

Interactive queries allow external systems to read from a state store over HTTP without going through Kafka. A Streams application exposes a ReadOnlyKeyValueStore via KafkaStreams.store(); a REST API layer wraps it for external access. This is the standard pattern for exposing running velocity totals or customer session data to an operational dashboard or a fraud investigation tool without creating a separate read path.

VelocityQueryResource.java (Kafka Streams interactive queries)java
@Path("/velocity")
public class VelocityQueryResource {

  private final KafkaStreams streams;
  private static final String STORE_NAME = "velocity-store";

  @GET
  @Path("/{customerId}")
  @Produces(MediaType.APPLICATION_JSON)
  public Response getVelocity(@PathParam("customerId") String customerId) {

    // locate the host that owns the partition for this key
    KeyQueryMetadata metadata = streams.queryMetadataForKey(
        STORE_NAME, customerId, Serdes.String().serializer());

    if (metadata.activeHost().equals(thisHostInfo)) {
      // query local store directly
      ReadOnlyKeyValueStore<String, VelocityRecord> store =
          streams.store(StoreQueryParameters.fromNameAndType(
              STORE_NAME, QueryableStoreTypes.keyValueStore()));
      VelocityRecord record = store.get(customerId);
      return record != null
          ? Response.ok(record).build()
          : Response.status(404).build();
    } else {
      // forward to the instance that owns the partition
      String remoteUrl = String.format("http://%s:%d/velocity/%s",
          metadata.activeHost().host(),
          metadata.activeHost().port(),
          customerId);
      return forwardRequest(remoteUrl); // delegate to HTTP client
    }
  }
}

Exactly-Once and Flink Checkpoints

Flink’s exactly-once processing guarantee is implemented via the Chandy-Lamport distributed snapshot algorithm. The JobManager periodically injects lightweight checkpoint barriers into every input partition. When an operator receives a barrier on all its input channels (aligned checkpointing) or immediately on the first barrier (unaligned checkpointing), it snapshots its state to the configured durable backend (S3/HDFS for RocksDB) and passes the barrier downstream. When all operators have checkpointed, the checkpoint is complete and the JobManager records the offset at which each Kafka source partition was positioned. On recovery from a failure, Flink resets all operators to their last checkpoint state, seeks each Kafka source back to the recorded offset, and replays events from there. The result: every event is processed exactly once in the operator state, and every output is written exactly once (with the Kafka transactional sink).

flink-checkpoint-config.yamlyaml
execution:
  checkpointing:
    interval: 60000          # 60 s — balance between recovery time and overhead
    mode: EXACTLY_ONCE
    min-pause-between-checkpoints: 30000  # prevent checkpoint storms
    timeout: 120000         # abort checkpoint if not complete in 2 min
    max-concurrent-checkpoints: 1
    externalized-checkpoint-retention: RETAIN_ON_CANCELLATION
    unaligned: false        # aligned: required for EXACTLY_ONCE with Kafka sink

  state-backend: rocksdb
  state-backend-incremental: true
  state-checkpoints-dir: s3://acmebank-flink/checkpoints/
  state-savepoints-dir: s3://acmebank-flink/savepoints/

  restart-strategy:
    type: exponential-delay
    initial-backoff: 1 s
    max-backoff: 2 min
    backoff-multiplier: 2.0
    reset-attempts-interval: 10 min

kafka-sink:
  semantic: EXACTLY_ONCE
  transaction-timeout-ms: 180000  # must exceed checkpoint timeout
Exactly-once is a guarantee about state, not about side effects

Flink’s exactly-once guarantee means each event updates operator state at most once. It does not mean that every write to an external system happens exactly once. An external database write, a REST call to a fraud alerting API, or a file write inside a processElement call is outside Flink’s transaction scope — these will be re-executed on recovery from a checkpoint. All external writes must be idempotent (keyed on event ID) or transactional (wrapped in the external system’s own transaction) to achieve end-to-end exactly-once semantics. The Kafka sink with EXACTLY_ONCE semantic handles Kafka outputs; everything else requires application-level idempotency.

Velocity Check Implementation

A SAMA-compliant velocity check — block a customer whose card generates more than N transactions within a rolling time window — requires per-key stateful aggregation with event-time semantics and state TTL. The implementation below uses a Flink KeyedProcessFunction with a ListState ring buffer and state TTL.

VelocityCheckFunction.java (Flink KeyedProcessFunction)java
public class VelocityCheckFunction
    extends KeyedProcessFunction<String, PaymentEvent, VelocityAlert> {

  private static final int  THRESHOLD  = 5;         // 5 transactions
  private static final long WINDOW_MS  = 600_000L;   // in 10 minutes
  private static final long STATE_TTL_MS = 1_200_000L; // 20 min TTL (2× window)

  private ListState<TxnStamp> buffer; // (eventTime, txnId) pairs

  @Override
  public void open(Configuration cfg) {
    StateTtlConfig ttl = StateTtlConfig
        .newBuilder(Time.milliseconds(STATE_TTL_MS))
        .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
        .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
        .build();

    ListStateDescriptor<TxnStamp> desc =
        new ListStateDescriptor<>("velocity-buffer", TxnStamp.class);
    desc.enableTimeToLive(ttl);
    buffer = getRuntimeContext().getListState(desc);
  }

  @Override
  public void processElement(PaymentEvent event,
                              Context ctx,
                              Collector<VelocityAlert> out) throws Exception {
    long now = event.getEventTime();
    long cutoff = now - WINDOW_MS;

    // evict entries outside the window
    List<TxnStamp> fresh = StreamSupport
        .stream(buffer.get().spliterator(), false)
        .filter(t -> t.ts >= cutoff)
        .collect(Collectors.toList());

    fresh.add(new TxnStamp(now, event.getTxnId()));
    buffer.update(fresh);

    if (fresh.size() >= THRESHOLD) {
      out.collect(VelocityAlert.builder()
          .customerId(ctx.getCurrentKey())
          .txnId(event.getTxnId())
          .windowCount(fresh.size())
          .windowMs(WINDOW_MS)
          .triggeringEvent(event)
          .build());
    }

    // schedule state cleanup timer to prevent leaks on inactive cards
    ctx.timerService().registerEventTimeTimer(now + STATE_TTL_MS);
  }

  @Override
  public void onTimer(long ts, OnTimerContext ctx,
                      Collector<VelocityAlert> out) throws Exception {
    // evict state for customers silent longer than STATE_TTL_MS
    long cutoff = ts - STATE_TTL_MS;
    List<TxnStamp> fresh = StreamSupport
        .stream(buffer.get().spliterator(), false)
        .filter(t -> t.ts >= cutoff)
        .collect(Collectors.toList());
    if (fresh.isEmpty()) {
      buffer.clear(); // explicit clear for inactive key
    } else {
      buffer.update(fresh);
    }
  }
}
  1. Define event time and Kafka source

    Configure the Kafka source with a WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(10)). This tells Flink to use the eventTime field embedded in the payment event, not the Kafka broker timestamp, for time-based operations. Bounded out-of-order allows events that arrive up to 10 seconds late relative to the watermark to still fall within their correct window.

  2. Key the stream by customer ID

    Apply keyBy(PaymentEvent::getCustomerId) before the process function. Keying ensures that all events for a given customer are routed to the same Flink task slot and share the same state partition. Cross-key state access is not possible in Flink and not needed for per-customer velocity checks.

  3. Configure RocksDB state backend with incremental checkpoints

    Set EmbeddedRocksDBStateBackend with enableIncrementalCheckpointing(true). For a velocity buffer with tens of millions of customer keys, the full checkpoint size without incremental mode can exceed 100 GB; incremental checkpoints typically reduce this to under 1 GB per checkpoint interval, which is sustainable within the 60-second checkpoint window.

  4. Set state TTL on the ListState descriptor

    Configure TTL at 2× the longest velocity window to ensure state is available for late events within the allowed lateness budget while still being evicted for inactive customers. Without TTL, a blocked or inactive card that never generates new events retains its velocity buffer in state indefinitely, causing unbounded RocksDB growth over months.

  5. Register cleanup timers

    The explicit onTimer handler is a belt-and-suspenders measure alongside TTL. TTL eviction is lazy (happens on next access or compaction); the timer forces active cleanup for truly inactive keys. Register the timer at now + STATE_TTL_MS on every event — the timer service deduplicates multiple registrations for the same timestamp automatically.

  6. Route late events to a side output

    Events that arrive after the watermark has advanced past their event time are considered late. Configure a OutputTag<PaymentEvent> as a late data side output on the process function. Late payment events should be logged and manually reviewed, not silently dropped — they may represent real transactions that the velocity check did not see in the main window.

  7. Validate with shadow mode before enforcing

    Run the velocity check in shadow mode (emit to a velocity-alerts-shadow topic, no enforcement action) for two weeks. Compare against the existing rule-engine output. Tune the threshold and window based on real false-positive rates before switching the authorisation gateway to act on Flink’s output.

Stream Enrichment

Most fraud and risk computations require enriching payment events with reference data that does not arrive in the payment event itself: customer segment, account risk tier, merchant category code, counterparty bank risk rating. This reference data changes slowly (hours to days) but must be available in milliseconds at the moment of scoring.

Flink provides two enrichment patterns. Async I/O calls an external service (Redis, HTTP endpoint) for each event asynchronously, with configurable parallelism and timeout. Suitable for low-volume enrichment with acceptable 2–5 ms latency per lookup. Broadcast State loads a slow-changing reference dataset once per operator and makes it available in-process for every event, with zero lookup latency. Suitable for relatively small datasets (merchant categories, risk tiers, country codes) that fit in memory and are updated infrequently.

For a payment stream enriched with merchant category codes: broadcast the merchant reference stream (from a Kafka topic updated daily) to all task slots via DataStream.broadcast(). In the KeyedBroadcastProcessFunction, the processBroadcastElement method updates the broadcast state when a new merchant record arrives; processElement reads from broadcast state to enrich each payment. The enrichment is always in-process — no network hop, sub-millisecond, consistent across all task slots.

Temporal join (also called a time-enriched join) is the correct pattern for enriching a payment stream with a reference data snapshot that was valid at the payment event time, not the current reference data. Flink’s SQL FOR SYSTEM_TIME AS OF syntax implements this: SELECT p.*, r.riskTier FROM payments p JOIN riskProfile FOR SYSTEM_TIME AS OF p.eventTime AS r ON p.customerId = r.customerId. This is essential for regulatory audit: the enrichment used at scoring time must be reconstructable from the data that was available at that moment, not from current reference data that may have changed.

Late Data and Watermarks

Event-time processing requires a mechanism to advance the notion of time in the stream and to decide when a window is “complete enough” to trigger computation. Watermarks are the mechanism: a watermark at time t asserts that no future events with event time less than t will arrive. Once the watermark passes a window’s end boundary, Flink fires the window and emits the result.

WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(10)) is the correct starting point for payment event streams. It generates a watermark that lags 10 seconds behind the highest observed event timestamp in the stream. This allows events that arrive up to 10 seconds out of order (due to network retransmissions, IPS gateway buffering, or micro-batch ingest delays) to still be included in their correct window.

Allowed lateness extends this further: events that arrive after the watermark has passed the window end but within the allowed lateness period update the window result and retrigger it. Use .allowedLateness(Duration.ofSeconds(30)) for IPS transfer events, which occasionally retransmit. After the allowed lateness period, truly late events are routed to a side output for manual review — never silently dropped on a payment topic.

RocksDB compaction under high write load can cause checkpoint timeouts

RocksDB’s background compaction consumes I/O and CPU at the same time Flink is trying to serialise state for a checkpoint. Under sustained high write throughput (millions of state updates per second), compaction and checkpointing compete for the same I/O bandwidth and the checkpoint timeout (configured at 2 minutes above) can be exceeded. Tune RocksDB compaction settings: set state.backend.rocksdb.writebuffer.size to 64 MB, state.backend.rocksdb.writebuffer.count to 4, and state.backend.rocksdb.compaction.style to LEVEL. Use dedicated SSDs for RocksDB working directories; shared network storage at the TaskManager host compounds the problem.

Pitfalls

Never increase operator parallelism on a running job without a state migration plan

Flink partitions state by key hash modulo the operator parallelism. When you change parallelism (e.g. from 12 to 24 task slots for a velocity check operator), the key-to-partition mapping changes. If you restore from a savepoint taken at parallelism 12 and restart with parallelism 24, Flink redistributes the state — but this requires a planned savepoint restore with the explicit new parallelism set, not an in-place scale-up. An unplanned parallelism change on a running job without a savepoint may corrupt or lose state entirely. Always take a savepoint, verify it, change parallelism in the job configuration, and restore from the savepoint. Test the procedure in a staging environment before production.

Growing state without TTL causes OOM in weeks

A velocity buffer state without TTL grows monotonically. Every unique customer key that ever transacted retains its state indefinitely. A bank with 2 million active cardholders and a 30-minute velocity buffer holding 10 transactions each consumes roughly 1–2 GB of state — manageable. After 12 months the same state grows to include churned customers with zero recent transactions, adding hundreds of millions of stale keys to RocksDB. This pattern has caused production OOM kills and data loss on stateful Flink jobs that worked fine in development but were never stress-tested with real churn rates. Always set TTL. Always verify RocksDB size growth in the first 30 days of production operation.

Keyspace skew drives uneven state and checkpoint latency

If one customer key generates 80% of transaction volume (a high-frequency aggregator, a batch payment sender, or a test account left in production), the Flink task slot that owns that key will be the slowest to checkpoint, highest in memory, and most likely to cause backpressure. Monitor per-task subtask metrics for numRecordsInPerSecond and currentOutputWatermark lag; a skewed key will show as a hot subtask. Apply a salting strategy (append a random suffix to the key, aggregate at a higher level) or pre-aggregate at the source.

Production Checklist

  1. Flink job uses EmbeddedRocksDBStateBackend with incremental checkpoints enabled for all stateful operators.
  2. Checkpoint interval 60 s; timeout 120 s; minimum pause between checkpoints 30 s; max concurrent checkpoints 1.
  3. Kafka sink configured with EXACTLY_ONCE semantic; Kafka transaction timeout exceeds checkpoint timeout by at least 50%.
  4. State TTL set on every state descriptor; TTL is at least 2× the longest aggregation window.
  5. Explicit event-time timer registered in processElement for active state cleanup in onTimer.
  6. WatermarkStrategy configured with bounded out-of-order tolerance matching the upstream source latency profile.
  7. Late events routed to a side output topic, not silently dropped; side output monitored for unexpected volumes.
  8. Parallelism change procedure documented and tested in staging: savepoint → update config → restore from savepoint.
  9. Per-subtask backpressure and checkpoint duration dashboards in Grafana; alert on checkpoint duration > 80% of timeout.
  10. Keyspace skew analysis run before production launch; hot keys identified and salted or pre-aggregated.
  11. RocksDB compaction settings tuned; dedicated SSD configured for TaskManager state directories.
  12. Velocity check validated in shadow mode for a minimum of two weeks before enforcement actions are enabled.