docs: add Distributed Computing category intro

Distributed Computing had no intro, so its meta description fell back to generic text and readers got no guidance on which tool to pick.

Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
Vinta Chen
2026-09-27 07:45:23 +08:00
co-authored by Claude
parent 2d9b988863
commit da7c0b7342
@@ -0,0 +1,22 @@
From a laptop to a cluster, Python distributed computing comes down to the workload: joblib runs loops, Dask scales pandas, Ray scales ML, and PySpark runs SQL.
How to choose:
- Parallel for loops on one machine: joblib
- Scaling pandas, NumPy, or scikit-learn code: Dask
- Training, tuning, and serving ML models, or GPU jobs: Ray
- Your own Python code with stateful workers on a cluster: Ray Core
- ETL and SQL on structured data, or a JVM shop: PySpark
- MPI programs on HPC clusters and supercomputers: mpi4py
joblib is a [package for parallel computing and disk-based caching](https://joblib.readthedocs.io/en/stable/) that leaves your code as unmodified as possible. You [write a parallel for loop](https://joblib.readthedocs.io/en/stable/user_guide/parallel.html#common-usage) as a generator expression: `Parallel(n_jobs=2)(delayed(sqrt)(i ** 2) for i in range(10))`. By default, joblib runs the calls in separate worker processes. When one machine isn't enough, [switch its backend](https://joblib.readthedocs.io/en/stable/user_guide/parallel.html#setting-up-joblib-s-backend-with-parallel-config), and the same loop runs on a Dask, Ray, or Spark cluster. scikit-learn's [`n_jobs` runs on joblib](https://scikit-learn.org/stable/computing/parallelism.html) too, so it follows the backend you pick.
Dask [scales pandas, scikit-learn, and NumPy workflows](https://docs.dask.org/en/stable/why.html) with minimal rewriting. A Dask DataFrame is [a collection of pandas dataframes](https://docs.dask.org/en/stable/) on different computers, and Dask Arrays parallelize NumPy. Dask [runs without any setup](https://docs.dask.org/en/stable/deploying.html#local-machine) on your laptop. Its `LocalCluster` follows the same interface as every other Dask cluster manager, so you swap it out when you're ready to scale up.
Ray is a [unified framework for scaling AI and Python applications](https://docs.ray.io/en/latest/ray-overview/index.html). Its libraries each distribute one ML task: Ray Data, Train, Tune, Serve, and RLlib. Ray Core runs your own code: [decorate a function with `@ray.remote`](https://docs.ray.io/en/latest/ray-core/walkthrough.html), call it with `.remote()`, and fetch the result with `ray.get()`. For workers that keep state between calls, decorate a class the same way to get an actor. Ray runs on one machine with `ray.init()`, and on several nodes once you [deploy a Ray cluster](https://docs.ray.io/en/latest/cluster/getting-started.html). Ray Data [suits GPU workloads for deep learning inference](https://docs.ray.io/en/latest/data/comparisons.html#how-does-ray-data-compare-to-other-solutions-for-offline-inference) better than Spark does, but unlike Spark, it has no SQL interface.
PySpark is the Python API for Apache Spark, for large-scale data processing. Dask's own comparison suggests Spark [when you prefer SQL, run mostly JVM infrastructure, or want an all-in-one solution](https://docs.dask.org/en/stable/spark.html#reasons-you-might-choose-spark). PySpark's docs [recommend DataFrames over RDDs](https://spark.apache.org/docs/latest/api/python/index.html), so Spark builds the most efficient query for you. Write it in [SQL or the DataFrame API](https://spark.apache.org/docs/latest/api/python/user_guide/sql.html), whichever you think in, and switch between the two as you go. Installing PySpark with pip is [for local use or as a client](https://spark.apache.org/docs/latest/api/python/getting_started/install.html) that connects to a cluster, not for setting up the cluster itself. The `pyspark` package also needs Java.
mpi4py provides [Python bindings for MPI](https://mpi4py.readthedocs.io/en/stable/), the Message Passing Interface, so your Python code runs across workstations, clusters, and supercomputers. Lowercase methods like `comm.send` pass any picklable Python object, and uppercase ones like `comm.Send` [pass NumPy arrays the fast way](https://mpi4py.readthedocs.io/en/stable/tutorial.html). Run your script with `mpiexec -n 4 python -m mpi4py script.py`, so an unhandled exception [aborts the whole MPI run instead of deadlocking](https://mpi4py.readthedocs.io/en/stable/mpi4py.run.html#exceptions-and-deadlocks). mpi4py runs on an MPI implementation like MPICH or Open MPI, and in production its docs [recommend a custom-built or system-provided one](https://mpi4py.readthedocs.io/en/stable/install.html).
Start on one machine: parallelism [brings extra complexity and overhead](https://docs.dask.org/en/stable/best-practices.html#start-small), and often you don't need it. When you do scale out, collect results at the end, since Dask's `compute()` and Ray's `ray.get()` both block until the work finishes. Call [`compute()` once](https://docs.dask.org/en/stable/best-practices.html#avoid-calling-compute-repeatedly) rather than in a loop, and [`ray.get()` as late as possible](https://docs.ray.io/en/latest/ray-core/tips-for-first-time.html#tip-1-delay-ray-get).