Real-Time Web Analytics Using Kafka's Stream Processing
Imagine your React site generates millions of events per day—clicks, views, purchases. You want to see DAU in real time, enrich orders with user data, and detect anomalies. But batch processing on Hadoop can't keep up: latency is hours. We are a team of engineers with 5+ years of experience in Kafka Streams, having delivered over 10 stream processing projects for high-load websites. Our solution is the Streams library, which runs inside your JVM application without separate infrastructure. Unlike Apache Flink or Spark Streaming, you don't need to deploy a cluster—just add a dependency to pom.xml. We guarantee reducing latency to 10–50 milliseconds—up to 99.9% faster than batch solutions. At a load of 50,000 events per second, infrastructure costs are cut by about 50% by eliminating the separate cluster—saving approximately $1,500 per month on average.
Architectural Overview
Kafka Streams reads topics, transforms, aggregates, joins data, and writes results back to Kafka or external systems via Kafka Connect. State is stored locally in RocksDB and replicated to changelog topics—this provides fault tolerance without an external database. Typical tasks: aggregating user events (DAU, funnels), enriching order streams with reference data, fraud detection, and materialized views from event-sourced data. We design the topology for your load—from thousands to 300,000 events per second.
Building the Processing Topology
Basic Topology
StreamsBuilder builder = new StreamsBuilder();
KStream<String, UserEvent> events = builder.stream(
"user-events",
Consumed.with(Serdes.String(), userEventSerde)
);
// Filter + transform
KStream<String, PageView> pageViews = events
.filter((userId, event) -> event.getType().equals("PAGE_VIEW"))
.mapValues(event -> PageView.from(event));
// Branching
Map<String, KStream<String, UserEvent>> branches = events.split(Named.as("branch-"))
.branch((k, v) -> v.getType().equals("PURCHASE"), Branched.as("purchases"))
.branch((k, v) -> v.getType().equals("CLICK"), Branched.as("clicks"))
.defaultBranch(Branched.as("other"));
branches.get("branch-purchases").to("purchase-events");
Aggregations with Window Functions
Task: count page views per user in a sliding 5-minute window:
KTable<Windowed<String>, Long> pageViewCounts = pageViews
.groupByKey(Grouped.with(Serdes.String(), pageViewSerde))
.windowedBy(
SlidingWindows.ofTimeDifferenceAndGrace(
Duration.ofMinutes(5),
Duration.ofSeconds(30) // grace period for late events
)
)
.count(Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("page-view-counts")
.withValueSerde(Serdes.Long())
);
// Publish results
pageViewCounts.toStream()
.map((windowedKey, count) -> KeyValue.pair(
windowedKey.key(),
new PageViewStat(windowedKey.key(), windowedKey.window().start(), count)
))
.to("page-view-stats", Produced.with(Serdes.String(), pageViewStatSerde));
Window choice depends on business logic. Comparison of window types:
| Window Type | Behavior | Use Case |
|---|---|---|
| Tumbling Windows | Fixed non-overlapping intervals | Counting events per minute |
| Hopping Windows | Overlapping intervals with fixed slide | Sliding average over 5 minutes updated every minute |
| Sliding Windows | Window shifts with each event | Real-time stats update on every event |
| Session Windows | Grouping by activity periods | User session analysis |
KTable and Materialized Views
KTable is a changelog stream where each new record with the same key overwrites the previous one. Used for reference data:
KTable<String, UserProfile> userProfiles = builder.table(
"user-profiles",
Materialized.as("user-profiles-store")
);
KStream<String, EnrichedEvent> enriched = events.join(
userProfiles,
(event, profile) -> EnrichedEvent.builder()
.event(event)
.userName(profile.getName())
.userSegment(profile.getSegment())
.build()
);
Application Configuration
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "site-analytics-processor");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.StringSerde.class);
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class);
// Performance
props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 4);
props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 10 * 1024 * 1024L); // 10MB
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000);
// Error handling
props.put(StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG,
LogAndContinueExceptionHandler.class);
// RocksDB state store
props.put(StreamsConfig.STATE_DIR_CONFIG, "/var/lib/kafka-streams");
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
// Graceful shutdown
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
Interactive Queries—Reading State Without Kafka
Allows reading the state store directly
ReadOnlyKeyValueStore<String, Long> store = streams.store(
StoreQueryParameters.fromNameAndType(
"page-view-counts",
QueryableStoreTypes.keyValueStore()
)
);
Long count = store.get(userId);
// For windowed store
ReadOnlyWindowStore<String, Long> windowStore = streams.store(
StoreQueryParameters.fromNameAndType(
"page-view-counts-windowed",
QueryableStoreTypes.windowStore()
)
);
WindowStoreIterator<Long> iterator = windowStore.fetch(
userId,
Instant.now().minus(Duration.ofMinutes(5)),
Instant.now()
);
Comparison to Other Frameworks
| Parameter | Kafka Streams | Flink | Spark Streaming |
|---|---|---|---|
| Infrastructure | Only JVM dependency | Separate cluster | Separate cluster |
| Latency | < 50 ms | < 100 ms | > 1 s |
| Scaling | Threads + instances | TaskManager | Executor |
| State | RocksDB + changelog | RocksDB/Flink State | Spark State |
| Cost (100k/s) | ~$500/month per instance | ~$2000/month | ~$1500/month |
Why Choose Kafka Streams for Web Analytics?
Flink requires cluster deployment (JobManager + TaskManagers), increasing cost and complexity. The Streams library works as a regular library—you run it alongside your API on the same server. For a site with up to 100,000 events per second, Streams handles it without separate infrastructure, saving over $1,500/month. To scale, increase thread count (NUM_STREAM_THREADS) or add instances—state is automatically balanced through Kafka.
Fault Tolerance and Monitoring
We use exactly-once semantics (processing.guarantee=exactly_once_v2) and state stores with changelog topics. On restart, the application restores state from Kafka. We configure a grace period for late events and a Dead Letter Queue for erroneous records. In production, we monitor metrics like process-rate, commit-latency, and rocksdb-block-cache-hit-ratio via JMX or Prometheus.
Our Process and Deliverables
- Analytics: audit topics, data schemas, and business requirements.
- Design: choose topology, window types, serialization (Avro + Schema Registry).
- Implementation: code topology, state stores, Interactive Queries.
- Testing: TopologyTestDriver, integration tests with Embedded Kafka.
- Deployment: Docker image with JMX Exporter, CI/CD pipeline.
What's Included
- Architectural pipeline diagram
- Configured topics and Schema Registry
- Topology code with unit tests
- Deployment and monitoring documentation
- Team training (1-day workshop)
Timeline
Basic topology: 3–4 days. Pipeline with aggregations and Interactive Queries: 6–9 days. Full production solution with monitoring, DLQ, and CI/CD: 2–3 weeks.
Common Pitfalls in Implementation
- Wrong window choice: Sliding is better than Tumbling for real-time.
- Missing grace period—late events get lost.
- Cache too large (CACHE_MAX_BYTES_BUFFERING)—increases commit latency.
- Ignoring serialization: Avro with Schema Registry is mandatory.
- Lack of RocksDB monitoring—cache hit ratio drops when memory overflows.
Ready to optimize your web analytics? According to the Kafka Streams documentation, processing latency can be as low as 10 ms. Contact us for a free consultation on your use case.







