Impala's planner makes a handful of decisions that dominate how fast a query runs and how much memory it needs: the order in which tables are joined, whether each join copies its smaller side to every executor or hash-partitions both sides across the cluster, whether an INSERT reshuffles and sorts rows before writing, and which replica of each block is scanned. It makes them from table and column statistics. When statistics are missing, stale or misleading, the decisions go wrong, and optimizer hints are the escape hatch that lets you pin one decision by hand.

This article is the reference for that escape hatch: every hint Impala supports, where each one goes in the SQL, what it changes in the plan, how to confirm it took effect, and the query options that are usually the better tool. Reading the plans themselves is covered in Impala query plans, which also works through a broadcast join that should have been a shuffle; here the worked example is the other common case, a partitioned INSERT. Hints are suggestions, not commands: since Impala 1.2.2 the planner can reorder joins in ways that make a hint irrelevant, which is why STRAIGHT_JOIN exists and why every hint must be checked in EXPLAIN.

Advertisement

Syntax and placement

Impala accepts hints inside comments that begin with a plus sign: block form /* +HINT */ and line form -- +HINT. Several hints can share one comment, separated by commas. The older square-bracket form such as [SHUFFLE] is still parsed for backward compatibility, but the documentation marks it deprecated and due for removal. The line-comment form is the one Hive tolerates, which matters if the same SQL file runs on both engines; Hive does not recognise the block-comment or bracket forms.

Placement is strict, and a misplaced hint is ignored rather than rejected, so the most common hint bug is a hint that never applied. Join hints go immediately after the JOIN keyword, before the right-hand table. STRAIGHT_JOIN goes immediately after SELECT. INSERT, UPSERT and CREATE TABLE AS SELECT hints go just before the SELECT (or right after INSERT). Scan scheduling hints go right after the table reference.

-- Join hints: right after JOIN; STRAIGHT_JOIN right after SELECT keeps your order
SELECT STRAIGHT_JOIN f.order_id, d.region
FROM sales_fact f
  JOIN /* +SHUFFLE */ customer_dim d ON f.customer_id = d.customer_id;

-- Same hint in line-comment form (also tolerated by Hive)
SELECT f.order_id, d.region
FROM sales_fact f
  JOIN -- +BROADCAST
  small_dim d ON f.dim_id = d.dim_id;

-- INSERT hints: before SELECT (or right after INSERT)
INSERT INTO sales_by_day PARTITION (sale_date)
  /* +SHUFFLE,CLUSTERED */
SELECT order_id, amount, sale_date FROM staging_sales;

-- Scan scheduling hints: after the table reference
SELECT count(*) FROM hot_lookup /* +SCHEDULE_CACHE_LOCAL,RANDOM_REPLICA */;

-- Deprecated bracket form: still parsed, slated for removal; do not write new code with it
-- ... JOIN [SHUFFLE] t2 ...
Where each Impala hint acts in the life of a querySQL texthints in commentsParserattaches hints to nodesJoin orderingSTRAIGHT_JOIN stops itDistributionBROADCAST / SHUFFLEInsert shapingSHUFFLE, CLUSTEREDSchedulerSCHEDULE_* replica choicefragmentsExecutorsscan, join, writeProfilecheck what ranQuery options act globally:DEFAULT_JOIN_DISTRIBUTION_MODE, BROADCAST_BYTES_LIMITStats feed the same planner steps. Fix stats first; a hint pins one decision for one query.
Each hint pins one planner or scheduler decision. Query options change the same decisions for a whole session.

Join distribution: BROADCAST versus SHUFFLE

A hash join has a build side (the right-hand input, loaded into a hash table) and a probe side (the left-hand input, streamed through it). With /* +BROADCAST */, the whole build side is sent to every executor that holds a part of the probe side. With /* +SHUFFLE */ (shown in EXPLAIN as a PARTITIONED join), both sides are hash-partitioned on the join key, so each executor builds a hash table for only its share.

The arithmetic decides which is right. Broadcast costs build-size times executors in network and memory, but never moves the probe side. Shuffle costs build-size plus probe-size in network, once, and divides build memory by the number of executors. With 20 executors and a 300 MB dimension against a 2 TB fact, broadcast moves 6 GB and holds 300 MB per node; shuffle would move over 2 TB. With a 40 GB build side, broadcast needs 40 GB of hash table on every node and will spill or fail against its memory limit; shuffle needs about 2 GB each.

