Skip to content
Md. Moudud Hassan
rss

11 min readArchitecture

Five Failures in a Minute Tells You Nothing

Most payment alerting starts the same way. Someone writes a consumer, counts failed transactions per merchant, and fires an alert at five failures in a minute.

It works for about a week. Then it starts paging at 2pm every day, because your largest merchant does 50,000 transactions an hour and a 0.4% failure rate is five failures a minute all afternoon — perfectly healthy, and indistinguishable from a crisis. So someone raises the threshold to twenty. And then a small merchant’s bank route breaks completely, fails all eleven of its transactions, and nobody hears about it.

The threshold was never the problem. The signal was.

Counting failures measures traffic

A raw failure count is a function of volume. Big merchant, more failures. Busy hour, more failures. The metric moves with business activity, and health is buried somewhere inside that.

The fix looks obvious: use a rate instead. Failures divided by attempts, per window.

That is closer, and it is still not enough on its own. A route that saw one attempt and failed it has a 100% failure rate. So does a route that saw 4,000 and failed them all. One is noise; the other is an incident. Rate alone will page you constantly about merchants who tried twice at midnight.

So the actual signal needs three parts, together:

  • a rate, not a count,
  • over a bounded window, not all time,
  • with a minimum volume before the rate is allowed to mean anything.

None of those are hard individually. What makes this a stream processing problem rather than a SELECT is that you need all three continuously, per key, across a firehose, with the state surviving restarts.

Pick the key that can act

Here is the decision that matters more than any code in this post: what are you grouping by?

Grouping by merchant tells you a merchant is having a bad time. That is worth knowing, but you usually cannot do anything about it automatically.

Grouping by route — the combination of bank or PSP, payment method, and card scheme — tells you something you can act on immediately. If DBBL / CARD / VISA is failing 80% of attempts while DBBL / MFS is fine, that is not a merchant problem. That is one integration degrading, and the response is mechanical: stop sending traffic down it.

That is the difference between an alert that wakes someone up and an alert that shifts traffic on its own. Design the key around the action you want to take, not the entity you happen to have.

public record PaymentOutcome(
        String transactionId,
        String merchantId,
        String bankCode,
        String paymentMethod,      // CARD, MFS, NETBANKING
        String scheme,             // VISA, MASTERCARD, ...
        Status status,             // AUTHORISED, DECLINED, TIMEOUT, ERROR
        Instant occurredAt) {

    public enum Status { AUTHORISED, DECLINED, TIMEOUT, ERROR }

    /** Declines are the customer's bank saying no. That is not our outage. */
    public boolean isTechnicalFailure() {
        return status == Status.TIMEOUT || status == Status.ERROR;
    }

    public String routeKey() {
        return bankCode + "|" + paymentMethod + "|" + scheme;
    }
}

Note isTechnicalFailure(). A DECLINED result — insufficient funds, expired card, fraud rules — is the system working correctly. Counting declines as failures means your alert tracks consumer behaviour: it will spike at the end of the month when people run out of money, which tells you nothing about your platform. Only timeouts and errors indicate something on your side or the route’s side is broken.

That distinction is unglamorous and it is most of the value. Plenty of alerting is useless purely because it never made it.

Do you actually need Kafka Streams?

Worth asking before adopting it, because Kafka Streams is not free.

A plain consumer with Redis handles counting well: INCR a key, set a TTL, check a threshold. If that is genuinely all you need, do that. It is simpler, easier to hire for, and easier to debug at 3am.

Kafka Streams earns its place when you need things Redis counters do not give you:

  • Windows with proper semantics — a Redis TTL is not a window. It cannot tell you “the minute from 14:03:00 to 14:04:00”, it can only tell you “roughly the last minute, from whenever this key was created.”
  • Event time rather than arrival time. If events arrive late — and in payments they do, because a bank confirmation can lag — you want them counted in the window they happened in, not the one they showed up in.
  • State that survives restarts and rebalances, backed by a changelog topic, without you building the recovery.
  • Replay. Reprocess last Tuesday from the log with a fixed rule and get the same answers.

If you need one of those, use Streams. If you need none of them, Redis is the better engineering decision, and choosing the smaller tool is not a failure of ambition.

The topology

