Most Impala cost work starts as a one-off clean-up: someone finds a dashboard that scans every partition, adds statistics, and the queue disappears. Six months later the bill is back, because new tables, teams and BI tools arrived unwatched. Tuning reduces cost once. Cost management keeps it down: a loop that meters every query, attributes it to an owner, compares spend with a budget, and changes behaviour or configuration when the two diverge.

This article builds that loop. It shows how to switch on Impala's own query log, how to turn raw query telemetry into a cost figure that survives argument, how to account for idle capacity, how to set budgets and alerts, how to enforce them without breaking legitimate work, and how to forecast when you genuinely need more nodes. The individual tuning levers, such as partition pruning, statistics and query limits, are covered in Impala cost optimization; here they appear only as actions the loop triggers.

The loop in one picture

The cost management loopQueriesBI, ETL, ad hocCoordinatorsadmission controlimpala_query_livein memory, recentimpala_query_logIceberg, persistentbatchCost model jobshares x time x priceOwner mappinguser and pool to teamShowback and budgetsdaily spend per teamAlerts and reviewsanomaly, forecastActionsfix query, clamp pool, scalechanges behaviour
Coordinators record every query; workload management moves telemetry from the live table to the persistent Iceberg log; a daily job turns it into cost per owner; budgets, alerts and reviews drive fixes, pool changes or capacity decisions, which change the next day's queries.

What you are actually managing

Impala has no price tag on a query. It has resources: executors with a fixed amount of memory that admission control may hand out, a number of admission slots, CPU cores, and storage behind them. Whatever you pay, whether a cloud rate per node-hour or the amortised cost of servers in a rack, you pay for that capacity whether queries use it or not. Cost management therefore has two jobs. The first is attribution: deciding which owner consumed which share of capacity. The second is utilisation: noticing how much capacity was paid for and used by nobody.

Both need durable data per query, joined to an owner. Profiles in the web UI diagnose one query well but are kept only for recent queries, which is useless for a monthly report.

Switching on the query log

Recent Impala releases include workload management, which writes telemetry for every completed query to system tables. It is switched on with the startup flag --enable_workload_management=true, set on every coordinator and on catalogd. Two tables appear in the sys database. sys.impala_query_live is an in-memory view of running and recently completed queries. sys.impala_query_log is an Apache Iceberg table that holds completed queries indefinitely. Coordinators batch rows into it every query_log_write_interval_s seconds (300 by default) or sooner when query_log_max_queued records (3,000 by default) are waiting. Set cluster_id to a distinct value per cluster so one shared log can tell clusters apart, and query_log_request_pool to send the log's own inserts to a maintenance pool.

The columns that matter for cost are resource_pool, db_user, start_time_utc, total_time_ms, event_completed_admission (milliseconds from start until admission finished), backends_count, cluster_memory_admitted, executor_slots, per_host_mem_estimate, pernode_peak_mem_max, bytes_read_total and compressed_bytes_spilled. The tables_queried and where_columns columns tell you which data each query touched and how it filtered it. Because the log is Iceberg, it needs maintenance like any other Iceberg table, and the full SQL and plan text, each allowed up to 16 MB by default, make it larger than people expect:

-- Weekly, from the maintenance pool. Check that your Impala version supports
-- these Iceberg statements before scheduling them.
ALTER TABLE sys.impala_query_log EXECUTE expire_snapshots(now() - interval 7 days);
OPTIMIZE TABLE sys.impala_query_log;

-- Startup flags worth lowering if nobody reads full plans from the log:
--   --query_log_max_plan_length=65536
--   --query_log_max_sql_length=262144

A cost model that survives argument

Wall time multiplied by hosts is the obvious cost model and it is wrong in a way that matters. Ten small queries running at once on the same twenty executors would each be charged the whole cluster, and the total would exceed what you paid. A defensible model charges each query for the share of capacity it reserved, because that share is what nobody else could use while it ran.

Admission control reserves two things per executor: memory, the query's memory limit or estimate, and slots, one per executor for an ordinary query or more when MT_DOP is raised. Each executor has a fixed amount of admittable memory and a fixed number of slots, which defaults to the number of cores. A query's share of one executor is the larger of its memory fraction and its slot fraction, the dominant share, because whichever runs out first is what blocks the next query. Multiply by executors and by time after admission, and you get node-equivalent hours, which sum to no more than the capacity you own. Queue time is excluded: a queued query reserves nothing.