Broadcast is the default when statistics are unavailable, because without row counts the planner cannot tell that the right side is large. That is the classic reason a hint is needed, and also the reason it usually should not be: compute stats and the planner will pick shuffle by itself. See Impala statistics and metadata.

Advertisement

STRAIGHT_JOIN: owning the join order

Impala's join-order optimisation puts the largest input on the probe side and joins the most selective tables early, again from statistics. If you add a join hint but the planner then swaps the two tables, your SHUFFLE now applies to a different build side, or the hint lands on a join that no longer exists in that shape. SELECT STRAIGHT_JOIN tells the planner to join in the order the tables appear in FROM, which is what makes join hints dependable. The documentation pairs them for that reason.

The rule when you take control: list the largest table first, then join the remaining tables in decreasing order of selectivity, so intermediate results shrink as early as possible. STRAIGHT_JOIN applies to the whole query block, so every join in it is now your responsibility. On a ten-table query that is a lot of decisions to freeze forever.

INSERT hints: SHUFFLE, NOSHUFFLE, CLUSTERED, NOCLUSTERED

Writing to a partitioned table is where hints still earn their keep. Each executor that receives rows for a partition opens a writer for that partition, and a Parquet writer buffers data in memory until it has a row group's worth. If every executor sees rows for every partition, a 365-partition insert on 20 executors opens up to 7,300 writers, each with its own buffer and each producing a small file.

/* +SHUFFLE */ adds an exchange that hash-partitions rows on the partition columns before the write, so each partition is written by exactly one executor: far fewer files and far less memory per node, at the cost of moving the data once. /* +NOSHUFFLE */ removes that exchange, which is right when the input is small or already partitioned appropriately and the extra network hop is pure cost.

/* +CLUSTERED */ adds a sort on the partition columns before the writer, so each executor writes one partition at a time and keeps only one writer open. It became available in Impala 2.8 and has been the default for HDFS tables since Impala 3.0, so on current versions you mostly meet it as /* +NOCLUSTERED */, which removes the sort. For Kudu tables the documentation recommends NOCLUSTERED, and since 2.10 combining NOCLUSTERED with NOSHUFFLE disables the exchange and sort that Kudu inserts otherwise add automatically.

-- Without hints on an older cluster, abridged:
WRITE TO HDFS [sales.sales_by_day, OVERWRITE=false, PARTITION-KEYS=(sale_date)]
00:SCAN HDFS [sales.staging_sales]

-- With /* +SHUFFLE,CLUSTERED */, abridged:
WRITE TO HDFS [sales.sales_by_day, OVERWRITE=false, PARTITION-KEYS=(sale_date)]
02:SORT
|  order by: sale_date ASC NULLS LAST
01:EXCHANGE [HASH(sale_date)]
00:SCAN HDFS [sales.staging_sales]

Scan scheduling hints

Each HDFS block (or Kudu tablet) has several replicas, and the scheduler picks which executor scans which replica. SCHEDULE_CACHE_LOCAL, SCHEDULE_DISK_LOCAL and SCHEDULE_REMOTE set the scheduler's replica-locality preference for that table; the documentation defines them as equivalent to the REPLICA_PREFERENCE query option set to CACHE_LOCAL, DISK_LOCAL or REMOTE, scoped to one table reference. RANDOM_REPLICA is the per-table form of the SCHEDULE_RANDOM_REPLICA option: it turns on random tie-breaking between equally good replicas, and can be combined with the others.

The use case is narrow: a small, hot table that many concurrent queries scan, where the deterministic choice keeps landing on the same node and makes it a hotspot. Spreading those scans with RANDOM_REPLICA flattens CPU across the cluster. On object storage such as S3 there is no locality to exploit, so these hints mostly do not apply; see Impala memory limits for the more usual reason one executor is overloaded.

Query options: the better hammer for most problems

