Tensor vs pipeline vs expert parallelism for MoE in mlx-lm
Parent: Mac local LLMs: Clusters, RDMA, exo and ds4 · Published reference · snapshot 2026-10-05
↓ Facts as markdownall context files
Tensor parallelism (TP): every layer's weight matrices are split across ranks; one all-reduce per sharded block. Pipeline parallelism (PP): whole layer ranges per rank, activations hop between ranks. Expert parallelism (EP, true meaning): whole experts live on different ranks and tokens are route...
These notes link each claim to its source. A source may be a research report hosted on this site rather than the primary document. A published reference means the content is available; it does not certify independent review or accuracy.Read the editorial policy and follow the sources before relying on a claim.
Facts
- Tensor parallelism (TP): every layer's weight matrices are split across ranks; one all-reduce per sharded block. Pipeline parallelism (PP): whole layer ranges per rank, activations hop between ranks. Expert parallelism (EP, true meaning): whole experts live on different ranks and tokens are routed to them by all-to-all dispatch/combine. In mlx-lm today, MoE models get TP (each expert's matrices are sliced) and PP only. EP does not exist in MLX or mlx-lm. [source]
- Terminology trap: the guruswami-ai mlx-benchmarks repo uses "EP" to mean Embarrassingly Parallel (one independent full copy per node, no communication). It says explicitly that Expert Parallelism is a different, unshipped concept (MLX PR #3158). [source]
- mlx-lm MoE TP (read from mlx_lm/models/deepseek_v3.py and qwen3_moe.py `shard()`): attention q/k/v (or q_b) projections `all-to-sharded`, o_proj `sharded-to-all`; dense MLP likewise. For MoE layers the router (gate) stays replicated on every rank; `shared_experts` and `switch_mlp.{gate,up,down}_proj` are sharded in place with all-to-sharded / sharded-to-all (an intermediate-dimension slice of EVERY expert on EVERY rank); the MoE block sets `sharding_group` and does a single `mx.distributed.all_sum` on its output. n_heads and n_kv_heads are divided by N. [source]
- Consequence: all ranks run the same routed experts for the same tokens, each computing a slice. Every rank reads 1/N of every active expert's bytes. There is no token routing between ranks, so no all-to-all. [source]
- mlx-lm PP (mlx_lm/models/pipeline.py `PipelineMixin.pipeline`): layers split in reverse, rank 0 gets the LAST layers and rank size-1 the first; even split with extra layers on low ranks, optional explicit `split`; non-owned layers are replaced by None so their weights are never loaded. Forward: each rank except the highest first `recv_like` from rank+1, runs its layers, `send`s to rank-1; the final hidden state is broadcast to all ranks with `all_gather` (kept in-graph via `mx.depends` on the last cache key). [source]
- Pipeline example (Awni Hannun, 2025-11-07): `mlx.launch --hosts first.ip,second.ip --env MLX_METAL_FAST_SYNCH=1 mlx-lm/mlx_lm/examples/pipeline_generate.py --model mlx-community/Kimi-K2-Thinking --prompt "..." -m 16384` ran the 1T int4-QAT Kimi K2 Thinking at about 15 tok/s on 2 M3 Ultras (about 3,500 tokens generated). [source]
- TP can only split dimensions evenly: Kimi K2 Thinking (MLX Q4) loads as TP2/TP3/TP4 but TP5 fails (12288 / 5 not an integer); Llama 405B (128 Q heads, 8 KV heads), Mixtral (32 Q heads), DeepSeek V3 (128 heads), Kimi K2.5 (64 heads) cannot do TP5. PP has no such constraint (PP5 ran fine). [source]
- PP loads faster and uses less memory than TP: only owned layers are loaded (Kimi K2 Thinking PP4 load 22.3 s, 145.9 GB peak per node vs TP4 99.9 s, 185.7 GB; PP5 5.2 s, 102.5 GB). [source]
- All-reduce count: the mlx-lm code does an all-reduce after attention (o_proj) and one after the MLP/MoE block, i.e. about 2 per layer, not 1 per layer as the mlx-benchmarks INTERCONNECTS doc states. [source]
- Ring-topology RDMA: RDMA also works in ring mode via `mlx.distributed_config`, so PP over a Thunderbolt ring does not need TCP/IP (guruswami-ai, Jan 2026); a ring of TB5 nodes avoids full-mesh cabling for PP. [source]
- `mx.distributed.all_to_all(x)` (split axis 0, chunk i to rank i) exists only in an unmerged PR (#3164, closed); MLX nn has shard_linear, shard_inplace, fully_shard, AllToShardedLinear, ShardedToAllLinear and quantized variants, nothing for expert placement. [source]
- 2026-02-23: 0xDaizz opened MLX PR #3158 "Add Expert Parallelism for MoE inference" (all_to_all, fused MoE dispatch/combine C++ primitives with CPU+Metal, blocking exchange_v, auto backend CPU for decode N<=64 and Metal for prefill N>=320, Python MixtureOfExperts layer). angeloskath asked for it to be split; sub-PR #3164 (all_to_all) was opened; 2026-08-15 zcbenz closed #3158: reviewing "takes as much time as writing it ourselves and at the moment this feature is not on our roadmap". #3164 is also closed. So EP for MLX stays out-of-tree. [source]
- 2025-11-07: pipeline_generate.py PR made 2-node Kimi K2 Thinking run in mlx-lm (pipeline-parallel). [source]
- 2026-01-12: first public PP-vs-TP RDMA benchmark on a 5-node M3 Ultra mesh (MLX discussion #2990). 2026-03 to 03-24: 290-point benchmark repo guruswami-ai/mlx-benchmarks (MLX 0.30.7, mlx-lm 0.30.8, macOS 26.3.1); their PRs added PP for Llama/Qwen2 and TP+PP for Mixtral. [source]
- 2026-03-05/06: mlx #3207 (RDMA file-transfer guide, moved to discussion #3481) and mlx-lm #955 (`mlx_lm share` docs proposal) document the PD-exhaustion workaround. [source]
- Kimi K2 Thinking (1T MoE, MLX Q4), batch 1, short context, 5x M3 Ultra 512 GB (MLX discussion #2990): PP5 14.45 tok/s, PP4 14.49, TP4 14.82, TP2 13.15 (TP5 impossible). TP4 vs PP4 = +2.3%. Author conclusion: RDMA makes TP all-reduce nearly free, and a PP ring on less RAM may match a TP mesh. Compare 2-node PP at 15 tok/s (Awni, Nov 2025) and existing-file Geerling exo TP 28.3 tok/s at 4 nodes: the 2x gap between mlx-lm PP/TP (about 14.5) and exo TP (28.3) for the same model is unexplained in sources. [source]
- Kimi K2.5 (1T, 614 GB at Q4) TP2 14.3 tok/s, TP4 16.1 tok/s (TP4 prompt 555 tok/s at 1K, TTFT 1.8 s; TP2 362 tok/s, 2.8 s) (mlx-benchmarks CHAKRA_CLUSTER). That is 1.13x for doubling nodes. [source]
- DeepSeek V3 (671B, 380 GB at Q4) fits one 512 GB M3 Ultra: 20.2 tok/s single node (mlx-benchmarks); with 37B active the single-node bandwidth model predicts it, the naive 671B model would predict about 0.9. Existing file's exo 21.1 tok/s at 1 node matches. [source]
- Mixtral 8x7B Q4 (47B/13B active): single 69.1 tok/s, PP2 62.5 (-10%), TP2 43.8 (-37%); prompt 740 -> 1100 tok/s with TP2. So for a MoE that fits one Mac, TP2 cuts decode by over a third. [source]
- Dense: Qwen 32B Q4 single 31.5, PP2 30.2 (-4%), TP2 21.6 (-31%), TP4 18.4 (-42%); prompt at 16K: single 289, PP2 538, TP2 547, TP4 958 tok/s. Llama 405B Q4 (202 GB, fits one 512 GB node): single 3.0 (prompt 28 tok/s at 1K, TTFT 37 s at 1K, about 10 min at 16K), TP2 4.3, TP4 6.4 (TTFT 10 s at 1K). [source]
- Quant interaction (Qwen 32B TP2): Q8 17.0 vs 18.4 single (93% retained), Q4 21.6 vs 31.5 (69%), Q2 25.0 vs 48.0 (52%). TP sync cost is roughly constant per layer, so smaller per-node weight reads make it dominant. [source]
- Single-node decode follows bandwidth/active bytes: about 620 GB/s effective of 819 GB/s on M3 Ultra (95-111% of that model; Q2 84% due to weaker dequant kernel). [source]
- Interconnect figures (mlx-benchmarks, M3 Ultra): RDMA TB5 sustained 5.3 GB/s per link (42.4 Gbps), about 50 us latency vs about 300 us TCP; NVLink 450 GB/s per GPU; 25 GbE about 200 us per all-reduce. Their modeled TP2 overhead on Llama 405B: TB5 RDMA about 36% of token time vs 72% on 25 GbE TCP (modeled, not measured). [source]
- llama.cpp RPC over TCP is latency-bound: reported 20 -> 2 tok/s going from wired gigabit to WiFi 6 (sharedllm.org); the worker's `-m` MB flag sets the backend memory pool it advertises, and overcommit crashes upload. [source]
- Secondary/unverified: exo on 8x M4 Pro Mac minis ran DeepSeek V3 671B at 5.37 tok/s (virge.io; no cluster interconnect or quant stated). [source]
- EP prototype numbers (PR #3158 author, 2-rank JACCL, E=384 D=7168 top_k=8, closed unmerged): decode EP 3.1x faster than TP (ratio 0.32), prefill TP 7% faster (EP/TP 1.07), geomean EP 1.8x faster. Single-author microbenchmark; not reproduced; PR rejected for review-cost reasons, not for the numbers. [source]
- Model fits one Mac: run single node for decode; TP/PP cost 4-42% decode. Distribute only for prompt speed (TP) or long-context memory headroom (PP2). [source]
- Model larger than one Mac (Kimi K2 class, ~1T): 2 nodes is the minimum; PP works with 2-5 nodes, TP needs divisor-friendly node counts (2, 3, 4). [source]
- Many users, model fits: independent copies per node ("EP" in mlx-benchmarks) give linear aggregate throughput with zero communication. [source]
- Pipeline long-context prefill limit: Metal ~60 s command-buffer timeout kills PP prefill at about 1,500 tokens for Kimi K2 Thinking PP4/PP5 (PP5 1,472 ok, 1,523 timeout; PP4 1,523 ok, 2,073 timeout) while TP4 handled 4,117 tokens (20.4 s, 2.45 tok/s overall). Llama 405B PP2 (63 layers/node) and PP4 (31 layers/node) both hit the timeout; Qwen 32B PP2 and Mixtral PP4 do not. Awni Hannun acknowledged the issue ("frustrating... no good general solution yet"); chunked prefill was not in mlx.distributed as of Jan 2026. Existing file's 7-8K figure is from a later patched stack. [source]
- Hardware-level kernel bug hypothesis: 21 kernel panics on an M3 Ultra were triggered in `route` (7), `ifconfig` (5), `kernel_task` (4), `python` (1); author attributes them to kext com.apple.driver.AppleThunderboltRDMA version 0.0.1 with a null-pointer dereference when ifconfig/route reconfigure TB5 interfaces while RDMA queue pairs are initialized/released. Both nodes of an RDMA pair often crashed together; recovery was a hard power cycle. Repeated Metal GPU timeouts also caused one node to kernel panic. [source]
- `mlx.distributed_config --auto-setup` run twice per boot (e.g. 5-node then 2-node subset) corrupts ARP tables and RDMA device mappings and yields `RTR failed with errno 60`; fix is reboot and run it exactly once, then use hostfile subsets for TP2/TP4. RDMA device names (rdma_enX) are ephemeral and change each reboot, so regenerate the hostfile and copy it to all nodes after each boot. [source]
- STP (spanning tree) shutting down TB5 interfaces, dynamic enX renaming, bridge0 coming up, ARP corruption: the hacks needed on macOS 26.2 for MLX's own distributed stack (exo handles several itself). [source]
- PD-exhaustion budget conflict: mlx #3207 reports `Couldn't allocate protection domain` after about 60 separate JACCL sessions (one per file); mlx-benchmarks reports the pool exhausting after 2-3 `mlx_lm.share` broadcasts or distributed runs, with `RTR errno 96` or `errno 22`, and says to make a large-model (>400 GB) broadcast the FIRST RDMA operation after a clean boot. Both agree: reboot is the only reclaim; prefer one N-node broadcast over N-1 two-node transfers. Different counts, different workloads; the real per-boot limit is unverified. [source]
- Large broadcast stalls (>400 GB): Llama 405B Q8 (402 GB, 85 files) and Kimi K2.5 614 GB stalled on a degraded mesh; after clean reboot Kimi K2.5 (606 files) reached 4 nodes in 2 min 6 s at 5.2 GB/s. [source]
- `mlx_lm share` has its own launcher: wrapping it in `mlx.launch` hangs. It uses `all_sum` (rank 0 contributes data, others zeros) so a 5-node broadcast runs as fast per link as a 2-node copy (213 GB: 37 s at 6.1 GB/s 2-node; 51 s 5-node). Cosmetic errors after success: PID-file cleanup `CalledProcessError` and `OSError: [Errno 66] Directory not empty` on the atomic rename when the target directory exists; verify safetensors counts on each node instead of trusting exit code. Use `python -u` to see tqdm over SSH. [source]
- Redundant bridge fix script from #3207 (needed after every reboot and OS update): destroy bridge0, delete it from system prefs, create per-port network services with networksetup, disable "Thunderbolt Bridge". [source]
- mx.distributed.send/recv with 16 MB or 4 MB chunks or `MLX_METAL_FAST_SYNCH=0` do NOT avoid the 5 s Metal timeout; only `stream=mx.cpu` does. [source]
- exo #1847 (Apr 2026, exo.app 0.3.68, 3x M3 Ultra, macOS 26.3.1): errno 22/2/60 can fire while idle with no inference; a 4-node report adds a SIGSEGV in jaccl::Connection::post_recv, with a proposed guard that rechecks SharedBuffer and QP state before ibv_post_recv so a corrupted QP raises an errno instead of crashing. [source]
- exo placement API: GET /instance/previews?model_id=... lists valid placements with `sharding` ("Pipeline" or tensor), `instance_meta` (MlxRing or MlxJaccl), `memory_delta_by_node` and `error`; POST /instance creates the chosen one. `uv run bench/exo_bench.py --model M --pp 128,512 --tg 128 --max-nodes 2 --instance-meta ring|jaccl|both --sharding pipeline|tensor|both --repeat 3 --json-out f` benchmarks placements. On Linux exo runs on CPU only; mlx-cuda12/13 extras exist (DGX Spark support merged). [source]
- llama.cpp RPC: `ggml-rpc-server` exposes ALL accelerators on the host (or a single CPU device if none); choose with `--device NAME` or CUDA_VISIBLE_DEVICES; backends can be mixed (CUDA and Metal workers in one run); `-ngl 99` spreads layers across local plus RPC backends. A symptom list includes `GGML_ASSERT(tensor->ne[0] % 512 == 0)` on the worker. [source]
- Other wrappers: LocalAI `mlx-distributed` backend (experimental) defaults to the ring (pipeline) backend with a hostfile of "ip:port" per rank; JACCL uses a device matrix JSON with null diagonal and rank 0 as TCP coordinator; `worker p2p-mlx` auto-discovers workers. [source]
- Is TP worth it on RDMA for a 1T MoE? mlx-lm benchmark (discussion #2990): TP4 only 2.3% above PP4, 14.8 vs 14.5 tok/s; mlx-benchmarks Kimi K2.5: TP2 14.3 vs TP4 16.1; vs exo+Geerling (existing file): 4-node 28.3 tok/s for Kimi K2 Thinking and Apple's "nearly 3x" for a 27B dense. Differences in framework (mlx-lm vs exo), model dimensions and tuning are plausible causes; no source isolates them. [source]
- "EP" meaning: mlx-benchmarks (embarrassingly parallel) vs vLLM/llm-d/Jarvislabs (expert parallelism, all-to-all). Never mix them. [source]
- PD exhaustion threshold: about 60 sessions (mlx #3207) vs 2-3 operations (mlx-benchmarks). [source]
- All-reduces per layer: 1 (mlx-benchmarks INTERCONNECTS) vs about 2 in mlx-lm code (attention o_proj plus MLP/MoE output). [source]
- Existing file says a patch stack removed reboots by May 2026; mlx-benchmarks (Mar 2026, kext 0.0.1) still treats reboot as mandatory; no source tests mlx 0.32.2 on current macOS. [source]
- Whether PR-style EP ever lands in MLX or lives in exo/other forks; the all_to_all primitive is not in mlx mainline. [source]
- Why mlx-lm TP4 on Kimi gives about 15 tok/s while exo TP4 gives 28 tok/s. [source]
- Whether expert-sliced TP (mlx-lm) loses to whole-expert EP for decode on JACCL outside the single microbenchmark. [source]
- Chunked prefill status for mlx-lm PP in current releases (7-8K PP2 per existing file implies partial mitigation). [source]
- True per-boot PD budget on macOS 26.4+. [source]
- mlx-lm tensor parallelism for DeepSeek V3 and Qwen3-MoE shards every expert's switch_mlp gate/up/down projections (all-to-sharded and sharded-to-all) on every rank and keeps the router replicated [source]
- In mlx-lm Qwen3-MoE the MoE block sets sharding_group and calls mx.distributed.all_sum on its output; heads and kv heads are divided by the group size [source]
- mlx-lm pipeline parallelism gives rank 0 the last layers and the highest rank the first layers, uses recv_like/send between ranks and an all_gather of the final hidden state [source]
- PipelineMixin splits layers evenly with extra layers on low ranks, accepts an explicit split, and sets non-owned layers to None [source]
- Awni Hannun ran Kimi K2 Thinking (int4 QAT, 1T) on 2 M3 Ultras at about 15 tok/s with pipeline_generate.py and MLX_METAL_FAST_SYNCH=1 [source]
- Expert parallelism for MLX (PR #3158, all_to_all, MoE dispatch/combine, MixtureOfExperts layer) was closed unmerged on 2026-08-15 with "not on our roadmap" [source]
- The all_to_all sub-PR #3164 for mx.distributed.all_to_all was closed, not merged [source]
- PR #3158 author benchmark on 2-rank JACCL: EP 3.1x faster than TP in decode (N<=64), TP 7% faster in prefill, EP 1.8x geomean [source]
- mlx-benchmarks says "Expert Parallelism" is a different concept from its "EP" (embarrassingly parallel independent copies) and calls it an advanced research topic (MLX PR #3158) [source]
- Kimi K2 Thinking MLX Q4 on 5x M3 Ultra: PP5 14.45, PP4 14.49, TP4 14.82, TP2 13.15 tok/s; TP5 fails because 12288 is not divisible by 5 [source]
- Kimi K2 Thinking PP4 loads in 22.3 s using 145.9 GB per node; TP4 loads in 99.9 s using 185.7 GB [source]
- RDMA works in ring mode set up by mlx.distributed_config, so pipeline parallelism over a TB5 ring does not need TCP/IP [source]
- PP prefill on Kimi K2 Thinking hit the Metal GPU timeout at 1,523 tokens (PP5) and 2,073 tokens (PP4) while TP4 handled 4,117 tokens [source]
- Awni Hannun called the GPU timeout "a frustrating one" with no good general solution yet (Jan 2026) [source]
- 21 kernel panics on an M3 Ultra were attributed to AppleThunderboltRDMA kext 0.0.1 triggered by ifconfig/route during TB5 mesh configuration (author's analysis, not Apple-confirmed) [source]
- Kimi K2.5 (614 GB Q4) TP2 14.3 tok/s and TP4 16.1 tok/s on M3 Ultras; TP4 TTFT 1.8 s at 1K prompt [source]
- Mixtral 8x7B Q4 single 69.1, PP2 62.5, TP2 43.8 tok/s on M3 Ultra [source]
- Qwen 32B Q4 single 31.5, PP2 30.2, TP2 21.6, TP4 18.4 tok/s; 16K prompt speed 289, 538, 547, 958 tok/s [source]
- Llama 405B Q4 single 3.0, TP2 4.3, TP4 6.4 tok/s; Llama 405B PP2 and PP4 hit the Metal 60 s timeout [source]
- Qwen 32B TP2 retains 93% of single-node decode at Q8, 69% at Q4 and 52% at Q2 [source]
- DeepSeek V3 671B Q4 (380 GB) runs on one 512 GB M3 Ultra at 20.2 tok/s because decode scales with about 37B active parameters [source]
- M3 Ultra effective decode bandwidth is about 620 GB/s of 819 GB/s peak in the benchmark model [source]
- TB5 RDMA sustained 5.3 GB/s per link (42.4 Gbps) with about 50 us latency vs about 300 us over TCP (author's measurements) [source]
- Running mlx.distributed_config --auto-setup twice in one boot corrupts ARP and RDMA device mappings and gives RTR errno 60 [source]
- mlx-benchmarks reports the RDMA protection-domain pool exhausting after 2-3 mlx_lm.share broadcasts or distributed runs, with errno 96 or 22, recoverable only by reboot [source]
- mlx #3207 reports PD exhaustion after about 60 per-file JACCL sessions and fixes it with a single session [source]
- A >400 GB model broadcast should be the first RDMA operation after a clean boot; Kimi K2.5 614 GB reached 4 nodes in 2 min 6 s at 5.2 GB/s after a clean reboot [source]
- mlx_lm share must be run directly, not under mlx.launch, which hangs it [source]
- mlx_lm share broadcasts via all_sum with rank 0 sending and others contributing zeros, so a 5-node broadcast runs as fast per link as a 2-node copy; 213 GB took 37 s (2 nodes) and 51 s (5 nodes) [source]
- mlx_lm share errors OSError Errno 66 Directory not empty and PID-file cleanup CalledProcessError are cosmetic when files transferred; verify safetensors counts [source]
- RDMA device names (rdma_enX) change after every reboot, so the hostfile must be regenerated and copied to all nodes [source]
- Smaller send/recv chunks and MLX_METAL_FAST_SYNCH=0 do not avoid the Metal 5 s timeout in JACCL file transfer; stream=mx.cpu does [source]
- A jaccl crash can occur at idle with no inference running (exo #1847) [source]
- A 4-node exo report attributes a SIGSEGV in the RDMA receive path to jaccl::Connection::post_recv and proposes a QP-state guard before ibv_post_recv [source]
- exo exposes GET /instance/previews listing placements with sharding, instance_meta (MlxRing, MlxJaccl) and memory_delta_by_node, and POST /instance to create one [source]
- exo_bench.py filters placements by --max-nodes, --instance-meta ring|jaccl|both and --sharding pipeline|tensor|both [source]
- exo runs on CPU only on Linux and offers mlx-cuda12/13 extras; DGX Spark CUDA 13 support was merged [source]
- llama.cpp ggml-rpc-server exposes all accelerators on its host, selectable with --device or CUDA_VISIBLE_DEVICES, and CUDA and Metal workers can be combined [source]
- llama.cpp RPC throughput fell from 20 to 2 tok/s moving from wired gigabit to WiFi 6 in one report [source]
- llama.cpp rpc-server -m sets the advertised backend memory pool in MB and overcommitting crashes the tensor upload [source]
- LocalAI mlx-distributed backend defaults to ring (pipeline) and supports JACCL (tensor) with a device-matrix hostfile and rank-0 coordinator [source]
- True expert parallelism assigns whole experts to different GPUs and uses all-to-all dispatch and combine; on NVLink/InfiniBand the all-to-all is manageable, on slow links it bottlenecks [source]
- exo on 8 M4 Pro Mac minis ran DeepSeek V3 671B at 5.37 tok/s (unverified secondary) [source]
- Neither Apple's MLX distributed docs nor WWDC26 sessions 232/233 mention expert parallelism or TP divisibility constraints [source]
- For MoE decode over Thunderbolt, sliced-expert TP needs all ranks to read every active expert, so it shrinks per-rank bytes by N but adds per-layer all-reduce latency; this is why a small-active MoE that fits one Mac loses decode speed when distributed [source]
Children
- No children recorded.