4 ms·
Having spent a lot of time in highly scalable distributed system land, I can perhaps give you few reasons: 1. Sometimes business needs require non-linear growt
by lemmsjid 6y ago
Having spent a lot of time in highly scalable distributed system land, I can perhaps give you few reasons:
1. Sometimes business needs require non-linear growth quickly. One day your perfectly optimized process now requires an n-squared algo across billions of records. Suddenly your one machine is tiny compared to what it used to be.
2. If you haven't scaled the workload horizontally, it can lead to a rewrite just to begin to distribute the load. When you hit the limits of a single machine, it is often a hard barrier that you cannot easily cross.
3. Distributed systems are distributing 3 main resources: memory, IO, and computation. IO especially quickly
becomes a bottleneck on a single system, but memory is also quite thorny because memory management itself can be a bottleneck on a vertically scaled system.
4. People with distributed systems do obsess over single-node QPS. If it takes 5k nodes to do work and you can optimize down to 2k nodes, you are saving a lot of money! However, it isn't that simple and this is where being properly distributed gives you cost leverage. You might find that 5 highly scaled machines are more costly than 50 commodity machines, especially in the cloud ecosystems.
5. Finally, and probably to your point and the GP's points, you kind of have to go with the flow when it comes to the level of abstraction people are writing the code at. It takes increasingly specialized knowledge to optimize a process that will run well on a 100gb process (e.g. virtual machine garbage collection issues). You are knowingly sacrificing efficiency for the nice higher level abstractions.
That said, I would emphasize that the abstractions become a smaller slice of the performance pie when you're dealing with algorithmic complexity. C won't magically make your Python algo O(1).
I don't disagree with your main point btw. I think a well done processing / data pipeline should have distributed systems available and single-node computational scenarios available, because there is an undeniable complexity gain when you reach for a distributed system immediately.
- closeparen 6y agoI'd argue that throughput that bad usually means you are doing something stupid and easy to fix, like making a blocking call in your single-threaded server process, doing per request what you should have done at startup, not indexing for your query pattern, etc. It's not like you need multiple person-years of Brendan Gregg-level wizardry to serve QPS in the low hundreds from <= 10 boxes. It should just happen, or else you should have a clear understanding of what about your workload is so intense.
- lemmsjid 6y agoTrue, I imagined while writing the above that people will be thinking of different use cases (or in this case pathological cases). It is very important to profile your system and rationalize all the things it is doing. I have multiple times, for example, seen production systems accidentally bottlenecked by having the wrong level of logging set, such that the primary cost of the system is trace logging, rather than anything having to do with customer value. I would argue, however, that it has little to do with the decision of whether or not to distribute your system. I've spent a lot of time dealing, for example, with bottlenecked single instance RDBMS instances that are handling load they shouldn't be handling. (For example those accidental recursive queries that are often a side effect of nice ORM abstractions) I totally agree with you, and have seen it happen, that people who do not understand the performance characteristics of their system can reach for a distributed solution before they've understood what their performance issue was. But I've also seen plenty of situations where they reach for bigger hardware for the same reasons. I think deciding whether or not to distribute means taking a disciplined approach to projecting the business needs of particular data entities you'll be dealing with. For example, if you have a users table, and it will ultimately store every human in the United States, then, well, that is quite do-able on today's single instance RDBMS systems, and you can project the theoretical growth over time. And if you need read load and HA, then you can go a replication route, or at least look at that first before doing something that reduces the quality of your transaction handling, like sharding. And then once the system is in place, taking a disciplined approach to profiling and quantifying the costs of the different aspects of the system and justifying their business value. For example having run large scale recommender systems, a typical decision might be to degrade the quality of an algorithm if it means saving a tremendous amount of money on processing.
- Olreich 6y agoTwo things: scaling well in a single node environment necessitates separating concerns such that scaling across nodes won’t suck as much. And, C often cuts down O(n^2) algorithms in python to O(1) because the algorithm is O(1), but Python was doing n^2 memory allocations behind the back.