Skip to content

Query engine

POST /api/v1/sql runs SQL through DataFusion. The query server reads the current Iceberg snapshot and can add uncommitted ingest rows from the optional WAL buffer.

The engine first looks for an exact metadata answer. If that is unavailable, it prunes files and rows before scanning. Compatible ordered queries can stop early, and eligible large plans can run across query replicas.

Each query replica caches Iceberg table metadata separately. The default SIGLAKE_ICEBERG_METADATA_CACHE_TTL_SECS time to live (TTL) is 5 seconds. During a background refresh, a replica can serve stale metadata up to the default 60-second cache-age ceiling. Two requests sent to different replicas can therefore see different committed generations. A single distributed query never mixes generations: each worker uses the generation pinned by its coordinator or refuses the shard. See Query replicas can briefly disagree across a commit for the cross-request limit and remedies.

Aggregate fast paths

Siglake stores counts and time buckets in Parquet footers and table side objects. The query server can use this metadata for supported shapes, including:

  • whole-table counts and group counts
  • date histograms and windowed counts
  • negation and integer-range counts
  • count(DISTINCT <column>)

A valid table-level side object reports a served_by value beginning with tier1_. These paths do not open data files and report rows_scanned: 0.

If table-level metadata is absent, the same query can use per-file metadata. This exact fallback reports served_by: "materialized". It can read footers or data pages even though rows_scanned remains zero, so that field alone does not show its I/O cost.

Readers compare an aggregate's total with the snapshot record count before using it. Missing, stale, or incomplete metadata causes an exact fallback instead of a partial answer. Floating-point values are excluded from typed group-count entries because their string form is not a stable lookup value.

The WAL buffer computes a delta for eligible shapes and folds it into the committed result. If the delta cannot be loaded within its limits, the request uses another execution path.

See Aggregate repair when a query stays on the materialized path.

Pruning

Scans can remove work at several levels:

Level Evidence
Partition Iceberg day(timestamp) partition values.
File Manifest bounds and file-level trigram blooms.
Row group Parquet statistics and row-group token blooms.
Row Page indexes, positional deletes, and optional inverted indexes.

An absent accelerator does not change results. The engine reads more data instead. This fallback rule matters most for bloom and inverted-index formats, where a false negative would omit rows.

stats.scan reports planned and completed work, including files_planned, files_read, row_groups_considered, row_groups_read, and rows_pruned_selection.

Ordered early-stop

A scan can advertise timestamp order when its files have usable bounds and compatible sort metadata. DataFusion can then satisfy ORDER BY timestamp ... LIMIT <n> without a full blocking sort.

Time-disjoint files are concatenated in order. Overlapping files use a k-way merge within per-partition and global fan-in limits. The defaults derive from the query pod's CPU count and stay within configured floors and ceilings. SIGLAKE_ORDERED_MERGE_MAX_FANIN and SIGLAKE_ORDERED_MERGE_GLOBAL_FANIN override the derived limits.

The scan follows the table's declared sort direction. For the other direction, it can decode files, row groups, and batches in reverse. This supports tables that still use the older descending timestamp order.

stats.scan.ordering explains whether ordering was advertised. The siglake_query_scan_output_ordering_total{outcome} counter uses these main outcomes:

Outcome Meaning
advertised The scan provides timestamp order.
filtered A pushed predicate makes the parallel pruned scan preferable.
fan_in One overlap cluster exceeds the partition fan-in limit.
global_fan_in All merge streams together exceed the global limit.
no_bounds A file lacks timestamp bounds.
mixed_sort_direction Files carry conflicting sort directions.
unknown_sort_order A file has unusable sort metadata.
order_walk_failed The file sort-order walk failed.
missing_sort_column The output schema lacks timestamp at the final check.

The response can also report not_projected, no_sort_order, non_identity_sort, missing_sort_field, or non_timestamp_sort. These describe the query projection or table declaration, so the counter does not use them as labels.

Ordered browses are slow maps each outcome to checks and remedies.

Distributed query

Query replicas have stable peer addresses. Any replica can coordinate an eligible request by dividing table files among workers and merging their results.

Small ordered browses stay on the coordinator by default. Replicas increase concurrent capacity for these queries rather than reducing one query's latency. Large scans can fan out. The selective-browse fan-out gate controlled by SIGLAKE_DIST_BROWSE_MIN_SCAN_ROWS is off by default.

Only events and managed user indexes can use file-shard fan-out. Joins and hot-cache table functions run locally.

Admission is per coordinator. The pod that receives the request reserves one share against its own budget, held through the merge; a shard request takes no reservation, because the coordinator already admitted the query. A worker is bounded instead by the process memory pool and the mid-flight rows-scanned breaker, and it spills sorts and aggregates to its own node-local scratch. What that costs at scale is in Distributed admission is per-coordinator.

To configure any of this - peer discovery, the coordinator's shard credential, spill scratch space, or the admission budget - see Distributed query.

One generation per fan-out

The coordinator pins each shard request to the table generation it planned against: the committed snapshot id, and the Iceberg schema id. A worker resolves both halves before it reads its file slice.