Several options change the same decisions for a whole session or resource pool without touching SQL. DEFAULT_JOIN_DISTRIBUTION_MODE (BROADCAST, the default, or SHUFFLE) only matters when a joined table is missing statistics; setting it to SHUFFLE makes the unknown case safe instead of fast-or-fatal. BROADCAST_BYTES_LIMIT caps the estimated size of a broadcast input, defaulting to 32 GB, with 0 disabling it; above the cap the planner picks a partitioned join. Runtime filters, which are often mistaken for something you hint, are controlled by options such as RUNTIME_FILTER_MODE; see runtime filters.

-- Session-wide alternatives to per-join hints
SET DEFAULT_JOIN_DISTRIBUTION_MODE=SHUFFLE;  -- only applies when a joined table lacks stats
SET BROADCAST_BYTES_LIMIT=2147483648;        -- 2 GB: estimated build sides above this are shuffled
SET RUNTIME_FILTER_MODE=GLOBAL;              -- runtime filters are options, not hints

-- The real fix most of the time
COMPUTE STATS sales.customer_dim;
COMPUTE INCREMENTAL STATS sales.sales_fact PARTITION (sale_date='2026-10-02');

Options can be set as defaults per admission-control pool, so a pool that serves ad hoc analysts can get conservative join settings without anyone editing their queries.

Worked example: a daily partitioned backfill

A team backfills a year of sales into sales_by_day, partitioned by sale_date, on a 20-executor cluster running Impala 2.8 or later but before 3.0, where clustering is not the default. The job takes 70 minutes, two executors hit their memory limit, and the table ends with about 60,000 Parquet files, most a few megabytes.

EXPLAIN shows a write fed directly by the scan: no exchange, no sort. Every executor scans part of the staging table, sees rows for most of the 365 dates, and opens a writer for each. Adding /* +SHUFFLE,CLUSTERED */ produces the second plan above: an exchange hashed on sale_date routes each date to one executor, and the sort lets it write dates one after another with one open writer.

After the change the job moves the staging data across the network once, but writes about 365 files of sensible size, memory per executor drops to a single writer buffer plus the sort, and the run finishes in well under half the time. The file count also speeds every later query on the table, because fewer files means fewer scan ranges and less metadata. On Impala 3.0 or later the same insert already gets the sort by default, and the remaining question is only whether the shuffle is worth it for this input.

Failure modes

  • A hint that never applied. Wrong position, deprecated syntax removed by an upgrade, or a client that strips comments before sending SQL. Check the plan, not the SQL.
  • A hint that outlived its reason. A BROADCAST added when a dimension had 10,000 rows still broadcasts it at 50 million. Every hint should carry a comment with the date, the table sizes and the ticket that justified it.
  • STRAIGHT_JOIN freezing a bad order. Data distributions shift, and a pinned order cannot adapt. Review pinned queries when table sizes change by an order of magnitude.
  • Hints masking missing stats. If a hint fixes one query, every other query on that table is still planned blind. Compute stats and remove the hint.
  • NOSHUFFLE on large inserts. Small files and writer memory blow-ups return. Use it only when the input is small or already partitioned.
  • Concurrency effects. A forced broadcast multiplies build memory by executor count, which reduces how many queries admission control can admit.

What to do next

  1. Run COMPUTE STATS (or incremental stats per partition) on every table in the slow query before considering a hint.
  2. Read EXPLAIN and the profile summary to find the one decision that is wrong: order, distribution, insert shape or scan placement.
  3. Try a session query option (DEFAULT_JOIN_DISTRIBUTION_MODE, BROADCAST_BYTES_LIMIT) before editing SQL.
  4. If you must hint, use the block-comment form in the documented position, pair join hints with STRAIGHT_JOIN, and confirm the change in EXPLAIN.
  5. For partitioned inserts, check for an exchange on the partition columns and a sort; add SHUFFLE when many executors write many partitions.
  6. Comment every hint with date, table sizes and reason, and grep for hints after each upgrade or large data change.
Key takeaway: Impala hints pin individual planner decisions: join distribution with BROADCAST or SHUFFLE, join order with STRAIGHT_JOIN, insert shape with SHUFFLE, NOSHUFFLE, CLUSTERED and NOCLUSTERED, and replica choice with the SCHEDULE hints. They are silently ignored when misplaced and silently wrong when data changes. Fix statistics first, prefer session or pool query options, use hints mainly for partitioned inserts and proven planner mistakes, and verify every hint in EXPLAIN.