4 ms·
I think the design could have been simpler with Kafka (which they touched on briefly): - Write counter changes to a Kafka topic with many partitions. The parti
by fire_lake 2y ago
I think the design could have been simpler with Kafka (which they touched on briefly):
- Write counter changes to a Kafka topic with many partitions. The partition key is derived from the counter name.
- Use Kafka connect to push all counter events to S3 for audit and analysis.
- Write a Kafka consumer that reads events in batches and updates a persistent store with the current count.
- Pick a good Kafka message lifetime to ensure that topic size is kept under control, but data is not lost.
This gives us:
- Fast reads (count is precomputed, but potentially stale)
- Fast writes (Kafka)
- Correctness (every counter is assigned exactly one consumer)
- Durability (all state is in Kafka or the persistent store)
- Scalable storage and compute requirements over time
If I were to really go crazy with this, I would shard each counter further and use CRDTs to compute the total across all shards.
- rshrin 2y agoYes, this is one of the approaches mentioned in the article and is indeed a valid approach. One thing to keep in mind is that we are already operating the TimeSeries service for a lot of other high ROI use cases within Netflix. There already exists a lot of automation to self-provision, configure, deploy and scale TimeSeries. There already exists automation to move data from Cassandra to S3/Iceberg. We somewhat get all that for free. The Counter service is really just the Rollup layer on top of it. The Rollup operational nuances are just to give it that extra edge when it comes to accuracy and reliability.
- fire_lake 2y agoDid you also consider AWS managed services? Like a direct write to Dynamo?
- rshrin 2y agoNot for this use case. Other use cases at Netflix use AWS Managed service when it makes sense from a use-case and cost perspective. In this case, using TimeSeries opens the door to a lot of other potential future use cases: 1. What was the count for counter X between times T1 and T2? 2. "I am going to re-run my batch job again from yesterday. Adjust the increments for this window and re-compute the final count". Although the #2 use-case requires lot of other nuances around Recounting, which we allude to but don't expand upon in the article (adjustable retention, multiple rollup checkpoints per counter, pushing back accept-limit for backfills etc.)