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:
- A manifest-based estimate can reject a request before data-file I/O.
- Per-tier ceilings clamp limits supplied in the request.
- A mid-flight breaker stops a scan after it crosses the row budget.
- 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_consideredwithcost.files_to_scanfor planning-time pruning. - Compare
stats.scan.files_plannedwithstats.scan.files_readfor execution-time pruning. - Compare planning and collection times in
stats.phases. - Check
stats.phases.distributed.shard_wall_microsfor a slow shard.
The HTTP API reference defines the complete response fields.