# Introduction to optimizing query performance

This document provides an overview of optimization techniques that can improve
query performance in BigQuery. In general, queries that do less work
perform better. They run faster and consume fewer resources, which can result in
lower costs and fewer failures.

## Query performance

Evaluating query performance in BigQuery involves several factors:

- [Input data and data sources (I/O)](https://docs.cloud.google.com/bigquery/docs/best-practices-performance-compute#reduce-data-processed): How many bytes does your query read?
- [Communication between nodes (shuffling)](https://docs.cloud.google.com/bigquery/docs/best-practices-performance-compute#reduce-data-processed): How many bytes does your query pass to the next stage? How many bytes does your query pass to each slot?
- [Computation](https://docs.cloud.google.com/bigquery/docs/best-practices-performance-compute#optimize-query-operations): How much CPU work does your query require?
- [Outputs (materialization)](https://docs.cloud.google.com/bigquery/docs/best-practices-performance-compute#reduce-query-output): How many bytes does your query write?
- [Capacity and concurrency](https://docs.cloud.google.com/bigquery/docs/best-practices-performance-overview#capacity-and-concurrency): How many slots are available and how many other queries are running at the same time?
- [Query patterns](https://docs.cloud.google.com/bigquery/docs/best-practices-performance-patterns): Are your queries following SQL best practices?

To evaluate specific queries or whether you are experiencing resource contention,
you can use [Cloud Monitoring](https://docs.cloud.google.com/bigquery/docs/monitoring) or the [BigQuery administrative resource charts](https://docs.cloud.google.com/bigquery/docs/admin-resource-charts) to monitor how your BigQuery jobs consume
resources over time. You can also use Gemini Cloud Assist to
[analyze your jobs](https://docs.cloud.google.com/bigquery/docs/use-cloud-assist#analyze_jobs).
If you identify a slow or resource-intensive query, you can
focus your performance optimizations on that query.

Some query patterns, especially those generated by business intelligence tools,
can be accelerated by using [BigQuery BI Engine](https://docs.cloud.google.com/bigquery/docs/bi-engine-intro).
BI Engine is a fast, in-memory analysis service that accelerates many
SQL queries in BigQuery by intelligently caching the data you use
most frequently. BI Engine is built into BigQuery,
which means you can often get better performance without any query modifications.

As with any systems, optimizing for performance sometimes involves tradeoffs.
For example, using advanced SQL syntax can sometimes introduce complexity and
reduce the queries' understandability for people who aren't SQL experts.
Spending time on micro-optimizations for noncritical workloads could also divert
resources away from building new features for your applications or from identifying
more important optimizations. To help you achieve the highest possible return on
investment, we recommend that you focus your optimizations on the workloads that
matter most to your data analytics pipelines.

## Optimization for capacity and concurrency

BigQuery offers two pricing models for queries:
[on-demand](https://cloud.google.com/bigquery/pricing#on_demand_pricing) pricing and
[capacity-based](https://cloud.google.com/bigquery/pricing#capacity_compute_analysis_pricing)
pricing. The on-demand model provides a shared pool of capacity, and pricing is
based on the amount of data that is processed by each query you run.

The [capacity-based](https://docs.cloud.google.com/bigquery/docs/reservations-workload-management) model is
recommended if you want to budget a consistent, monthly expenditure
or if you need more capacity than is available with the on-demand model.
When you use capacity-based pricing, you allocate dedicated query processing
capacity that is measured in [slots](https://docs.cloud.google.com/bigquery/docs/slots).
The cost of all bytes processed is included in the capacity-based price.
In addition to fixed [slot commitments](https://docs.cloud.google.com/bigquery/docs/reservations-workload-management#slot_commitments),
you can use [autoscaling slots](https://docs.cloud.google.com/bigquery/docs/slots-autoscaling-intro),
which provide dynamic capacity based on your query workload.

The performance of queries that are run repeatedly on the same data can vary,
and the variation is generally larger for queries using on-demand slots than
it is for queries using slot reservations.

During SQL query processing, BigQuery breaks down the
computational capacity required to execute each stage of a query into
slots. BigQuery automatically determines the
number of queries that can run concurrently as follows:

- On-demand model: number of slots available in the project
- Capacity-based model: number of slots available in the reservation

Queries that require more slots than are available are [queued](https://docs.cloud.google.com/bigquery/docs/query-queues)
until processing resources become available. After a query begins execution,
BigQuery calculates how many slots each query stage
uses based on the stage size and complexity and the number of slots
available. BigQuery uses a technique called [fair scheduling](https://docs.cloud.google.com/bigquery/docs/slots#fair_scheduling_in_bigquery)
to ensure that each query has enough capacity to progress.

Access to more slots doesn't always result in faster performance for a query.
However, a larger pool of slots can improve the performance of large or complex
queries, and the performance of highly concurrent workloads. To improve query
performance, you can [modify your slot reservations](https://docs.cloud.google.com/bigquery/docs/reservations-tasks)
or set a higher limit for [slots autoscaling](https://docs.cloud.google.com/bigquery/docs/slots-autoscaling-intro).
You can also use Gemini Cloud Assist to
[manage your reservations and capacity with natural language prompts](https://docs.cloud.google.com/bigquery/docs/use-cloud-assist#administer_bigquery).

## Query plan and execution graph

BigQuery generates a [query plan](https://docs.cloud.google.com/bigquery/query-plan-explanation) each time that you run a query.
Understanding this plan is critical for effective query optimization.
The query plan includes execution statistics such as bytes read and slot time
consumed. The query plan also includes details about the different stages of execution, which can help you
diagnose and improve query performance. The [query execution graph](https://docs.cloud.google.com/bigquery/docs/query-insights)
provides a graphical interface for viewing the query plan and diagnosing query performance issues.

You can also use the [`jobs.get` API method](https://docs.cloud.google.com/bigquery/docs/reference/rest/v2/jobs/get)
or the [`INFORMATION_SCHEMA.JOBS` view](https://docs.cloud.google.com/bigquery/docs/information-schema-jobs)
to retrieve the query plan and timeline information. This information is used by
[BigQuery Visualizer](https://github.com/GoogleCloudPlatform/professional-services/tree/master/tools/bq-visualizer),
an open source tool that visually represents the flow of execution stages in a
BigQuery job.

When BigQuery executes a query job, it converts the declarative
SQL statement into a graph of execution. This graph is broken up into a series
of query stages, which themselves are composed of more granular sets of
execution steps. BigQuery uses a heavily distributed parallel
architecture to run these queries. The BigQuery stages model the
units of work that
many potential workers might execute in parallel. Stages communicate with one
another through a
[fast, distributed shuffle architecture](https://cloud.google.com/blog/products/gcp/in-memory-query-execution-in-google-bigquery).

To estimate how computationally expensive a query is, you can look at the
total number of slot seconds the query consumes. The lower the number of slot
seconds, the better, because it means more resources are available to other
queries running in the same project at the same time.

The query execution graph can help you understand how
BigQuery executes queries and if certain stages dominate resource
utilization. For example, a `JOIN` stage that generates far more output rows than
input rows might indicate an opportunity to filter earlier in the query.
However, the managed nature of the service limits whether some details are
directly actionable. For best practices and techniques to improve query
execution and performance, see
[Optimize query computation](https://docs.cloud.google.com/bigquery/docs/best-practices-performance-compute).

## What's next

- Learn how to troubleshoot query execution issues using the [BigQuery audit logs](https://docs.cloud.google.com/bigquery/docs/reference/auditlogs).
- Learn other [cost-controlling techniques](https://docs.cloud.google.com/bigquery/docs/controlling-costs) for BigQuery.
- View near real-time metadata about BigQuery jobs using the [`INFORMATION_SCHEMA.JOBS` view](https://docs.cloud.google.com/bigquery/docs/information-schema-jobs).
- Learn how to monitor your BigQuery usage using the [BigQuery System Tables Reports](https://github.com/GoogleCloudPlatform/bigquery-utils/tree/master/dashboards/system_tables).