Netflix TechBlog | Medium
Follow
Netflix’s Distributed Counter Abstraction
Netflix's Distributed Counter Abstraction is a service designed to store and query large volumes of temporal event data with low millisecond latencies. It supports two main categories of use cases: Best-Effort and Eventually Consistent. The Best-Effort counter uses EVCache for high throughput and low latency within a single region, but lacks cross-region replication and consistency guarantees. The Eventually Consistent counter uses a durable queuing system like Apache Kafka for accurate and durable counts, but can lead to delays and challenges in rebalancing partitions. Netflix's approach combines logging each counting activity as an event and continuously aggregating these events to meet requirements for auditing and recounting.