Impala is called an MPP engine, for massively parallel processing, and the label carries a specific design rather than a vague promise of speed. Every executor in the cluster runs copies of the same plan pieces against its own slice of the data, rows stream between those copies over the network without being written to disk, and one coordinator stitches the answer together. Nothing is shared between executors except the storage underneath and the plan they were handed.
This page is about how the work and the data get divided: fragments and instances, who reads which bytes, what crosses the network at each exchange, and where parallelism stops helping. By the end you should be able to estimate a plan's network traffic on paper and predict its straggler.
What MPP means in Impala
Strictly, MPP means shared-nothing: each worker owns its CPU, memory and local disks, and workers cooperate only by sending messages. Classic MPP databases also own their storage, so each node holds a fixed shard of every table. Impala departs from that in one important way. Storage is separate and shared: HDFS, Ozone, S3, ADLS or Kudu hold the files, and any executor can read any file. What stays shared-nothing is the execution. An executor never reads another executor's memory, never waits on a shared lock, and never coordinates with its peers except by exchanging row batches.
Because data is not pinned to nodes, the scheduler chooses who reads what on every query, so executors can be added without rebalancing data. Because execution is shared-nothing and streamed, a query is only as fast as its slowest piece, latency is low, and there is no stage boundary from which to resume after a failure. Compared with stage-by-stage engines such as Spark or Hive on Tez, which write shuffle output to disk, Impala trades mid-query fault tolerance for latency, which is why it suits interactive queries better than hours-long batch jobs.
From plan to fragments to instances
The planner first builds a single-node plan: scans at the leaves, then joins, aggregations, sorts and a final result node. It then turns that into a distributed plan by deciding where data must move and inserting an exchange at each such point. Everything between two exchanges becomes a plan fragment, the unit that gets shipped to executors. A fragment is a template; an instance is one running copy of it on one host, with its own operators, memory and threads.
The number of instances is decided per fragment. A fragment containing a scan gets one instance on each executor that was assigned scan ranges for it. A fragment that receives a hash-partitioned exchange usually gets one instance per participating executor, so every host owns a share of the key space. The root fragment runs as a single instance on the coordinator. With intra-node parallelism switched on (covered below) a host can run several instances of the same fragment.
In a distributed EXPLAIN, each PLAN FRAGMENT header names a fragment and its host and instance counts, and each EXCHANGE node marks where one fragment feeds the next.
Dividing the data: scan ranges and their assignment
Data is divided before any operator runs. The coordinator lists the files the query touches, after partition pruning, and cuts them into scan ranges: byte ranges of a file that a scanner can read independently, roughly one per storage block for splittable formats. The scheduler then assigns every scan range to exactly one executor. That assignment is the first and most important act of parallelism, because it fixes how much reading each host does.
For HDFS-style storage the scheduler prefers a host that holds a replica of the block, since a local read avoids the network. The REPLICA_PREFERENCE query option sets how strongly: CACHE_LOCAL (the default) prefers replicas in the HDFS cache, then local disk, then remote; DISK_LOCAL treats cached and disk replicas equally; REMOTE treats all three equally, which spreads load at the cost of locality.
Object stores have no locality at all, so every range is remote. Early releases placed remote ranges randomly, which meant a rerun of the same query read the same file from different hosts and missed every per-host cache. Since Impala 3.2 the placement of remote ranges is deterministic (IMPALA-7928): the same file range tends to land on the same small set of executors, which keeps the file handle cache and the local data cache warm across runs. The trade-off is a little less perfect balance in exchange for cache hits.
A query cannot use more scan parallelism than it has ranges, and imbalance in bytes becomes imbalance in time, so file sizing is a parallelism decision, not a storage detail.
Moving the data: exchanges
An exchange has a sender side, at the top of one fragment, and a receiver side, at the bottom of the next. Senders serialise row batches, optionally compress them, and stream them over RPC; receivers deserialise and hand them to their parent operator. There are three routing patterns you will see in plans:
| Plan label | Routing | Typical use | Network cost |
|---|---|---|---|
EXCHANGE [UNPARTITIONED] | every sender to one receiver | final merge at the coordinator; small results | size of the data, into one host |
EXCHANGE [HASH(k)] | row goes to the instance owning hash(k) | partitioned joins, distributed aggregation | about the size of the data, spread over all hosts |
EXCHANGE [BROADCAST] | every row to every receiver | broadcast join build side | size of the data times the number of receivers |
Scans, filters and pre-aggregations run in parallel with no communication. Every exchange is a synchronisation point in disguise: receivers finish only when every sender has, so the slowest sender sets the pace for the next fragment. Runtime filters travel from join build sides back to scans so that fewer rows reach the exchange at all.
Worked example: estimating network cost
Take a 20-executor cluster and this query: total revenue per customer region for last quarter, from a sales fact table and a customers dimension.
SELECT c.region, SUM(s.amount) AS revenue
FROM sales s JOIN customers c ON s.customer_id = c.id
WHERE s.sale_date >= '2026-07-01' AND s.sale_date < '2026-10-01'
GROUP BY c.region;Assume partition pruning leaves 2 billion sales rows, the scan keeps two columns, and a shuffled sales row costs about 16 bytes on the wire. The customers table has 50 million rows of about 24 bytes each, roughly 1.2 GB. These are illustrative numbers, but the arithmetic is the real method.
Partitioned join. Both sides are hash-exchanged on the customer id. Sales moves 2,000,000,000 x 16 bytes, about 32 GB, and customers about 1.2 GB, for about 33 GB in total, spread so that each host sends and receives roughly 1.6 GB. Each host builds a hash table of only its share of customers, about 60 MB.
Broadcast join. Sales stays where it was scanned and customers is sent to all 20 hosts: 1.2 GB x 20, about 24 GB of traffic, and each host holds the whole 1.2 GB hash table. Network traffic is lower, memory is twenty times higher, and every host must receive the complete build side before probing begins.
Neither is universally right. With 200 executors the broadcast costs 240 GB and the partitioned plan still about 33 GB; with a 100 MB dimension the broadcast is trivially cheap. The planner makes this choice from table statistics, which is why missing stats so often produce slow distributed plans. The aggregation then follows a two-phase pattern every MPP engine uses: each instance pre-aggregates locally, so 2 billion rows collapse to a few dozen regions per host; those partial sums cross a HASH(region) exchange for the merge; and the tiny result crosses an unpartitioned exchange to the coordinator. The expensive data never travels twice.
F02:PLAN FRAGMENT [UNPARTITIONED] hosts=1 instances=1
EXCHANGE [UNPARTITIONED]
F01:PLAN FRAGMENT [HASH(c.region)] hosts=20 instances=20
AGGREGATE [FINALIZE] group by: c.region
EXCHANGE [HASH(c.region)]
F00:PLAN FRAGMENT [RANDOM] hosts=20 instances=20
AGGREGATE [STREAMING] group by: c.region
HASH JOIN [INNER JOIN, BROADCAST] runtime filter -> s.customer_id
|-- EXCHANGE [BROADCAST] (from the customers scan fragment)
SCAN HDFS sales partitions=92/1460The fragment above is abridged and hand-annotated to show the shape, not verbatim Impala output; read your own plan with EXPLAIN at EXPLAIN_LEVEL=2 to see the real one.
Parallelism inside a host: MT_DOP
Impala historically ran one instance of each fragment per host and parallelised only the scan: several scanner threads feed a single instance whose joins and aggregations run on one thread, leaving cores idle when a large join or aggregation dominates.
The MT_DOP query option is the newer model. Set above zero, it runs that many instances of each fragment per host, each with its own operators, so joins and aggregations are parallel inside the host too. The documented range is 0 to 64 and the default is 0 for SELECT, while COMPUTE STATS on Parquet tables sets it to 4 automatically. Impala 3.4 and earlier supported it only for some plan shapes, so check what your release allows before relying on it.
-- try intra-node parallelism for one heavy query and compare
SET MT_DOP=8;
SELECT ... ;
PROFILE; -- in impala-shell: compare per-host time and peak memory
SET MT_DOP=0; -- back to the default for the sessionIt is not free: each extra instance has its own hash tables and buffers, so memory per host grows with the degree of parallelism. Raise it for CPU-bound queries on hosts with spare cores and memory, and measure before making it a default.
Skew and stragglers
An MPP query finishes when its slowest instance finishes. Twenty hosts that each take 10 seconds give a 10-second query; nineteen that take 10 and one that takes 60 give a 60-second query and nineteen idle hosts. The query SUMMARY shows this directly: for each operator compare Avg Time with Max Time. A ratio near one means balanced work; a large ratio means one instance carried more than its share.
The usual causes, in rough order of frequency:
- Hot keys in a hash exchange. If one customer id or a NULL key owns a large share of rows, the instance that owns its hash bucket does that share of the join or aggregation alone.
- Uneven scan ranges. A few huge unsplittable files, or partitions of wildly different size, give some hosts far more bytes to read.
- Remote reads on some hosts. When only part of the cluster has local replicas, the rest read over the network and lag.
Fixes follow from the cause: salt or pre-aggregate hot keys, filter NULL join keys early, compact small files, compute statistics, and take a persistently slow host out of service.
The serial tail
Some work cannot be spread. Planning and scheduling run on the coordinator before any executor starts. The root fragment, which merges partial results and returns rows to the client, is one instance. A global ORDER BY without a limit sends every row through a final merge on one host; with a LIMIT, each instance keeps only its own top-n and the merge stays small. Fetching a large result through a slow client keeps the query, and its resources, open until the client catches up.
This is Amdahl's law: if 2 seconds of a 10-second query are serial, no number of executors brings it below 2 seconds. Dedicated coordinators keep that work off the executors; pruning and modest result sizes keep it small.
Failure modes
Streaming execution has a sharp failure model. Fragment instances keep their state in memory and there is no materialised stage output to replay, so when an executor dies or becomes unreachable mid-query, the queries running fragments on it fail. Newer releases can transparently retry a failed query from the start on the remaining executors; that is a whole-query retry, not a fragment restart, so a long query loses all of its work.
- Memory multiplies with fan-out. A broadcast build side or a high MT_DOP repeats memory on every host or instance; per-host limits fire even when the cluster as a whole has plenty.
- Too many small files. Hundreds of thousands of tiny ranges make coordinator planning and scheduling the bottleneck.
- Stale statistics. Wrong row estimates pick broadcast for a table that grew tenfold, and the first sign is a memory-limit error on every host at once.
What to do next
- Run
EXPLAINatEXPLAIN_LEVEL=2on your three most expensive queries and label every exchange as unpartitioned, hash or broadcast. - For each join, estimate broadcast and partitioned network bytes on paper as in the worked example, and check whether the planner's choice matches.
- Open the
SUMMARYof a slow query and list every operator whose Max Time is more than twice its Avg Time; trace each to a hot key, uneven files or a slow host. - Check file counts and sizes for your largest tables, and schedule compaction where most files are far below the block size.
- Run
COMPUTE STATSon tables that changed significantly, then re-check the plans that misbehaved. - Trial
MT_DOPon one CPU-bound query, record time and peak memory per host, and only then decide whether to widen its use.
Keep going with Impala Architecture, Impala Join Strategies, Impala Query Plans, Impala Runtime Filters and Impala Admission Control.