Zalando's Ad Platform team replaced a 7-year-old homegrown in-memory stream join with Apache Flink to match ad auction events with user interactions in near-real time within a 15-minute window. The homegrown Java app lacked checkpointing and proper stream partitioning, causing state loss on restarts and requiring overprovisioning. After evaluating alternatives (Nakadi SQL, Spark), they chose Flink with a RocksDB state backend and 3-minute incremental checkpoints. A key optimization was replacing Flink's built-in Interval Join with a custom KeyedCoProcessFunction using direct RocksDB point lookups, eliminating costly seek operations that consumed 15–30% CPU. Extensive RocksDB tuning (write buffers, compaction sizes, LZ4 compression, SSD-optimized settings) was needed to stabilize state growth. Infrastructure challenges included Kinesis connector limitations, Karpenter pod evictions interfering with Flink's autoscaler, and OOM kills requiring careful JVM memory parameter tuning. After a 4-week shadow pipeline validation and 1-week A/B test, they cut average pod count from 20 to 5, reduced memory from 320GB to ~100GB, and halved EC2 costs from ~€80 to ~€30/day, while improving event match rate by 0.5%.

13m read timeFrom engineering.zalando.com
Post cover image
Table of contents
BackgroundHomegrown solutionAlternativesArchitecture: What We BuiltThe Road to Production: What Actually HappenedMigrationResultsWhat's NextCredits
7K Impressions