-- Node-equivalent hours per owner per day. Set the two constants to your executors:
-- 128 GB admittable memory per executor, 32 admission slots per executor.
WITH q AS (
  SELECT to_date(start_time_utc)                                   AS day,
         resource_pool, db_user, backends_count,
         GREATEST(total_time_ms - event_completed_admission, 0) / 3600000.0 AS run_h,
         cluster_memory_admitted / NULLIF(backends_count, 0) / POWER(1024, 3) AS gb_per_host,
         executor_slots,
         bytes_read_total
  FROM sys.impala_query_log
  WHERE start_time_utc >= date_sub(now(), 35)
)
SELECT day, resource_pool, db_user,
       COUNT(*)                                                      AS queries,
       SUM(run_h * backends_count *
           GREATEST(gb_per_host / 128, executor_slots / 32))         AS node_eq_h,
       SUM(bytes_read_total) / POWER(1024, 4)                        AS tib_read
FROM q
GROUP BY day, resource_pool, db_user;

Two caveats. A query whose client fetches rows slowly keeps its memory reserved until the last row is fetched, so it is charged while its CPU idles, which usefully exposes slow BI drivers. And the log has no CPU-time column in the documented table, so the model measures reservation, not consumption. That is the right basis for a cluster whose size is set by peak reservations, and the over-reservation report below shows where reservation and consumption diverge.

Attribution and idle capacity

Database users are not owners. Map them to teams from your directory, and map service accounts explicitly, because a shared BI account can hide a dozen teams; where one is unavoidable, give each team its own resource pool and attribute by pool. Then decide what to do with capacity nobody reserved. Spreading idle cost in proportion to use makes the bill add up to what you paid and gives everybody a reason to care about utilisation:

# showback.py - daily owner costs from the node-equivalent rollup, idle spread by share.
import csv, collections

NODE_HOUR_PRICE = 3.10      # blended executor cost per node-hour (cloud rate or amortised hardware)
NODES, HOURS = 20, 24

team_of = {r["db_user"]: r["team"] for r in csv.DictReader(open("user_team.csv"))}
used = collections.Counter()
for r in csv.DictReader(open("rollup_yesterday.csv")):      # output of the SQL above
    used[team_of.get(r["db_user"], "UNOWNED")] += float(r["node_eq_h"])

capacity = NODES * HOURS
attributed = sum(used.values())
idle = max(capacity - attributed, 0.0)
print(f"utilisation {attributed / capacity:.0%}, idle {idle:.0f} node-h")
for team, nh in used.most_common():
    direct = nh * NODE_HOUR_PRICE
    spread = idle * (nh / attributed) * NODE_HOUR_PRICE
    print(f"{team:16} {nh:7.1f} node-h  direct ${direct:8.2f}  idle ${spread:8.2f}")

Budgets and alerts

A budget turns a report into a control. Give each team a monthly figure in node-equivalent hours or money, agreed with whoever pays. Alert on two signals: a day that is abnormal compared with the team's own history, which catches a broken job or a new dashboard on its first day, and a month-to-date forecast that crosses the budget, which catches slow drift. A robust anomaly test uses the median and the median absolute deviation, so one past spike does not mask the next:

# spend_alerts.py - run daily after showback; daily is the last 35 daily costs, oldest first.
import statistics

def alerts(team, daily, budget, day_of_month, days_in_month):
    out = []
    history, today = daily[:-1], daily[-1]
    med = statistics.median(history)
    mad = statistics.median(abs(x - med) for x in history) or 1.0
    if (today - med) / (1.4826 * mad) > 4:
        out.append(f"{team}: ${today:,.0f} yesterday, typical ${med:,.0f}")
    mtd = sum(daily[-day_of_month:])
    forecast = mtd / day_of_month * days_in_month
    if forecast > budget:
        out.append(f"{team}: forecast ${forecast:,.0f} against budget ${budget:,.0f}")
    return out

Send alerts to the owning team with the top three query fingerprints behind the change; an alert without a cause is ignored within a week.

Finding what to fix

The log can tell owners what to fix. Group queries by fingerprint, the SQL text with literals replaced, and rank fingerprints by total node-equivalent hours, not by the duration of a single run. A three-second tile that runs thousands of times a day usually outranks a monthly report:

import re

def fingerprint(sql):
    s = re.sub(r"--[^\n]*", " ", sql)
    s = re.sub(r"'(?:[^']|'')*'", "?", s)          # string literals
    s = re.sub(r"\b\d+(?:\.\d+)?\b", "?", s)        # numbers
    s = re.sub(r"\s+", " ", s).strip().lower()
    return re.sub(r"in \(\s*\?(?:\s*,\s*\?)*\s*\)", "in (?)", s)   # IN lists of any length

Then add three diagnostic columns per fingerprint. Over-reservation is per_host_mem_estimate / pernode_peak_mem_max: a ratio of ten means the query reserves ten times what it uses, so the pool admits a tenth of the concurrency it could, usually because statistics are missing or no MEM_LIMIT is set. Spill, any non-zero compressed_bytes_spilled, means the reservation is too small and the query pays in disk time. Pruning is visible in where_columns: a query on a date-partitioned table whose filter columns do not include the partition column almost certainly reads every partition.

An enforcement ladder

