3 ms·
> However, we discovered after some time that the custom Python implementation for those workers was dropping up to 5% of the events. This was mostly due to the
by pul 9y ago
> However, we discovered after some time that the custom Python implementation for those workers was dropping up to 5% of the events. This was mostly due to the nature of how reading happens with Kinesis: every stream has multiple shards (ours up to 50!) and each reading client would use a so-called shard iterator to keep track of where it was reading last. Since the used machines could always crash, be recycled, or scaled down, we needed to save those shard iterators in some serialized format to Redis and share them across machines and process boundaries. Since we had so many shards, every once in awhile we would skip events and hence lose them.
I've never worked with Kinesis, but in Kafka you'd store offsets specifically to solve this issue. When one of the members of a consumer group would drop out, the partition (read: shard) would automatically be reassigned to another member. This gives an at least once delivery guarantee, combined with idempotent actions gives effectively once semantics. No need to loose any messages. What was the issue that the dubsmash engineers were solving here?
- alexatkeplar 9y agoWith Kinesis, you would just use the Kinesis Client Library (https://github.com/awslabs/amazon-kinesis-client-python https://github.com/awslabs/amazon-kinesis-client-python) which would automatically handle committing the offsets to DynamoDB. Home-rolling a checkpoint-free event pipeline is a rookie mistake; it's a pity they didn't come across our Snowplow project (Apache 2.0 event pipeline running on Kinesis, Kafka and NSQ, https://github.com/snowplow/snowplow/ https://github.com/snowplow/snowplow/).