Skip to content
Pivot is in early development and is not production ready. See Project status.

Architecture

Pivot separates query execution, data access, and server metadata into distinct layers. This separation allows the same execution engine to operate across different data sources, including Iceberg and Delta Lake, and across deployment models ranging from a local shell to a distributed server cluster. It also allows multiple deployments to share the same data sources, metadata layer, or both.

These responsibilities are split across five core components:

  • Dispatch execution pool - Executes physical plans across a pool of workers responsible for computation and I/O.
  • Planner - Translates SQL queries into optimized physical execution plans.
  • Catalog - Connects the planner to named datastores and maintains the datastore snapshots used by each query.
  • Datastore - Manages schemas, tables, table versions, and the underlying data files.
  • Metastore - Manages server-level configuration, including users, authentication, and datastore connections.

The Dispatch execution engine is responsible for receiving a query’s physical execution plan and reliably executing it across a set of Dispatch workers: a collection of workers, one per CPU core, that perform the computation, networking, and I/O required to execute a query.

A physical execution plan runs across the Dispatch worker pool The query's physical execution plan is dispatched to a pool with one worker per CPU core. Workers 1, 2, and N each perform computation, networking, and other I/O. Physical execution plan operators + data flow Dispatch pool one worker per CPU core Worker 1 CPU core 1 Computation Networking + other I/O Worker 2 CPU core 2 Computation Networking + other I/O … Worker N CPU core N Computation Networking + other I/O

Once a query has been dispatched to the pool, each CPU worker is responsible for executing its share of the physical plan and scheduling operations to make efficient use of the CPU caches. If a worker has just produced an output that is still hot in cache, for example, it will prefer to continue with the next operator in the pipeline rather than switch to unrelated work whose data is not cache-resident.

Workers try to preserve this locality whenever possible, but not at the expense of parallelism. Cross-worker work stealing allows idle workers to pick up work from busier ones, keeping execution balanced and preventing a single worker from becoming a bottleneck while other CPU cores sit idle.

When multiple queries run concurrently, workers try to balance how CPU resources are shared between them. Scheduling favors cache locality when possible, while ensuring that no single query monopolizes the available CPU time or causes other queries to starve.

Beyond CPU work, the Dispatch worker also handles IO requests and prioritizes between them. For example, a “Materialize” operator that enriches existing data with additional columns, takes precedence over the input operator feeding it, to prevent a scenario where the channel between the two fills up and takes all available memory.

Prioritizing materialization drains buffered input Input reads produce record batches that wait in a channel for materialization to fetch additional columns. Prioritizing materialization drains the channel and releases memory held by waiting batches. Input reads new data Waiting batches Materialize fetch columns

Pivot’s planner uses a fork of DuckDB to parse SQL, resolve table and column references, and optimize the logical plan. Pivot then translates that plan into its own execution operators.

A query passes through four stages:

  1. Parsing - DuckDB’s PostgreSQL-derived parser checks the SQL syntax and builds an abstract syntax tree (AST) representing the statement’s expressions and clauses, such as projections, filters, joins, and ordering.
  2. Binding and logical planning - the binder resolves table and column references against the query’s catalog snapshot, resolves aliases and function calls, and checks expression types. It produces a logical plan describing the operations needed to answer the query.
  3. Logical optimization - the optimizer rewrites the plan to reduce unnecessary work. This includes evaluating constant expressions in advance, pushing filters closer to table scans, removing unused columns, choosing join order using available row-count estimates, and more.
  4. Physical translation and compilation - Pivot converts the optimized logical plan into its own operators and expressions, then applies additional refinements, such as pushing eligible limits into grouped aggregations. It selects execution implementations for scans, joins, and aggregations and connects the operators into a parallel dataflow ready to run on the Dispatch worker pool.
From SQL to a Pivot execution plan SQL passes through four stages in the query planner: parsing into an AST, binding and logical planning using the catalog snapshot, logical optimization, and physical translation and compilation into a dataflow for the Dispatch pool. The dashed arrow supplies catalog metadata to binding. Query planner SQL 1. Parse AST 2. Bind + plan Logical plan 3. Optimize Optimized plan 4. Translate and compile Dataflow Catalog snapshot schemas + tables Dispatch pool

The planner resolves a statement against the query’s catalog snapshot and pushes projections and predicates into the scan. Partition values, file statistics, row-group statistics, and dictionaries can eliminate work before any column data is decoded.

The catalog is Pivot’s bridge between the planner and the datastores. It keeps track of the datastores available on a server and routes table lookups to the right one. Through this interface, the planner discovers schemas and tables, resolves column names and types, and obtains metadata such as row-count estimates without needing to understand how each datastore persists its data.

Each query opens a catalog transaction. The first time it accesses a datastore, the transaction captures that datastore’s snapshot and reuses it for the rest of the query. Planning, scans, and late materialization therefore agree on the same table versions and file sets, even if a background refresh discovers newer data while the query is running. Snapshots are taken independently for each datastore, not as one atomic snapshot across all datastores.

