PrepZone Logo
PrepZone

Event Sourcing on a Kafka Log

Rebuild VaultCommerce order state from the event log — snapshots, replay, and temporal queries.

Why this matters

  • VaultCommerce order read model rebuilds from OrderPlaced, PaymentCaptured, OrderShipped events.
  • Snapshots every 100 events avoid replaying 5 years of history on cold start.
  • Temporal queries ('what was inventory on date X?') need full event history or snapshots.
  • Pairs with CQRS — write side emits events, read side projects to Postgres/ES.
Partition log
Segment 0offsets 0–999
Segment 1offsets 1000–1999
Active segmentoffsets 2000+
Consumer offset 1850Points into segment 1
Each partition is an append-only sequence of segments. Consumers track position by offset.

Event store patterns on Kafka

Compacted topic keyed by orderId holds latest snapshot OR full event stream on separate topics. VaultCommerce uses non-compacted event stream + Postgres projection as source of truth for queries. Kafka retention = 90 days; archive to S3 for compliance.

Key points

  • Event sourcing — state derived by applying ordered events
  • Projection — materialized view built by consuming events
  • Snapshot — cached aggregate state at a point in time
  • CQRS — separate write model (commands/events) and read model
  • Replay — reprocess history to rebuild or fix projections

Snapshots and replay

Snapshot service consumes until offset N, writes OrderSnapshot to compacted topic. New service instances load snapshot then replay delta. VaultCommerce replays staging topics in Testcontainers for integration tests.

VaultCommerce rollout checklist

Before promoting changes that touch the VaultCommerce order and payment event backbone, 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. Events are facts; projections are derived state. 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.

Java
// VaultCommerce order projection — rebuild from event stream
public OrderView rebuild(String orderId, long upToOffset) {
    OrderView view = OrderView.empty(orderId);
    try (KafkaConsumer<String, DomainEvent> consumer = createReplayConsumer()) {
        consumer.assign(List.of(new TopicPartition("vaultcommerce.orders.events.v1", 0)));
        consumer.seek(new TopicPartition("vaultcommerce.orders.events.v1", 0), 0);
        while (true) {
            var records = consumer.poll(Duration.ofSeconds(1));
            for (var r : records) {
                if (!r.key().equals(orderId)) continue;
                view = view.apply(r.value());
                if (r.offset() >= upToOffset) return view;
            }
        }
    }
}

Quick recall

Everything you need if you only revisit this box.

  1. Events are facts; projections are derived state.
  2. Snapshots bound replay time on cold start.
  3. Kafka + OLTP projection is VaultCommerce's hybrid.

Test yourself

Answer these before moving on — recall is what makes it stick.