The schema half matters after an additive migrate-schema. That migration commits a schema change and no data snapshot, so the schema id moves while the snapshot id stands still. A snapshot-only pin was therefore satisfied by a worker whose metadata cache predated the migration, which planned its shard against the narrow column set and failed a fragment that named the new column. With the schema id pinned, a distributed query reads a migrated column on every shard as soon as the migration commits, the same as a single-pod query. A worker refreshes its metadata once if it lags, and it serves a historical schema out of the metadata it retains, so a worker running ahead of the coordinator answers the generation it was asked for rather than refusing.

If either half cannot be resolved, because the snapshot has expired, the schema is unknown, or neither is visible after a metadata refresh, the worker returns 503 with Retry-After and reason shard_pin_unresolved. The coordinator does not retry that shard against a different generation or run it locally. The request therefore returns one generation's result or an error, never a mixture of generations.

The schema id is optional on the wire, so a rolling upgrade stays compatible in both directions. A coordinator that omits it gets the snapshot-only enforcement it expects, and an older worker ignores it. An older worker cannot enforce the schema half, so a mixed-version fleet keeps the pre-upgrade behavior until every replica runs the new image.

This consistency rule can reduce availability while caches and snapshot expiry race. siglake_query_shard_pin_total{outcome="miss"} reports failures. See A 503 with reason: shard_pin_unresolved for operational checks.

Metadata fast paths and response format

For records responses, the coordinator tries eligible count, group-count, distinct-count, negation-count, and histogram metadata paths before it classifies a distributed plan. These paths also fold in the WAL delta.

An ndjson request skips that coordinator metadata battery. Mergeable aggregates then run on shards and are combined by the coordinator. Use ndjson, a shape outside the fast paths, or /api/v1/sql/distributed when you need to test fan-out itself.

The classifier chooses one of these plans:

Plan Coordinator behavior
Scan Concatenates shard rows, then reapplies the limit.
OrderedScan Merge-sorts limited shard results, then reapplies the limit.
Aggregate Combines shard aggregate values.
OrderedAggregate Combines complete shard groups before applying global ordering and limit.
Local Runs the whole query on one replica.

OrderedAggregate cannot send a top-N grouping unchanged to each shard. A group outside every shard's local top N can still be in the global top N after the counts are combined.

Workers receive a ScanShard, run POST /api/v1/sql/shard, and return Arrow IPC. /api/v1/sql/local always stays on one replica, while /api/v1/sql/distributed forces coordination.

Mergeability rules

count, sum, min, and max can merge across shards, with or without GROUP BY. The coordinator sums partial counts and sums, and applies min or max to the partial values.

avg, DISTINCT aggregates, window functions, joins, and queries with OFFSET run locally. An aggregate over a subquery with an inner LIMIT also runs locally. These shapes do not have a merge rule that preserves SQL semantics.

stats.phases.distributed reports the selected mode, each shard's wall time, and coordinator merge time. Because shards run concurrently, the slowest shard matters more than the sum of their times.

Caching

Query caches use the table snapshot as part of their identity where the result depends on committed data. The cache layers include Parquet footers, aggregate results, windowed results, live file lists, decoded scan batches, and complete small SQL responses.

The standing rule is that a result cache entry must be a function of its table snapshot and query inputs. A time-to-live value alone is not a safe invalidator for a leading-edge query.

The benchmarking caveat

SIGLAKE_QUERY_RESULT_CACHE=off disables complete-result memoization. Use it when measuring query execution rather than cache hits. It is a measurement setting, not a general performance recommendation.

Guardrails

The query server applies four different limits because no single estimate is exact for every plan:

  1. A manifest-based estimate can reject a request before data-file I/O.
  2. Per-tier ceilings clamp limits supplied in the request.
  3. A mid-flight breaker stops a scan after it crosses the row budget.
  4. A wall-clock deadline bounds the complete request.

dry_run: true and /api/v1/sql/explain return the estimate without running the query. A request can lower its limits, but it cannot raise the configured ceilings.

The estimate assigns one of these classes:

Class Estimated upper range
small 100 MB and 1 million rows
medium 10 GB and 100 million rows
large 100 GB and 1 billion rows
huge Above the large range

cost.exact says whether the source statistics were exact. See the guardrail settings for the active ceilings.

Execution tiers

Priority Behavior
interactive Runs on the shared interactive runtime and returns the result in the request.
batch Creates an asynchronous job on a separate runtime. The client polls the job endpoint.

The runtimes are separate, but both use the query admission budget. Admission can reject a batch submission before it creates a job. A running batch job keeps its reservation until it finishes or cancellation reaches its owner.

DataFusion's process memory pool covers scans, sorts, aggregates, joins, and the WAL-buffer part of a query. Under pressure, reads reduce concurrency and operators can configure spilling. A request that cannot reserve memory returns 503 with Retry-After.

Audit

The query server submits completed queries to a bounded audit writer. Retained rows include the subject, endpoint, SQL, format, priority, timing, status, cost, and truncation state. These rows live in the query_audit Iceberg table. See Audit logging for the drop conditions.

Audit retention is an operator task. See Retention and deletes before using siglake audit-rotate.

Response accounting

Successful SQL responses include cost, stats.scan, stats.phases, and the x-siglake-server-micros header.

  • Compare cost.files_considered with cost.files_to_scan for planning-time pruning.
  • Compare stats.scan.files_planned with stats.scan.files_read for execution-time pruning.
  • Compare planning and collection times in stats.phases.
  • Check stats.phases.distributed.shard_wall_micros for a slow shard.

The HTTP API reference defines the complete response fields.