The network indeed becomes the bottleneck. In 2 main ways:
1. Reading data from cloud storage is very expensive. Here’s a blogpost where we talk about some of the optimizations we’ve done in that area: https://blog.getdaft.io/p/announcing-daft-02-10x-faster-io
2. During a global shuffle stage (e.g. sorts, joins, aggregations) network transfer of data between nodes becomes the bottleneck.
This is why the advice is often to stick with a local solution such as DuckDB, Polars or Pandas if you can keep vertically scaling!
However, horizontally scaling does have some advantages:
- Higher aggregate network bandwidth for performing I/O with storage
- Auto-scaling to your workload’s resource requirements
- Scaling to large workloads which may not fit on a single machine. This is more common in Daft usage because we also work with multimodal data such as images, tensors and more for ML data modalities.
Hope this helps!
Depends on your workload size, if the compute time is less than time it takes to transfer the data back and forth, then it might be a bad idea to use it indeed.
Typically these solutions should only be tested after maxing out vertical scaling, before applying for horizontal scaling.
Network is one hell of a destroyer when it comes to advantages gained from distributed computing.