Spring Boot’s Kafka Streams support is a KafkaStreamsConfiguration bean plus @EnableKafkaStreams, and then topologies are ordinary @Bean methods taking a StreamsBuilder.

@Configuration
@EnableKafkaStreams
public class StreamsConfig {

    @Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME)
    public KafkaStreamsConfiguration kStreamsConfig() {
        Map<String, Object> props = new HashMap<>();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "route-health-monitor");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        // Count events in the minute they OCCURRED, not the minute they arrived.
        props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG,
                  PaymentEventTimestampExtractor.class);

        // Alerting drives automated traffic shifting, so duplicates are expensive.
        props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);

        // Restarts are much cheaper with a warm standby.
        props.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 1);

        return new KafkaStreamsConfiguration(props);
    }
}

The aggregate itself is a small record carrying both halves of the ratio. You cannot compute a rate from a count alone, so attempts have to be aggregated alongside failures:

public record RouteWindow(long attempts, long failures) {

    static RouteWindow empty() {
        return new RouteWindow(0, 0);
    }

    RouteWindow add(PaymentOutcome outcome) {
        return new RouteWindow(
            attempts + 1,
            failures + (outcome.isTechnicalFailure() ? 1 : 0));
    }

    double failureRate() {
        return attempts == 0 ? 0.0 : (double) failures / attempts;
    }
}

And the topology:

@Bean
public KStream<String, PaymentOutcome> routeHealth(StreamsBuilder builder) {

    KStream<String, PaymentOutcome> outcomes =
        builder.stream("payment.outcomes.v1",
                       Consumed.with(Serdes.String(), outcomeSerde));

    outcomes
        // Re-key to the thing we can actually act on.
        .selectKey((k, v) -> v.routeKey())
        .groupByKey(Grouped.with(Serdes.String(), outcomeSerde))

        // One-minute tumbling windows, allowing 30s for stragglers.
        .windowedBy(TimeWindows.ofSizeAndGrace(
                Duration.ofMinutes(1), Duration.ofSeconds(30)))

        .aggregate(
            RouteWindow::empty,
            (routeKey, outcome, agg) -> agg.add(outcome),
            Materialized.<String, RouteWindow, WindowStore<Bytes, byte[]>>as("route-health")
                        .withValueSerde(routeWindowSerde))

        // Emit ONE result per window, when it closes — not one per event.
        .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded()))

        .toStream()
        .filter((windowedKey, w) -> w.attempts() >= MIN_ATTEMPTS   // volume guard
                                 && w.failureRate() >= FAILURE_THRESHOLD)
        .map((windowedKey, w) -> KeyValue.pair(
                windowedKey.key(),
                RouteAlert.of(windowedKey, w)))
        .to("payment.route-alerts.v1", Produced.with(Serdes.String(), alertSerde));

    return outcomes;
}

With MIN_ATTEMPTS = 50 and FAILURE_THRESHOLD = 0.25, that reads as: in any given minute, if a route saw at least 50 attempts and a quarter of them failed technically, raise an alert. Both numbers are business decisions, and both should be configurable without a deploy — you will change them.

The two lines that decide whether this works

Most of the topology is mechanical. Two calls carry nearly all the operational weight, and both are easy to leave out.

suppress, or a thousand alerts

Without suppress, a windowed aggregation emits a new result on every input record. Your alert topic receives an update per payment — the same window, re-emitted with a slightly higher count each time. Downstream, that is either a thousand duplicate alerts or a deduplication problem you now own.

Suppressed.untilWindowCloses(...) holds the result until the window plus its grace period has elapsed, then emits exactly one final value per key. That is what makes the output an alert rather than a stream of partial counts.

The trade is latency. A one-minute window with 30 seconds of grace means an alert arrives up to 90 seconds after the trouble started. For automatically shifting traffic away from a degraded bank route that is acceptable. For fraud interception it may not be, and then you want shorter windows, or the Processor API with a punctuator so you can emit early warnings and a final verdict.

Note also that untilWindowCloses requires a strict buffer config. unbounded() is honest about what it does — it buffers in memory and will not silently drop results. If memory is a genuine concern, maxBytes(...).shutDownWhenFull() fails loudly instead of quietly discarding data, which is the right behaviour for anything alerting-related.

Event time, or windows that lie