Enforcement should escalate, because the first response to an overspend should be information, not a cancelled query. A ladder that works in practice:

  1. Inform. Daily showback and weekly top fingerprints go to every team, whether or not they are over budget.
  2. Nudge. Anomaly and forecast alerts go to the owning team with causes attached, and a cost review is booked if the forecast stays over budget for a week.
  3. Limit. The team's pool gets tighter defaults: a smaller maximum memory per query with clamping, scan-byte and execution-time limits, and a lower maximum number of running queries. These settings are covered in admission control and memory limits.
  4. Isolate. Heavy or experimental work moves to its own pool or executor group, with a fixed memory share, so it can only slow itself down.
  5. Fund or stop. If the workload is legitimate and still over budget, the budget changes or capacity is bought; if not, the job is switched off.

Publish the ladder early: surprise limits get routed around with new service accounts, which destroys attribution.

Worked example: one tile, $2,600 a month

A cluster has 20 executors with 128 GB of admittable memory and 32 slots each, at a blended $3.10 per node-hour: 480 node-hours, or $1,488, a day. The first week of metering shows 300 node-equivalent hours reserved per day, 62.5 percent utilisation. The BI team owns 130 of them, and one fingerprint, a sales tile run about 6,000 times a day, accounts for 28.

Its numbers explain why. Statistics are missing, so the estimate is 18 GB per host, while pernode_peak_mem_max is about 700 MB, an over-reservation of roughly 25. Its memory share is 18/128, about 0.14, far above its slot share of 1/32. Each run lasts 6 seconds on 20 executors: 6/3600 x 20 x 0.14 x 6,000 is about 28 node-equivalent hours, or $87 a day, around $2,600 a month for one tile. Its where_columns list event_ts but not the partition column dt.

The fix is the usual one: COMPUTE STATS, a pool MEM_LIMIT of 2 GB, and a redundant partition predicate in the dashboard's template. Runs drop to 3 seconds and the dominant share becomes the slot share, 1/32: 3/3600 x 20 x 0.031 x 6,000 is about 3.1 node-equivalent hours, $9.70 a day. More important for the platform, freeing roughly 25 node-equivalent hours of reservation lets the BI pool's queue disappear, which was the reason the team had asked for five more nodes.

Forecasting capacity

Average utilisation is the wrong signal for buying capacity, because queries queue at the peak, not on average. Compute reserved node-equivalents per hour from the log, take the 95th percentile of the busiest hours each week, and plot it against capacity together with queue time, the gap between start_time_utc and admission. Scale only when the peak percentile nears capacity and queue time breaks a latency target. If peaks are short, moving batch work into troughs or autoscaling executor groups beats buying nodes that idle most of the day.

Failure modes

  • The log is not maintained. Snapshots and small files accumulate, queries against the log slow down and the log itself becomes a cost. Schedule expiry and optimisation from day one.
  • Shared accounts. One BI service user holds 60 percent of spend and no team owns it. Route by pool or require per-team connections.
  • Charging wall time times hosts. Totals exceed the bill, owners stop trusting the numbers. Charge reserved shares.
  • Budgets without causes. Alerts say spend is up but not why, so they are muted. Always attach fingerprints.
  • Enforcement first. Surprise limits cancel a finance close job; teams create new accounts to escape them. Escalate in published steps.

Trade-offs

ChoiceGainsCosts
Showback onlyNo friction, builds trustRelies on goodwill
Chargeback to budgetsReal incentive to fixDisputes about the model
Idle cost spread by shareBill adds up, utilisation mattersSmall teams pay for others' peaks
Idle cost kept centralSimple, fair to light usersNobody owns over-provisioning
Reservation-based modelMatches what sizes the clusterIgnores CPU actually burned
Hard pool limitsBounded spendCancelled legitimate work

What to do next

  1. Enable workload management on every coordinator and catalogd, set cluster_id and query_log_request_pool, and schedule snapshot expiry and optimisation of the log.
  2. Build the user-to-team mapping, including every service account, and give shared BI tools one pool per team.
  3. Run the node-equivalent rollup daily and publish showback, with utilisation and an UNOWNED row.
  4. Agree monthly budgets with each team and wire the anomaly and forecast alerts to the owners, with top fingerprints attached.
  5. Rank fingerprints weekly by total cost and add over-reservation, spill and pruning columns; fix the top five.
  6. Publish the enforcement ladder before applying any new pool limit.
  7. Track the weekly peak-hour percentile of reserved capacity and queue time before approving any request for more nodes.
Key takeaway: Impala cost management is a loop, not a tuning pass. Turn on workload management so every query lands in sys.impala_query_log, charge each query for its dominant share of admitted memory and slots after admission, map users to owning teams, spread idle capacity so the bill adds up, alert on anomalies and budget forecasts with causes attached, fix the most expensive fingerprints first, escalate enforcement in published steps, and buy capacity only when peak reservations and queue time demand it.