5 ms·
As someone unfamiliar with db management, is it really less operational overhead to have to physically scale your hardware than using a distributed option with
by multifascia 6y ago
As someone unfamiliar with db management, is it really less operational overhead to have to physically scale your hardware than using a distributed option with more elastic scalability capabilities?
- stu2010 6y agoRelational databases enable some very flexible data access patterns. Once you shard, you lose a lot of that flexibility. If you move away from a relational model, you lose even more flexibility and start having to do much more work in your application layer, and usually start having to use more resources and developer time every step of the way. The productivity enabled by having one master RDBMS is a big deal, and if they can buy commodity servers that satisfy their requirement, this seems like a fine way to operate.
- hinkley 6y agoIf I had a billion dollars, I'd put a research group together to study the prospects of index sharding. That is, full table replication, but individual servers maintaining differing sets of indexes. OLAP and single request transactions could be routed to specialized replicas based on query planning, sending requests to machines that have appropriate indexes, and preferably ones where those indexes are hot.
- ddorian43 6y agoThe problem is the network. You need billion dollars to fix the network so it's as fast as local ram/nvme.
- zinekeller 6y ago... and even if you somehow solved that, the law of physics hits you hard. Latency can be a real performance killer generally and is doubly true in database-type computing.
- hinkley 6y agoAre you guys talking about single server, non redundant databases? That’s not even apples and oranges. More like watermelons and blueberries.
- hinkley 6y agoWhy would that be the case? In this case we have already accepted that multiple servers will be involved. That means the limitations of the networking are a given.
- jandrewrogers 6y agoThis has been done in popular commercial databases for decades, and is thoroughly researched. As far as I know, these types of architectures are no longer used at this point due to their relatively poor scalability and write performance. I don't think anyone is designing new databases this way anymore, since it only ever made sense in the context of spinning disk. The trend has been away from complex specialization of data structures, secondary indexing, etc and toward more general and expressive internal structures (but with more difficult theory and implementation) that can efficiently handle a wider range of data models and workloads. Designers started moving on from btrees and hash tables quite a while ago, mostly for the write performance. Write performance is critical even for read-only analytical systems due to the size of modern data models. The initial data loading can literally take several months with many popular systems, even for data models that are not particularly large. Loading and indexing 100k records per second is a problem if you have 10T records.
- hinkley 6y agoWe pick up algorithms from 30 years ago all the time. Nobody’s as bad as the fashion industry, but we sure do try. Part of it is short memories, but part of it is how the cost inequalities in our hardware shifts back and forth as memory or storage or network speeds fall behind or sprint ahead.
- jasonwatkinspdx 6y agoI agree this is an under appreciated strategy. Someone in my family worked for a hedge fund where one of their simple advantages was they just ran MS SQL on the biggest physical machine available at any given moment. Lots of complexity dodged by just having a lot of brute capacity.
- anewaccount2021 6y agoAs stated in the post, there are read replicas. Assuming their workload is primarily reads, this buys them a decent amount of redundancy.
- tanelpoder 6y agoAlso worth noting that scalability != efficiency. With enough NVMe drives, a single server can do millions of IOPS and scan data at over 100 GB/s. A single PCIe 4.0 x4 SSD on my machine can do large I/Os at 6.8 GB/s rate, so 16 of them (with 4 x quad SSD adapter cards) in a 2-socket EPYC machine can do over 100 GB/s. You may need clusters, duplicated systems, replication, etc for resiliency reasons of course, but a single modern machine with lots of memory channels per CPU and PCIe 4.0 can achieve ridiculous throughput... edit: Here's an example of doing 11M IOPS with 10x Samsung Pro 980 PCIe 4.0 SSDs (it's from an upcoming blog entry): https://twitter.com/TanelPoder/status/1352329243070504964 https://twitter.com/TanelPoder/status/1352329243070504964
- riku_iki 6y ago> 16 of them (with 4 x quad SSD adapter cards) in a 2-socket EPYC machine can do over 100 GB/s. It is more interesting if actual CPU can handle such traffic in context of DB load: encode/decide records, sort, search, merge etc.
- tanelpoder 6y agoYes, with modern storage, throughput is a CPU problem. And CPU problem for OLTP databases is largely a memory access latency problem. For columnar analytics & complex calculations it's more about CPU itself. When doing 1 MB sized I/Os for scanning, my 16c/32t (AMD Ryzen Threadripper Pro WX) CPUs were just about 10% busy. So, with a 64 core single socket ThreadRipper workstation (or 128-core dual socket EPYC server), there should be plenty of horsepower left.
- tanelpoder 6y agoAs I mentioned memory access latency - I just posted my old article series about measuring RAM access performance (using different database workloads) to HN and looks like it even made it to the front page (nice): https://news.ycombinator.com/item?id=25863093 https://news.ycombinator.com/item?id=25863093
- 6y ago
- speedgoose 6y agoThey may rely on the ACID properties of their database. Which makes everything simpler, easier, and safer. https://dev.mysql.com/doc/refman/8.0/en/mysql-acid.html https://dev.mysql.com/doc/refman/8.0/en/mysql-acid.html
- deleted 6y ago[deleted]
- baby 6y agoI’m guessing that since it is for registration and all, the usage might be write-driven, or at least equally balance between writes and reads. In addition, you really care about integrity of your data so you probably want serializability, avoid concurrency and potential write/update conflicts, and to only do the writes on a single server. For this reason it sounds to me that partitioning/sharding is the only way to really scale this: have different write servers that care about different primary keys.
- sumtechguy 6y agoThat really depends on your software. Something like a NOSQL style it is kind of built in that it will be distributed. But that backs the compute cost back into the clients. Each node is 'crap' but you have hundreds so it does not matter. Something like SQL server it comes down to how fast you can get the data out of the machine to clone it somewhere else (sharding/hashing, live/live backups, etc). This is disk, network, CPU. Usually in that order. In most of the ones I ever did it was almost always network that was the bottleneck. Something like a 10gb network card (was state of the art neato at the time, I am sure you can buy better now) you were looking at saturation of 1GB per second (if you were lucky). That is a big number. But depending on your input transaction rate and how the data is stored it can drop off dramatically. Put it local to the server and you can 10x that easy. Going out of node costs a huge amount of latency. Add in the req of say 'offsite hot backup' and it slows down quickly. In the 'streaming' world like kafka you end up with a different style and lots of small processes/threads which live on 'meh' machines but you hash it and dump it out to other layers for storage of the results. But this comes at a cost of more hardware and network. Things like 'does the rack have enough power', 'do we have open ports', 'do we have enough licenses to run at the 10GB rate on this router'. 'how do we configure 100 machines in the same way', 'how do we upgrade 100 machines in our allotted time'. You can fling that out to something like AWS but that comes at a monetary cost. But even virtual there is a management cost. Less boxes is less cost.