Troubleshooting: Expensive queries

View as Markdown

This guide helps you find queries that use a lot of CPU or memory on a cluster, and how to reduce their cost.

A query that can’t be answered from an existing index makes the cluster build a temporary dataflow, compute the result, and then drop the dataflow. The cluster does this work on every execution, alongside the work of maintaining its indexes and materialized views. Expensive queries slow down other queries on the same cluster, and can make the cluster run out of memory.

Common causes

  • No usable index: Frequent queries that join, aggregate, or filter on columns without an index build a dataflow on every execution.
  • Ad-hoc queries on a production cluster: Exploratory queries, such as large joins or full scans, compete for CPU and memory with the indexes and materialized views that serve your application.
  • Large results: Queries without filters or limits read and return much more data than needed.

Diagnosing the issue

Find the most expensive queries

Queries with the standard execution strategy built a temporary dataflow. To find the ones that consumed the most time over the last day, grouped by query text, query the statement log:

SELECT
  left(sql, 60) AS sql,
  cluster_name,
  count(*) AS executions,
  round(sum(extract(epoch FROM finished_at - began_at)), 3) AS total_seconds,
  max(result_size) AS max_result_bytes
FROM mz_internal.mz_recent_activity_log
WHERE execution_strategy = 'standard'
  AND cluster_name NOT LIKE 'mz_%'
  AND began_at > now() - INTERVAL '1 day'
GROUP BY sql_hash, sql, cluster_name
ORDER BY total_seconds DESC
LIMIT 10;
                             sql                              | cluster_name | executions | total_seconds | max_result_bytes
--------------------------------------------------------------+--------------+------------+---------------+------------------
 SELECT o.customer_id, t.total FROM orders o JOIN order_total | quickstart   |          1 |         0.081 |              234
 SELECT count(*) FROM orders                                  | quickstart   |          1 |         0.026 |               20
 SELECT * FROM order_counts WHERE customer_id = 5             | quickstart   |          1 |         0.015 |               21

Look for two patterns:

  • A query with many executions runs often enough that its total cost adds up, even if each execution is fast. Make it a fast path query.
  • A query with a high total_seconds for few executions is an expensive ad-hoc query. Isolate it or reduce the data it reads.

Canceled and failed statements have no execution_strategy, so this query doesn’t include them. To find long-running statements that were canceled or failed, filter on finished_status IN ('canceled', 'error') instead.

Find expensive queries that are running now

Dataflows for queries that are running are currently named oneshot-select-<id>. This name is not a stable interface and can change between releases. To see how much CPU time each one has used, run the following on the cluster that runs the queries:

SET cluster = quickstart;

SELECT
  mdo.name,
  mse.elapsed_ns / 1000 * '1 MICROSECONDS'::interval AS elapsed_time
FROM mz_introspection.mz_scheduling_elapsed AS mse,
  mz_introspection.mz_dataflow_operators AS mdo,
  mz_introspection.mz_dataflow_addresses AS mda
WHERE mse.id = mdo.id
  AND mdo.id = mda.id
  AND list_length(mda.address) = 1
  AND mdo.name LIKE 'Dataflow: oneshot-select-%'
ORDER BY mse.elapsed_ns DESC;
             name              |  elapsed_time
-------------------------------+-----------------
 Dataflow: oneshot-select-t176 | 00:00:35.980703
 Dataflow: oneshot-select-t190 | 00:00:14.433187

To find the query text and cancel a running query, see Find running queries.

Compare with the cost of indexes and materialized views

To see how the CPU and memory of your indexes and materialized views compare, run EXPLAIN ANALYZE CLUSTER on the cluster:

EXPLAIN ANALYZE CLUSTER CPU, MEMORY;

If the indexes and materialized views account for most of the cluster’s CPU and memory, the queries are not the main cost. To see which operators of an index or materialized view use the most resources, see EXPLAIN ANALYZE and Dataflow troubleshooting.

Check the query plan

Run EXPLAIN on an expensive query to see what the dataflow does:

EXPLAIN SELECT o.customer_id, t.total
FROM orders o
JOIN order_totals t USING (customer_id)
WHERE o.id < 10;

Look for operators that read full collections, such as joins without a matching index, or Read operators on large sources or materialized views.

Resolution

Make frequent queries fast path

A query that reads from an index and applies only filters and projections doesn’t build a dataflow. EXPLAIN shows Explained Query (fast path) for these queries.

  • Create an index on the columns that the query filters on.

  • Move joins and aggregations into a view, and index the view on the lookup key. For example:

    CREATE VIEW order_totals AS
      SELECT customer_id, sum(amount) AS total
      FROM orders
      GROUP BY customer_id;
    
    CREATE INDEX order_totals_idx ON order_totals (customer_id);
    
    -- Fast path lookup
    SELECT * FROM order_totals WHERE customer_id = 5;
    
  • Run the query on the same cluster as the index. Indexes are local to a cluster.

For more techniques, see Optimization.

Isolate ad-hoc queries

Run exploratory queries on a separate cluster, so they can’t slow down or crash the cluster that serves your application:

CREATE CLUSTER adhoc (SIZE = '25cc');
SET cluster = adhoc;

For guidance on how to split work across clusters, see Operational guidelines.

Return less data

  • Add filters, and use temporal filters on timestamp columns. Materialize can skip over old data in storage that doesn’t match the filter.
  • Add a LIMIT clause to exploratory queries.
  • Select only the columns you need.

The max_query_result_size configuration parameter makes a query fail if its result exceeds the limit. It doesn’t limit the memory that the temporary dataflow uses to compute the result.

Size up the cluster

If the queries are already optimized, size up the cluster to give it more CPU and memory.

Back to top ↑