Sizing an Impala cluster means deciding how many executor hosts to buy or rent and what each one should look like: how much memory, how many cores, how much local disk for caching and spilling. Get it wrong one way and queries queue in admission control or die with memory-limit errors; get it wrong the other way and you pay for idle machines. The answer comes from the workload, not from a rule of thumb, and this page shows how to derive it.
The method is bottom-up. Measure what your queries actually use, decide how many must run at once, and compute a requirement for each resource separately: memory, CPU, disk and network. The cluster size is the largest of those requirements, plus headroom and failure tolerance. Coordinator counts, catalog metadata and statestore growth are a separate control-plane problem covered in scaling Impala; this page is about the executors that do the work.
Four resources, one binds first
Each executor runs one impalad daemon that scans data, joins, aggregates and sorts. For any given workload one resource runs out first:
| Resource | Runs out when | Symptom | Sized from |
|---|---|---|---|
| Memory | Concurrent queries need more than the daemon limit | Queries queue or fail admission; spilling | Per-host peak memory x concurrency |
| CPU | Scans and joins cannot keep up with the data volume | Long scan and join times at low memory use | Bytes scanned per second at the target latency |
| Local disk | Remote reads repeat, or spills have nowhere to go | Slow reads from object storage; scratch errors | Hot working set; peak spill volume |
| Network | Shuffles of large joins and aggregations saturate links | Exchange operators dominate profiles | Bytes exchanged per query |
Memory is the binding constraint for most interactive Impala workloads because admission control reserves memory per host before a query starts. CPU binds for scan-heavy reporting over wide tables, and local disk binds when data lives in object storage. The model used here: memory sets RAM per host and how many queries each pool may run at once, CPU sets the host count, and disk sets local storage per host.
Collect the inputs
Sizing is only as good as its inputs, so gather them from the existing system or a representative pilot rather than guessing. For each workload class (dashboards, ad hoc analysis, batch reports) collect:
- Per-host peak memory of typical and 95th-percentile queries. The query profile records peak memory per node; the planner's per-host estimate in
EXPLAINoutput is a guess, often far off, so prefer measured peaks. - Bytes scanned per query after partition pruning and runtime filters, and the share served from local cache versus remote storage.
- Peak concurrency: the number of queries running at once at the busiest time of day, not the daily average.
- Target latency per class, such as 95% of dashboard queries under 5 seconds.
- Spill volume from profiles of the largest sorts, joins and aggregations.
- Growth: expected change in data volume and users over the planning horizon, usually 12 to 18 months.
If there is no existing system, run a pilot on a few nodes with real data and the twenty most important queries, and take the numbers from those profiles. A pilot of three or four executors is enough to measure per-core scan rates, but it overstates per-host memory for partitioned joins and aggregations, whose state is divided across hosts, while broadcast build sides and scan buffers do not shrink. Re-measure per-host peaks on a partial cluster closer to the target size, and leave headroom for network behaviour of large joins, which a small pilot cannot show.
Memory: per-host peak times concurrency
Admission control reserves memory on every host a query runs on, so the arithmetic is per host. For each pool, the memory needed on one host is roughly the per-host peak of a representative query multiplied by how many of those queries run concurrently. Sum across pools, then add headroom for the planner's errors and for queries that are larger than typical.
That figure becomes the daemon's memory limit, set with the -mem_limit startup flag, either in bytes or as a percentage of host RAM. The host needs more RAM than the limit: the operating system, the page cache that helps local reads, monitoring agents and the daemon's own embedded JVM all live outside it. A common planning split is to give impalad around 70 to 80 per cent of RAM on a dedicated host and leave the rest; check your version's default and set the flag explicitly so nobody inherits a surprise.
Per-query limits matter as much as the total. A pool's maximum per-query memory caps how much one query may reserve on each host, which protects concurrency from a single runaway query; queries that need more must spill. Detailed knobs are in Impala memory limits and pool design is in admission control.
CPU: from scan throughput
CPU sizing starts from scan throughput, because scanning and decoding Parquet or ORC is usually the largest CPU cost. Measure how many bytes per second one core scans and filters for your data and queries; it varies widely with compression, column count, predicate selectivity and whether data is read from local cache. Treat any figure you have not measured on your own data as an assumption.
Given that rate, the cores needed for one query to finish in its target time is bytes scanned divided by the product of per-core rate and target seconds. Multiply by concurrency for the cluster total, then divide by usable cores per host. Joins, aggregations and expression evaluation add CPU on top of scanning; a pilot profile shows their share, and a multiplier of 1.5 to 2 times scan cost is a reasonable starting assumption until you measure.
Impala parallelises a scan across hosts and, within a host, across scanner threads and, with the MT_DOP query option, across multiple fragment instances. More hosts help only if there are enough files and row groups to keep every core busy; a table with 200 files cannot use 1,000 cores in one scan. File layout is therefore part of sizing.
Disk: data cache and scratch
Executors use local disk for two very different things, and they should be sized separately.
The data cache. When tables live in object storage or remote HDFS, the data cache stores file ranges it has recently read on the executor disks. Enable it with --data_cache=dir1,dir2:quota, where the quota applies to each directory, so --data_cache=/data/0,/data/1:500GB allows up to 1 TB on that host. Size the total across the cluster to the hot working set, the data queried repeatedly over a day or a week, not the whole table size. A cache that holds the last 30 days of a fact table queried mostly for the last 7 days is mostly wasted. See the Impala data cache for behaviour and metrics.
Scratch space for spilling. Operators that exceed their memory spill to directories listed in --scratch_dirs. Size scratch from the largest spill seen in profiles multiplied by the number of large queries that can spill at once, and put it on fast local disks separate from the cache where possible, so a big spill does not evict cached data or compete for the same device. Spill to disk explains when spilling helps and when it only hides an undersized cluster.
Node shape: fewer large hosts or more small ones
Once you know the totals, choose how to divide them into hosts. The same 1,536 cores can be twenty-four 64-core hosts or forty-eight 32-core hosts, with RAM per host following from the memory check, and the choice changes behaviour:
| Aspect | Fewer, larger hosts | More, smaller hosts |
|---|---|---|
| Per-query memory ceiling | Higher per host, so big joins spill less | Lower per host; large broadcast joins hurt |
| Loss of one host | Removes a large share of capacity | Removes a small share |
| Exchange traffic | Less, because more work stays local | More fragments, more network |
| Scan parallelism | Fewer hosts to spread small tables over | More hosts, but only if files allow it |
| Coordinator load | Fewer fragments to schedule | More fragments per query |
Broadcast joins copy the small side to every host, so the per-host ceiling matters most for workloads with large dimension tables. Failure tolerance pushes the other way: if you can only afford to lose 10% of capacity when one host dies, you need at least ten. A balanced starting point for mixed workloads is a host size where the largest common query fits in its per-host limit without spilling, and enough hosts that one failure costs under 10% of capacity.
Worked example: two pools
Take a cluster serving two pools, with numbers measured in a pilot. Dashboards: per-host peak 2 GB, 40 concurrent at peak, each scanning 20 GB in under 5 seconds. Analysts: per-host peak 12 GB, 8 concurrent, each scanning 400 GB in under 60 seconds. Pilot scan rate: 150 MB per second per core (an assumption for this example; measure yours). Candidate host: 256 GB RAM, 32 usable cores.
import math
pools = {
# per-host peak GB, concurrency, GB scanned, target seconds
"dashboard": (2, 40, 20, 5),
"analyst": (12, 8, 400, 60),
}
host_ram_gb, host_cores = 256, 32
daemon_fraction = 0.75 # mem_limit as a share of host RAM
mem_headroom = 1.3 # estimate error, larger-than-typical queries
cpu_multiplier = 1.7 # joins and aggregation on top of scanning
scan_mb_per_core_s = 150 # measured in the pilot; replace with yours
per_host_mem = sum(peak * conc for peak, conc, _, _ in pools.values()) * mem_headroom
mem_limit_gb = host_ram_gb * daemon_fraction
print(f"needed per host {per_host_mem:.0f} GB vs limit {mem_limit_gb:.0f} GB")
cores = sum(
conc * (gb * 1024) / (scan_mb_per_core_s * secs) * cpu_multiplier
for _, conc, gb, secs in pools.values()
)
cpu_hosts = math.ceil(cores / host_cores)
print(f"cores needed {cores:.0f} -> {cpu_hosts} hosts")
hosts = cpu_hosts + 1 # survive one host failure at peak
print(f"executors: {hosts}")Memory first: 40 x 2 plus 8 x 12 is 176 GB per host, and with 1.3 headroom that is about 229 GB. A 256 GB host at 75% gives a 192 GB limit, so this host shape is too small for memory; either the analyst pool must accept lower concurrency or spilling, or hosts need 384 GB. Adding hosts helps only partly: partitioned join and aggregation state spreads thinner, but broadcast build sides and scan buffers are needed on every host, so treat this as a RAM-per-host decision and re-measure peaks at larger scale.
CPU next: dashboards need 40 x 20,480 / (150 x 5) x 1.7, about 1,857 cores; analysts need 8 x 409,600 / (150 x 60) x 1.7, about 619 cores. That is about 2,476 cores, or 78 hosts of 32 cores, plus one for failure: 79. With 384 GB hosts the memory check passes at 288 GB of limit. The decision is 79 hosts of 384 GB, and the dashboards dominate CPU because their latency target is tight. Relaxing the dashboard target to 8 seconds would remove about 20 hosts, which is the kind of trade to put in front of the people who set the targets.
Validate before you buy
A sizing calculation is a hypothesis. Validate it before buying the full fleet by replaying a day of real queries at production concurrency against a partial cluster and scaling the result, then again on the full cluster before go-live. While it runs, watch admission queue length and wait time per pool, the share of queries that spill and how much, per-host memory against the limit, CPU utilisation per host, cache hit ratio, and exchange time in the slowest profiles.
Read the results by resource. Queues with idle CPU mean memory or pool limits are too tight. High CPU with no queueing means you are CPU-bound and more hosts will help. Heavy exchange time means the network or the join strategy is the problem, and adding hosts can make it worse. Re-run the sizing each quarter with fresh profiles; workloads drift faster than hardware plans.
Common sizing mistakes
| Mistake | What happens | Fix |
|---|---|---|
| Sizing memory from EXPLAIN estimates | Admission rejects or over-admits; OOM at runtime | Use measured per-host peaks from profiles |
| Average instead of peak concurrency | Queues every morning at the busiest hour | Size to peak; use pools to shape bursts |
| mem_limit close to physical RAM | Host swaps or the kernel kills impalad | Leave room for OS, page cache, JVM and agents |
| Cache sized to the whole table | Expensive disks, low hit ratio gains | Size to the hot working set |
| Scratch on the cache disk | Spills evict cache, both get slower | Separate devices or quotas |
| No failure headroom | One dead host breaks latency targets at peak | Add capacity for at least one host |
| Ignoring file layout | More hosts do not speed scans | Compact small files; enough row groups to parallelise |
What to do next
- Split the workload into classes and pools, and write a latency target and peak concurrency for each.
- Collect per-host peak memory, bytes scanned and spill from real profiles, at the median and 95th percentile.
- Measure per-core scan throughput on your data in a pilot; do not borrow a number.
- Compute RAM per host from memory and the host count from cores with the script above.
- Pick a host shape where the largest common query fits without spilling and one failure costs under 10% of capacity.
- Set mem_limit, data cache quota and scratch directories explicitly on every executor.
- Replay a real day at production concurrency and compare queueing, spill and CPU against the plan.
- Repeat the calculation quarterly and before any large onboarding.