Architecture
System architecture
Section titled “System 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.
Dispatch execution pool
Section titled “Dispatch execution pool”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.
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.
Planner
Section titled “Planner”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:
- 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.
- 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.
- 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.
- 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.
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.
Catalog
Section titled “Catalog”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.
Datastore
Section titled “Datastore”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_amountFROM analytics.sales.orders AS ordersJOIN 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.
Metastore
Section titled “Metastore”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:
Sharing a metastore across deployments
Section titled “Sharing a metastore across deployments”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.
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.
Deployment Architecture
Section titled “Deployment Architecture”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.
Standalone
Section titled “Standalone”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.
pivot open s3://my-bucket/pivot-dataA 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.
Server
Section titled “Server”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.
pivot server --config pivot.yamlThe 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.
