Rebuilding the Data Path Behind Shopify’s BFCM Live Map
Shopify’s Black Friday and Cyber Monday (BFCM) live map has become a yearly tradition: a real-time, external-facing visualization of sales across over 1.7 million merchants. For the 2021 edition, the Shopify Data Platform Engineering team was tasked with not only scaling the system that powers it, but also adding richer insights — all within weeks of the event. The solution centered on introducing Apache Flink to the critical path, a move that allowed the team to process a significantly higher volume of data, maintain 100 percent uptime, and eliminate manual interventions.
Where the Old System Stumbled
The 2020 live map was backed by a bespoke, stateful streaming service built in Go, internally named Cricket. It consumed messages from relevant Kafka topics and pushed metrics to a React frontend via Redis. While functional, the architecture had inherent weaknesses that made scaling for 2021 risky.
The primary concern was volume. At peak during the 2021 BFCM weekend, the checkout Kafka topic saw roughly 50,000 messages per second. That topic carries far more events than the live map needs, along with fields the map never uses. Cricket was processing all of it, raising the specter of an overload. The path between Cricket and the frontend also had issues: Redis was used to queue and broadcast messages to browsers, a process that was inefficient and complex, and browser long polling could hang, causing the map’s arc visuals to momentarily disappear under load.
The metrics themselves have relaxed requirements compared to systems like chat. They can tolerate some data loss—the arc visuals are already a sample of orders since the browser can’t draw them all—and they only require the latest value, not a full history of broadcasts. Queuing the entire broadcast history was excessive.
For 2021, the team sought to add two new metrics on top of the existing sales and orders per minute and carbon offset:
- Product trends: The top 500 product categories by change in sale volume over the last six hours.
- Unique shoppers: The count of unique buyers per shop, aggregated over time.
Load tests showed that adding these metrics would make Redis a severe bottleneck, with the increase in messages and connections causing the live map’s visuals to disappear. With more data forecasted for the event, the existing architecture was not going to hold.
A Flink-Based Redesign
Rather than building from scratch, the team had to iterate on the existing system. The key move was placing Flink on the critical path. Flink’s job was to filter out irrelevant checkout events before they reached Cricket, dramatically reducing the volume Cricket had to process—down to one percent of the event stream. This resolved the core scaling issue while allowing Cricket to continue computing the original metrics.
To meet high availability requirements, the Flink jobs were deployed using a combination of cross-region sharding and cross-region active-active. Deduplication was handled in Cricket. For the new metrics, Flink was the source of computation, with Cricket acting as a relay to the frontend.
The new product trends metric relied on Shopify’s product categorization algorithm to emit the top 500 categories with sales quantity changes every five minutes. For a given product, the change was calculated as:
change = SUM(prior 1hr sales quantity) / MEAN(prior 6hr sales quantity) - 1
Results and a Few Rough Edges
The decision to offload computation to Flink proved sound. During the BFCM weekend, the Flink jobs ran with 100 percent uptime, no backpressure, and were free of manual intervention. As a safety net, the new metrics were also built as batch jobs on the existing Spark infrastructure; they worked, but were ultimately unused because Flink met or beat expectations.
Not everything went off without a hitch. The method for fetching messages from Redis and serving them to users caused high CPU loads on the machines. Cricket’s faster production rate, combined with the large memory footprint of the new product trends metric, clogged Redis. A small number of users saw some arc visuals blip into existence and then vanish. The Production Engineering team stepped in, dropped unnecessary Redis state, and resolved the backlog within two hours. The impact was minimal, and the live map shipped with the new metrics intact, processing significantly more events than the previous year.
Streamlining for the Next BFCM
The Flink-based system that powered this year's live map proved that mission-critical streaming applications are viable on the platform, with internal teams assembling sophisticated pipelines in just weeks. Beyond BFCM, the same streaming infrastructure could replace batch processing in other analytics products and visualizations, where data freshness is currently a limitation.
The next iteration aims to simplify the architecture by pushing more complexity into Flink itself. The goals are to eliminate the custom stateful stream processor, remove the system bottleneck, and reduce back-pressure handling to a single point instead of managing it across streaming jobs, Cricket, and the Web tier.
One promising design under consideration looks like this:
In this approach, Flink jobs would produce all metrics and periodically snapshot them to a database or key-value store. The Web tier would synchronize its in-memory cache on a schedule and serve polling requests directly from browsers. The design is straightforward and meets both scalability and complexity requirements.
The team plans to pursue this direction, building on the momentum from a successful BFCM where the platform handled record-breaking sales without incident.



