Troubleshooting: Slow queries
View as MarkdownThis guide helps you find out why a query takes longer than expected to return results, and how to fix it.
Common causes
- Lagging dependencies: A materialized view, index, or source that the query reads from is behind. Materialize waits for it to catch up before it returns a consistent result.
- No usable index: The query can’t be answered from an existing index, so Materialize builds a temporary dataflow or reads from storage for every execution.
- Busy cluster: Other work on the cluster, such as maintaining indexes and materialized views or running other queries, is using the CPU.
- Transactions: All statements in a transaction run at the same timestamp, so a fast object can wait for a slower one.
- Large results or client distance: Transmitting a large result, or a long network path between the client and Materialize, adds latency after the query has executed.
Diagnosing the issue
Find slow queries
List the slowest queries of the last hour from the statement log:
SELECT
left(sql, 60) AS sql,
cluster_name,
execution_strategy,
finished_at - began_at AS duration
FROM mz_internal.mz_recent_activity_log
WHERE statement_type = 'select'
AND finished_status = 'success'
AND cluster_name NOT LIKE 'mz_%'
AND began_at > now() - INTERVAL '1 hour'
ORDER BY duration DESC
LIMIT 10;
sql | cluster_name | execution_strategy | duration
--------------------------------------------------------------+--------------+--------------------+--------------
SELECT o.customer_id, t.total FROM orders o JOIN order_total | quickstart | standard | 00:00:00.081
SELECT count(*) FROM orders | quickstart | standard | 00:00:00.026
SELECT * FROM order_counts WHERE customer_id = 5 | quickstart | standard | 00:00:00.015
SELECT * FROM order_totals WHERE customer_id = 5 | quickstart | fast-path | 00:00:00.005
The execution_strategy column shows how Materialize executed the query:
| Strategy | Meaning |
|---|---|
fast-path |
The cluster read the result directly from an existing index, without building a dataflow. |
persist-fast-path |
The cluster read the result directly from storage, without building a dataflow. Materialize uses this for small LIMIT queries. |
standard |
The cluster built a temporary dataflow to compute the result, then dropped it. |
constant |
Materialize computed the result without a cluster. |
The query only lists successful statements. Canceled and failed statements
have no execution_strategy. To include them, remove the finished_status
filter.
The Query history tab in the Materialize console shows the same information, and lets you filter and sort statements by duration.
Break down where the time goes
mz_statement_lifecycle_history
records when each query reaches each stage of its lifecycle. Use it to split a
query’s duration into stages:
WITH events AS (
SELECT
statement_id,
max(occurred_at) FILTER (WHERE event_type = 'execution-began') AS began,
max(occurred_at) FILTER (WHERE event_type = 'optimization-finished') AS optimized,
max(occurred_at) FILTER (WHERE event_type = 'storage-dependencies-finished') AS storage_ready,
max(occurred_at) FILTER (WHERE event_type = 'compute-dependencies-finished') AS compute_ready,
max(occurred_at) FILTER (WHERE event_type = 'execution-finished') AS finished
FROM mz_internal.mz_statement_lifecycle_history
WHERE occurred_at > now() - INTERVAL '1 hour'
GROUP BY statement_id
)
SELECT
left(a.sql, 40) AS sql,
a.execution_strategy,
e.optimized - e.began AS optimization,
greatest(e.storage_ready, e.compute_ready) - e.optimized AS dependency_wait,
e.finished - greatest(e.storage_ready, e.compute_ready) AS execution,
e.finished - e.began AS total
FROM mz_internal.mz_recent_activity_log AS a
JOIN events AS e ON e.statement_id = a.execution_id
WHERE a.statement_type = 'select'
AND a.cluster_name NOT LIKE 'mz_%'
AND a.began_at > now() - INTERVAL '1 hour'
ORDER BY total DESC
LIMIT 10;
sql | execution_strategy | optimization | dependency_wait | execution | total
------------------------------------------+--------------------+--------------+-----------------+--------------+--------------
SELECT o.customer_id, t.total FROM order | standard | 00:00:00.006 | 00:00:00.001 | 00:00:00.074 | 00:00:00.081
SELECT count(*) FROM orders | standard | 00:00:00.007 | 00:00:00 | 00:00:00.019 | 00:00:00.026
SELECT * FROM order_counts WHERE custome | standard | 00:00:00.003 | 00:00:00 | 00:00:00.012 | 00:00:00.015
SELECT * FROM order_totals WHERE custome | fast-path | 00:00:00.003 | 00:00:00.001 | 00:00:00.001 | 00:00:00.005
Use the stage that dominates to pick a resolution:
| Stage | What happens | Resolution |
|---|---|---|
optimization |
Materialize parses, plans, and optimizes the query, and picks a timestamp. With real-time recency enabled, this includes waiting for the latest upstream offsets. | Simplify the query, or index a view that performs the complex part. In self-managed deployments, see Check environmentd. |
dependency_wait |
Materialize waits for the sources, tables, materialized views, and indexes the query reads from to catch up to the chosen timestamp. | Fix lagging dependencies. |
execution |
The cluster computes the result and returns it to Materialize. Sending the rows to the client is not included. | Use an index, reduce cluster load, or return less data. |
If total is small but your client reports a much higher latency, the time is
spent outside of Materialize. See Reduce client-side
latency.
Check for lagging dependencies
To see how far behind each of your objects is, query
mz_wallclock_global_lag:
SELECT o.name, o.type, l.lag
FROM mz_internal.mz_wallclock_global_lag AS l
JOIN mz_catalog.mz_objects AS o ON o.id = l.object_id
WHERE o.id LIKE 'u%'
ORDER BY l.lag DESC NULLS LAST
LIMIT 10;
A lag of a few seconds is expected. Lag that is much higher, or that keeps growing, means the object can’t keep up with its inputs. To see lag visually, open the object’s workflow graph in the console: click Clusters, select the cluster, select the object under Materialized Views or Indexes, then open the Workflow tab.
To find why an object is lagging, see Freshness troubleshooting and Dataflow troubleshooting. For a lagging source, see Troubleshoot ingestion.
Check the query plan
Run EXPLAIN on the query, on the cluster where the
query runs:
EXPLAIN SELECT * FROM order_totals WHERE customer_id = 5;
Explained Query (fast path):
→Map/Filter/Project
Project: #0, #1
→Index Lookup on materialize.public.order_totals (using materialize.public.order_totals_idx)
Lookup values: (5)
Used Indexes:
- materialize.public.order_totals_idx (lookup)
Target cluster: quickstart
Explained Query (fast path) means that no dataflow is built: an Index Lookup or Indexed operator reads from an index, and ReadStorage reads from
storage. A plan without (fast path) means Materialize builds a temporary
dataflow on every execution, even if the plan lists Used Indexes.
Check cluster utilization
A cluster near 100% CPU delays every query that runs on it:
SELECT c.name AS cluster, r.name AS replica, u.process_id, u.cpu_percent, u.memory_percent
FROM mz_internal.mz_cluster_replica_utilization AS u
JOIN mz_catalog.mz_cluster_replicas AS r ON r.id = u.replica_id
JOIN mz_catalog.mz_clusters AS c ON c.id = r.cluster_id
WHERE c.name = 'quickstart';
You can also see CPU and memory for each cluster under Clusters in the console.
Resolution
Fix lagging dependencies
- To find and fix the cause of lag in a materialized view or index, see Freshness troubleshooting and Dataflow troubleshooting. For a lagging source, see Troubleshoot ingestion.
- Avoid chaining materialized views where you don’t need to. Each materialized view in a chain adds a small amount of lag to the next one.
- If you don’t need strict serializable
results, use the
serializableisolation level. Materialize can then serve results at the latest timestamp that all inputs can already serve, instead of waiting for the lagging object. The whole result is then as stale as the most lagging input. - If you run queries inside a transaction, see Avoid transactions.
- If the dependencies can’t keep up with their inputs, size up the cluster that maintains them.
Use an index
- If the objects the query reads from don’t have an index, create one on the key the query filters or joins on. See Optimization.
- Run the query on the same cluster as the index. Indexes are local to a cluster.
- Make the index key match how you query the data. For example, a filter on
customer_idneeds an index oncustomer_id. - Move joins and aggregations out of the query and into a view, then index the view. The query becomes a lookup against an existing index.
Reduce cluster load
- Move queries that build dataflows onto a separate cluster, so that they don’t compete with indexes and materialized views for CPU. Indexes are local to a cluster, so a query that uses an index on the original cluster can become more expensive on the new one. See Expensive queries.
- Size up the cluster.
Return less data
- Filter results with temporal filters. Materialize can skip over old data in storage that doesn’t match the filter.
- Add a
LIMITclause to exploratory queries. A query that selects from a single source, table, or materialized view with only simple filters and projections, no ordering, and aLIMITplusOFFSETbelow 25 reads directly from storage.EXPLAINshowsExplained Query (fast path)for these queries. - Select only the columns you need. A large
result_sizeinmz_recent_activity_logadds time to transmit the result.
Avoid unnecessary transactions
All statements in a transaction run at the same timestamp, and that timestamp must be valid for every object the transaction may access. As a result, a query against a fast object can wait for a slower object in the same schema.
- Don’t use transactions for single statements.
- Check whether your SQL library or ORM wraps every query in a transaction, and disable that behavior.
Reduce client-side latency
Run clients in the same cloud region as your Materialize region. For example,
if your Materialize region is in AWS us-east-1, run your client in AWS
us-east-1. To reuse connections instead of opening a new one for each query,
see Connection pooling.
Self-managed deployments
In self-managed deployments, you also operate the components between the client
and the cluster: balancerd, which terminates TLS and proxies connections, and
environmentd, which plans every query and dispatches it to a cluster. Either
can add latency that doesn’t show up in cluster metrics.
Make sure statement logging is enabled
The steps above need the statement log. Operator Helm chart versions earlier than v26.40.0 disable statement logging. To check:
SHOW statement_logging_max_sample_rate;
If the result is 0, the statement log is empty. To enable statement logging
and understand its cost, see Query
History.
Attribute latency to a component
Compare three measurements over the same time window. Each one covers a smaller part of the query path, so the gap between two of them tells you where the time goes.
| Measurement | Covers |
|---|---|
| Latency reported by your client | The full round trip: client, network, balancerd, environmentd, and the cluster. |
finished_at - began_at in mz_recent_activity_log |
environmentd and the cluster. |
mz_compute_peek_duration_seconds |
From when environmentd sends the query to the cluster until the result arrives. This includes waiting for dependencies and, for standard queries, building the temporary dataflow. |
To compute the average statement log latency over the last minute:
SELECT
count(*) AS statements,
round(avg(extract(epoch FROM finished_at - began_at)) * 1000, 1) AS avg_ms,
round(max(extract(epoch FROM finished_at - began_at)) * 1000, 1) AS max_ms
FROM mz_internal.mz_recent_activity_log
WHERE began_at > now() - INTERVAL '1 minute'
AND statement_type = 'select'
AND finished_status = 'success'
AND cluster_name NOT LIKE 'mz_%';
Then compare:
- Client latency is much higher than the statement log: The time is spent
in the client, the network, or
balancerd. See Checkbalancerdand Check the client. - Statement log latency is much higher than peek duration: The time is
spent in
environmentdbefore the query reaches the cluster, for example in optimization or timestamp selection. See Checkenvironmentd. - Peek duration is high: Use the lifecycle breakdown to see whether the query waits for dependencies or for execution on the cluster, then follow the steps in Resolution.
mz_compute_peek_duration_seconds has an instance_id label that holds the
cluster ID. Filter on it to compare the same clusters as the statement log
query, which excludes system clusters.
To collect mz_compute_peek_duration_seconds and other Prometheus metrics, see
Grafana.
Check environmentd
Every query passes through environmentd. Watch the CPU of the environmentd
pod while queries are slow. If it is saturated, or throttled by a CPU limit,
while cluster CPU is low, give environmentd more CPU with
environmentdResourceRequirements in the Materialize custom resource. See
Materialize CRD field
descriptions.
Applying the change rolls out a new environmentd. With the v1alpha1 CRD,
also set requestRollout to a new UUID, or the operator does not roll out the
change. See Modifying the custom
resource.
Check balancerd
balancerd handles every byte of every connection.
- Don’t set CPU limits on
balancerdpods, and run at least two replicas. - Its memory use grows with the number of connections. Check the pods’ restart
count: a
balancerdpod that is killed for running out of memory causes connection errors and latency spikes.
container_cpu_cfs_throttled_periods_total metric from cAdvisor instead of CPU
usage.
Check the client
- Don’t set CPU limits on load-testing or client processes. A throttled client queues requests and reports the queueing time as query latency.
- Don’t measure through
kubectl port-forwardor SSH tunnels. They cap throughput far below what the deployment can serve. - Run the client in the same region and network as Materialize.