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:

ResourceRuns out whenSymptomSized from
MemoryConcurrent queries need more than the daemon limitQueries queue or fail admission; spillingPer-host peak memory x concurrency
CPUScans and joins cannot keep up with the data volumeLong scan and join times at low memory useBytes scanned per second at the target latency
Local diskRemote reads repeat, or spills have nowhere to goSlow reads from object storage; scratch errorsHot working set; peak spill volume
NetworkShuffles of large joins and aggregations saturate linksExchange operators dominate profilesBytes 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.

Query profilesper-host peak, bytes, spillConcurrency targetsper pool, at peakLatency targetsper workload classRAM per hostsum(peak x conc) / shareHost countbytes / (rate x secs)Cache diskhot working setScratch diskspill x parallel spillsExecutors = CPU hosts + 1each with RAM from the memory checkValidate with a replaythen adjustper hostper host
Executor sizing: each resource gets its own requirement from measured inputs; CPU sets the host count and memory sets RAM 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 EXPLAIN output 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:

AspectFewer, larger hostsMore, smaller hosts
Per-query memory ceilingHigher per host, so big joins spill lessLower per host; large broadcast joins hurt
Loss of one hostRemoves a large share of capacityRemoves a small share
Exchange trafficLess, because more work stays localMore fragments, more network
Scan parallelismFewer hosts to spread small tables overMore hosts, but only if files allow it
Coordinator loadFewer fragments to scheduleMore 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

MistakeWhat happensFix
Sizing memory from EXPLAIN estimatesAdmission rejects or over-admits; OOM at runtimeUse measured per-host peaks from profiles
Average instead of peak concurrencyQueues every morning at the busiest hourSize to peak; use pools to shape bursts
mem_limit close to physical RAMHost swaps or the kernel kills impaladLeave room for OS, page cache, JVM and agents
Cache sized to the whole tableExpensive disks, low hit ratio gainsSize to the hot working set
Scratch on the cache diskSpills evict cache, both get slowerSeparate devices or quotas
No failure headroomOne dead host breaks latency targets at peakAdd capacity for at least one host
Ignoring file layoutMore hosts do not speed scansCompact small files; enough row groups to parallelise

What to do next

  1. Split the workload into classes and pools, and write a latency target and peak concurrency for each.
  2. Collect per-host peak memory, bytes scanned and spill from real profiles, at the median and 95th percentile.
  3. Measure per-core scan throughput on your data in a pilot; do not borrow a number.
  4. Compute RAM per host from memory and the host count from cores with the script above.
  5. Pick a host shape where the largest common query fits without spilling and one failure costs under 10% of capacity.
  6. Set mem_limit, data cache quota and scratch directories explicitly on every executor.
  7. Replay a real day at production concurrency and compare queueing, spill and CPU against the plan.
  8. Repeat the calculation quarterly and before any large onboarding.
Key takeaway: Size Impala executors from measured workload, not rules of thumb. Compute per-host memory as peak per-host query memory times peak concurrency plus headroom and keep mem_limit well below RAM; compute cores from bytes scanned, measured per-core scan rate and the latency target; size the data cache to the hot working set and scratch to peak spill. Let CPU set the host count plus one for failure and memory set RAM per host, choose a host shape that fits the biggest common query, and validate with a replay before buying.