Why this matters
- VaultCommerce fraud-detection runs as an embedded Kafka Streams app in the payment service JVM.
- No separate cluster — scales with application pods; state in local RocksDB.
- DSL (
stream.groupByKey().aggregate()) covers 80% of VaultCommerce use cases. - Exactly-once v2 processing guarantee with
processing.guarantee=exactly_once_v2.
Topology building blocks
KStream for record-per-event; KTable for changelog (latest per key). vaultcommerce.analytics.clicks.v1 → filter bot traffic → groupBy sessionId → count → vaultcommerce.fraud.signals.v1. Application ID vaultcommerce-fraud-detector defines consumer group and state store namespace.
Key points
- application.id — consumer group + state store prefix; unique per topology
- KStream — unbounded sequence of records
- KTable — changelog stream materialized as key-value table
- Topology — DAG of processors wired to Kafka topics
- exactly_once_v2 — transactional read-process-write guarantee
Deployment model
One stream thread per pod typically; num.stream.threads for parallelism within JVM. VaultCommerce runs 4 pods × 1 thread = 4 tasks (must match partition count). Rebalance on pod crash restores state from changelog.
VaultCommerce rollout checklist
Before promoting changes that touch vaultcommerce.analytics.clicks.v1 fraud topology, run the staging KRaft cluster (Kafka 3.7+, Schema Registry 7.x) through a 10k events/min soak test. Compare producer request latency p99 and consumer lag per group against the pre-deploy baseline. Kafka Streams embeds in your JVM. Document the change in the internal topic registry, attach Grafana screenshots to the change ticket, and keep an engineer on lag dashboards for 30 minutes after production rollout — roll back the service release before altering broker-level settings if lag or under-replicated partitions spike.
// VaultCommerce fraud topology — Kafka Streams 3.7+
StreamsBuilder builder = new StreamsBuilder();
KStream<String, ClickEvent> clicks = builder.stream("vaultcommerce.analytics.clicks.v1");
clicks.filter((k, v) -> v.productId() != null)
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)))
.count()
.toStream()
.mapValues((windowedKey, count) -> new FraudSignal(windowedKey.key(), count))
.to("vaultcommerce.fraud.signals.v1");
KafkaStreams streams = new KafkaStreams(builder.build(), streamsConfig);
streams.start();
Quick recall
Everything you need if you only revisit this box.
- Kafka Streams embeds in your JVM.
- application.id scopes state and consumer group.
- Changelog topics back RocksDB state stores.
Test yourself
Answer these before moving on — recall is what makes it stick.