As Wix Kafka usage grew to 5B messages per day, over 20K topics and more than 100K leader partitions serving 2000 microservices. We decided to migrate from self-operated single cluster per data-center to a managed cloud service (Like Amazon MSK or Confluent Cloud) with a multi-cluster setup.
This talk is about how we automatically migrated all Greyhound consumers and producers even while they were handling regular production traffic with 0 down time and the lessons we learned along the way.
These lessons include:
1. clusters by SLA - choose your clusters wisely
2. Automation, Automation, Automation - all the process has to be completely automated at such scale
3. Prefer a gradual approach - E.g. migrate topics in small chunks and not all at once. Reduces risks if things go bad
4. and many more...