14 ms·
How do you load balance across replicas? How do you shard across replicas?
by LordHumungous 6y ago
How do you load balance across replicas? How do you shard across replicas?
- jpgvm 6y agoLoad balancing across read replicas is usually handled by your connection bouncer, say pgBouncer/pgPool/etc though you may also do some amount of more complex both L3 and L7 balancing if you get really big. Sharding is usually a matter of actually splitting the masters. There are many techniques for achieving this. If you want the database to do all the work you will probably want to use something like Citus for PostgreSQL or Vitess for MySQL. You can also build bespoke topologies using PostgreSQL logical or MySQL binlog replication. Failing that you can do application level sharding if you don't want the database doing anything fancy for you and manage each shard as an independent database cluster. By the time you actually need to do this you will be able to afford one of these options. :) In the meantime you will save a ton of CPU, storage and development time vs a "NoSQL" store as databases like PostgreSQL are inherently more efficient for all but the simplest of KV access patterns.
- deleted 6y ago[deleted]
- LordHumungous 6y ago> By the time you actually need to do this you will be able to afford one of these options. :) What if I need to do this now? Why would I build a distributed postgres snowflake that takes 10 hours to spin up a new replica, requires that I implement my own sharding, instead of using a datastore that is designed to handle all of these things at scale?
- jpgvm 6y agoComes down to your data model. If its inherently relational it's still the best play. Scale and performance are much more tractable problems than integrity and consistency. One you can measure and be sure of, the other you need a Phd to fully understand all the edge conditions that need covering. There are some pretty decent NoSQL stores now for simpler access patterns. As long as you stay away from nonsense like MongoDB and stick to real databases like Cassandra/ScyllaDB/BigTable/etc you will do fine. These stores are a fraction as flexible as PostgreSQL/MySQL but do allow scale-out storage and fast primary key lookups and scans. Good for when the size of your data is well in excess of 1TB+ and you don't need anything complex or consistency.
- donor20 6y agoReality - these folks don't need to "do this now". Yes, Visa may need this. Guess what, 1TB+ of transaction data paying 30 center + 2% PER LINE - you'll be able to afford to do something reliable and scalable. Folks don't realize, noSQL is not actually that scalable except in very narrow ways. And you can spin up pretty good scale SQL stuff with things look AWS RDS, including backups, replicas, snapshots to go back in time etc (noSQL doesn't support a lot of this).
- donor20 6y agoA lot of folks underestimate what one box can do. Memory / core counts have gone crazy on just one box. Local storage also has gone crazy. 4TB memory on a single node dual CPU machine, CPU's with 32 cores per CPU+? Read only replicas are pretty trivial as well.
- LordHumungous 6y agoA lot of people underestimate what a high scale workload means.
- jpgvm 6y agoAnd even more people talk about high scale workloads with no clue what they actually look like. :) I routinely work with 10TB+ PostgreSQL clusters, 10TB+ BigTable clusters and 500TB+ BigQuery projects all in my current day job. I'm in Data Infra btw so this is sort of my bread and butter. In the past I have worked with 100TB+ Cassandra clusters, 50TB+ MySQL+Vitess and countless other stores like MongoDB, RethinkDB, Voldemort, TokyoCabinet and probably tons I have forgotten. It's highly unlikely one actually works with and manipulates these volumes of data on a regular basis and doesn't respect SQL stores and the JVM (the literal king of Big Data).
- LordHumungous 6y agoI don't understand the point you are trying to make.