When ShopWave needed to process 2 billion daily events in real time for personalized recommendations and fraud detection, batch processing was no longer an option. Their existing Spark-based pipeline had a 45-minute delay between event ingestion and insight delivery, which meant fraud was detected long after the damage was done and product recommendations were always slightly stale.
We chose Apache Flink over alternatives like Kafka Streams and Spark Structured Streaming for several reasons: true event-time processing with watermark support, exactly-once state semantics, and the ability to handle both streaming and batch workloads. The learning curve was steeper, but the operational benefits at scale were worth it.
The architecture centered on a three-layer design. The ingestion layer used Kafka with schema registry for type-safe event streams. The processing layer ran Flink jobs for real-time aggregations, pattern detection, and feature computation. The serving layer used a combination of Redis for hot data and ClickHouse for analytical queries, giving both sub-millisecond and complex analytical access patterns.
Fraud detection was the most technically challenging component. We implemented a complex event processing pattern in Flink that correlated user behavior across multiple event types (page views, cart additions, checkout attempts, payment submissions) within sliding time windows. The system learned normal behavior patterns per user segment and flagged anomalies in real time, reducing fraud losses by 73% in the first quarter.
The key operational challenge was managing Flink's state at scale. With 2 billion events per day, checkpoint sizes grew to several terabytes. We tuned RocksDB state backend parameters extensively, implemented incremental checkpoints, and designed our state schema to minimize serialization overhead. We also built custom monitoring that tracked checkpoint duration, state size growth rate, and backpressure metrics, alerting us before problems became outages.
Topics