The catalog also exposes a read-only “virtual” system tables datastore, so users can query metadata about the current running pivot instance.

A datastore is a collection of schemas and tables exposed to Pivot through a common interface. Each configured datastore implementation handles table discovery, metadata, snapshots, and supported read and write operations for its underlying storage or table format.

Pivot currently supports/exposes the following datastores:

  • Pivot - A collection of Delta Lake tables.
  • Iceberg - Tables stored using the Apache Iceberg table format.
  • System - Internal tables exposing information about the running Pivot instance.

A server can expose multiple datastores, with tables addressed as datastore.schema.table.

For example, a query can join orders in a Pivot datastore named analytics with customer details in an Iceberg datastore named lake:

SELECT
orders.order_id,
customers.customer_name,
orders.total_amount
FROM analytics.sales.orders AS orders
JOIN lake.crm.customers AS customers
ON orders.customer_id = customers.customer_id;

Here, analytics and lake are datastore names, while sales and crm are schemas within those datastores. Pivot reads each table through its datastore and executes the join in the same query engine.

The metastore is a collective of configurations and secrets that declare the metadata a pivot instance needs for it to run.

These include:

  • Datastores - their names, implementations, locations, and settings, including which datastore is the default for unqualified table names.
  • Secrets - credentials for accessing object storage, scoped to the locations they apply to. Multiple datastores can use the same secret.
  • Users - which users are authored to access pivot, and how do they authenticate.

Thanks to the separation of datastores and metastore, different instances / deployments of pivot might access different /overlapping datastores with different permission models and with different configurations:

Two Pivot instances with different metastores share one datastore Object storage holds two datastores: analytics, a Pivot datastore, and lake, an Iceberg datastore. A Pivot shell started with pivot open on the analytics location uses an ephemeral metastore and reads and writes only analytics. A Pivot server started with pivot server and a metastore from pivot.yaml reads and writes analytics, reads lake, and serves SQL clients such as the analyst user over the Postgres wire. Object storage S3 · GCS datastores analytics Pivot · Delta Lake format s3://company-data/pivot/ pivot manifest + delta log + parquet lake Iceberg s3://lakehouse/warehouse/ metadata + manifests + parquet read + write read + write read Pivot shell metastore $ pivot open s3://company-data/pivot/ Pivot server metastore $ pivot server --config pivot.yaml Postgres wire SQL clients backend · psql · BI tools $ psql -h pivot.internal -U analyst

Upcoming: PostgreSQL-backed metastores are under development and are not available in the current release. The configuration may change before the feature is merged.

A PostgreSQL-backed metastore allows running multiple Pivot instances that share a single source of truth for datastore definitions, users, and other configuration.

Keeping this metadata in a simple to setup PostgreSQL endpoint makes a Pivot cluster easy to deploy, manage, and coordinate compared to requiring a dedicated coordination service such as ZooKeeper.

A Pivot cluster behind a load balancer, sharing one metastore and one datastore The datastore in object storage sits above the cluster and is what every instance reads and writes. A PostgreSQL metastore sits beside the cluster and supplies configuration and identity to every instance. The cluster holds N Pivot server instances. A load balancer routes each connection to one instance, and SQL clients connect to the load balancer over the Postgres wire. Datastore s3://company-data/pivot/ data plane read / write Metastore PostgreSQL · control plane datastores locations users scram-sha-256 credentials storage keys config + identity Pivot cluster Instance 1 server … Instance N server any instance Load balancer one address for the whole cluster pivot.internal:5432 Postgres wire SQL clients backend · psql · BI tools

Any user or datastore created on one instance is immediately visible to all other instances. Scaling rules and other cluster-wide configuration can be coordinated from a single central location.

Pivot can run as a standalone SQL shell on top of a datastore or as a server accepting client connections. Both modes use the same Catalog, Planner, and Dispatch pool, running together in a single process. Cluster deployments with a shared PostgreSQL metastore are planned to be released soon.

pivot open starts an interactive SQL shell with its own query engine. It opens one datastore, local or in object storage, and uses an ephemeral metastore for the session. Queries execute on the machine running the shell, even when the data is stored remotely, and tables remain in the datastore after the shell exits.

Terminal window
pivot open s3://my-bucket/pivot-data

A local datastore can be opened by only one process at a time, so use a remote datastore when several instances need the same data. See the CLI reference for options, credentials, and shell commands.

pivot server runs a persistent SQL endpoint that accepts connections from applications, BI tools, and PostgreSQL clients. Each connection submits its queries to the server’s query engine; concurrent connections share that instance’s compute resources. A server can expose multiple datastores and execute queries across them.

Terminal window
pivot server --config pivot.yaml

The server runs in the foreground; connect from another terminal with any PostgreSQL client. The configuration file sets the instance’s resources, its endpoint, the datastores it serves, and its users.

Separate servers can share the same remote datastores, each with its own configuration, and see each other’s commits through background refresh. Run compaction and vacuum in only one of them per datastore, as described in datastore maintenance.