By default Kafka Streams uses the record’s broker timestamp — roughly when it was produced. In a payment platform that is often not when the payment happened. A bank confirmation callback can arrive seconds or minutes after the transaction it describes.

If you window on arrival time, a burst of late confirmations all land in the current window and look like a spike that never occurred. Meanwhile the minute where the trouble actually happened looks quiet.

public class PaymentEventTimestampExtractor implements TimestampExtractor {

    @Override
    public long extract(ConsumerRecord<Object, Object> record, long partitionTime) {
        if (record.value() instanceof PaymentOutcome outcome && outcome.occurredAt() != null) {
            return outcome.occurredAt().toEpochMilli();
        }
        // Never return a negative timestamp — Streams will throw and kill the task.
        return partitionTime;
    }
}

Falling back to partitionTime rather than -1 matters. A negative return from a TimestampExtractor causes Streams to throw, and one malformed record can take down the whole application.

Event time is also what makes the grace period meaningful: it is how long you are willing to wait for stragglers before declaring the window final. Set it from measured lateness, not intuition — look at the actual distribution of occurredAt versus arrival for a day, and pick a percentile you can live with.

Closing the loop

An alert that only lands in Slack is a half-finished feature. Since we deliberately keyed by route, the consumer can act:

@KafkaListener(topics = "payment.route-alerts.v1", groupId = "route-governor")
public void onRouteAlert(RouteAlert alert) {
    log.warn("Route {} degraded: {}/{} failed ({}%) in window ending {}",
             alert.routeKey(), alert.failures(), alert.attempts(),
             Math.round(alert.failureRate() * 100), alert.windowEnd());

    // Stop sending traffic down a route we know is broken.
    routeGovernor.demote(alert.routeKey(), Duration.ofMinutes(5));

    notifier.paymentsOnCall(alert);
}

demote is where the value is. If an alternate route exists for that method and scheme, shift traffic to it. If none does, at least fail fast rather than making every customer wait for a 30-second timeout on a route you already know is dead.

This pairs naturally with a circuit breaker, and the two see different things. A circuit breaker sees one service instance’s view of its own calls. This sees the whole platform’s view of a route, across every instance, over a defined window. The breaker reacts faster; the stream sees the pattern. Use both — and make the demotion time-boxed so recovery is automatic rather than a manual step someone forgets.

What production will teach you

Some things that are obvious only in hindsight.

State stores are not free. Every windowed aggregation is a RocksDB store on local disk, backed by a changelog topic in Kafka. Lose the local disk — a pod reschedules onto a new node — and the store restores by replaying that changelog. With large state this can take minutes, during which the instance is not processing. NUM_STANDBY_REPLICAS_CONFIG = 1 keeps a warm copy elsewhere and turns most restarts into a fast failover. On Kubernetes, a StatefulSet with a persistent volume avoids most restores entirely.

Rebalances stop the world. Adding or removing an instance triggers a rebalance, and with a stateful topology that can mean moving partitions and their state. Static group membership (group.instance.id) plus a sensible session.timeout.ms prevents routine pod restarts from triggering a full reshuffle.

Exactly-once is not free either. exactly_once_v2 adds transactional overhead and increases end-to-end latency, since output becomes visible only on commit. It is the right default here because a duplicate alert can demote a healthy route. For a dashboard feed it would be needless cost — decide per topology, not per platform.

Monitor the monitor. This application is now on the critical path for knowing whether payments work. Alert on its consumer lag, on its state (REBALANCING for an extended period is a problem), and on the absence of expected output. A stream processor that has silently died looks exactly like a platform with no failures — which is the most dangerous failure mode in the whole design, because everything appears fine.

The short version

  1. Count is the wrong signal. It measures traffic. Use a rate.
  2. A rate needs a minimum-volume guard, or one attempt failing reads as 100%.
  3. Key by what you can act on — the route, not the merchant.
  4. Separate declines from technical failures. A decline is the system working.
  5. suppress or drown — otherwise every event re-emits the window.
  6. Use event time, and set the grace period from measured lateness.
  7. Feed the output back into routing, not just a chat channel.
  8. Monitor the stream processor itself. Silence is ambiguous.

The underlying shift is small but it changes what you build: stop asking did this payment fail? and start asking what does this failure mean, given everything else happening right now? The first is a boolean you already have. The second needs state, a window, and a denominator — and that is the entire reason stream processing exists.