Scaling¶
Each role scales on a different signal, so size them separately. KEDA covers the ingester and query; the legacy CPU horizontal pod autoscalers (HPAs) cover the ingester and compactor. This page gives the shape of each tier and the sizing rules the benchmark arc produced.
Ingest¶
Ingest scales horizontally, close to linearly. Measured accept throughput on AWS m6i-class nodes:
| Ingesters | Sustained rows/s |
|---|---|
| 1 | ~108 K |
| 2 | ~200 K |
| 3 | ~324 K |
| 4 | ~430 K |
Budget roughly 50,000 events per second (EPS) per pod-CPU at default
write-ahead log (WAL) segment thresholds. A single ingester at limits.cpu: 1
sustains about 55,000 EPS on log-line events.
Scale out rather than up. Past limits.cpu: 2 returns go sub-linear, because
writer tasks serialise per tenant. If you must push a single tenant harder on
one pod, raise --ingest-backpressure-shards /
ingester.backpressure.shards to add writer tasks per tenant.
The ingester and compactor share one WAL tree. Use ReadWriteMany (RWX) whenever
WAL consumers may run on different nodes, which includes scaling ingesters and
enabling the query freshness buffer without co-locating their pods.
ReadWriteOnce (RWO) requires every WAL consumer on one node, and the chart does
not co-locate them for you. For RWO, set matching nodeSelector mappings under
ingester and compactor to select one uniquely labelled node, add the same
mapping under query when query.walBuffer.enabled is true, and leave
antiAffinity.enabled false.
Autoscale ingest on request rate and backpressure queue depth.
Size drain and compactor pods¶
Compaction scales horizontally only with catalog claim, and the chart enforces
that rather than advising it. compactor.replicas above 1 fails the install
while compactor.catalogClaim.enabled is false, and so does
autoscaling.compactor.maxReplicas above 1 when
autoscaling.compactor.enabled is true.
compactor:
replicas: 3
catalogClaim:
enabled: true
extraArgs: ["--role", "drain"]
wal:
mirror:
enabled: true # required: the claim drain reads segments from the bucket
Replicas without the claim share no claim table. Each one runs the whole drain and maintenance loop against the same tables, they lose Iceberg optimistic-concurrency commit races to each other, and the layout stops converging. Nothing reports it as a misconfiguration, which is why the chart refuses the combination instead.
The claim needs the mirror, and that direction is refused too:
compactor.catalogClaim.enabled: true with wal.mirror.enabled: false fails
the install. In claim mode the drain reads the wal_segments catalog table and
never the local sealed/ directory, and rows land there only from an ingester
that mirrors its segments. Each refusal names the value to set: turn the claim
on, turn the mirror on, or hold the compactor at one pod.
Set dead-worker reclaim cadence, stale-claim age and committed-object retention with the catalog-claim settings.
That recipe gives you three drain-only replicas. Run one separate compactor
process against the same Iceberg warehouse and catalog with siglake compactor
--role maintenance; do not run one maintenance process per drain. Maintenance
reclusters files, expires snapshots, samples gauges, folds aggregates and
executes delete tasks, with the lease electing work per table. The chart renders
only one compactor Deployment, so a split topology needs a maintenance
Deployment you manage separately. For a single compactor process, keep the
default --role combined.
The sizing rule the benchmark arc produced: once drains keep up with accept,
wall-clock throughput equals sustained throughput. A round with 3 ingesters and
8 drains finished in 27.5 minutes, where a round with more ingesters and fewer
drains took 36.5. If wall time matters to you, add drains until
siglake_compactor_sealed_pending stops growing.
The packaged compactor limit is 1 GiB and 2 CPU. At that size, safe compaction
concurrency is one bin at a time, so the default lags hardest while a backfill
drives overlap depth up. With no explicit override, Siglake derives bin
concurrency from the cgroup memory limit and the available cores. The chart
pins that choice with compactor.binConcurrency; raise it only alongside
compactor.resources.limits.memory and CPU. Budget a 1 GiB process reserve
plus roughly 4 GiB per concurrent bin, with at least one CPU per bin.
Query pods derive their caches and execution pool from their memory limit too; use the Helm query memory budget when sizing them. Each query pod also requests 10Gi of node ephemeral storage for spill scratch, so read Query spill and ephemeral storage before sizing node disks.
The Helm chart cannot scale a multi-pod compactor from its backlog. It refuses
autoscaling.compactor.enabled: true and
autoscaling.compactor.customMetric.enabled: true together with
compactor.catalogClaim.enabled: true. The claim-mode gauge reports the whole
shared queue of sealed, unclaimed segments, not one pod's share. The Pods
metric algorithm would make the requested count depend on the current replica
count. Filesystem mode stays capped at one pod, so the custom metric cannot
scale that mode either.
Use siglake-operator for backlog-driven compactor scaling. With
spec.autoscaling.compactor.max above 1, set
spec.autoscaling.compactor.target to the number of sealed segments per
worker. The operator divides the shared queue depth by that target once to
choose the replica count.
The catalog is on the commit path. At high commit rates RDS becomes the
bottleneck before the compactor does, so watch
siglake_iceberg_commit_duration_seconds.
Size query pods and understand distributed fan-out¶
Treat query.resources.limits.memory as each query pod's sizing budget. The
default 4 GiB limit gives the pod a 1 GiB object cache, 512 MiB of metadata
caches, a 1.25 GiB DataFusion pool and about 1.25 GiB for the process and
transient allocations. It gives the two text-index caches nothing: they take
only what is left once the pool can still reserve one compacted file's decode
working set, so a pod at the floor deserializes a Puffin index on every text
query. A 5 GiB pod derives 400 MiB of them. Size for text search accordingly,
or override the two byte limits and accept a smaller pool.
Sorts and aggregates spill when they reach the pool bound. If a query needs to
exceed the 8 GiB spill cap, Siglake refuses it with ResourcesExhausted as a
retryable 503. If you need more spill room, raise query.spill.maxBytes,
query.spill.sizeLimit and query.resources.limits.ephemeral-storage
together. Keep the defaults ordered at 8 GiB, 10 GiB and 12 GiB. The Helm
query memory budget derives the memory split and
the spill settings explain the
three storage ceilings.
Set replicas from concurrent-query demand. One 2026 browse measurement scaled through 16 concurrent queries and then plateaued when decode saturated the CPU. Use that result as a starting point to test your workload, not as a fixed limit. Aggregate-heavy workloads can scale further because they do less decode work.
Ordinary log search scales horizontally by replication, not fan-out. Query replicas form a StatefulSet and any replica can answer a request. Add replicas to multiply concurrent-query throughput; do not expect them to reduce the latency of one ordinary search.
If query.replicas > 1 and bearer-token or OIDC authentication is on,
configure query.distributed.coordinatorToken for
fan-out. Helm refuses to render without it.
A second pod also needs the shared batch job
store, which is the default. With
query.jobs.persistent: false each pod keeps its own in-memory job table, and
the Service routes a status, result or cancel request to a pod that has never
heard of the job. Helm fails the install when query.replicas, or an enabled
keda.query.maxReplicas, is above 1 with the store off, so the in-memory
opt-out is available only to a tier you hold at one pod.
Small-LIMIT browses and Tier-1 metadata aggregates classify as Local, so
the replica that receives the request answers it. The default
SIGLAKE_DIST_SCAN_LOCAL_MAX_LIMIT=100000 keeps bounded scans local, because
coordination costs more than it saves below that limit. The selective-browse
fan-out gate, SIGLAKE_DIST_BROWSE_MIN_SCAN_ROWS, ships off.
The 2026-08-18 measurements showed 705 QPS on one replica and 2,269 QPS on three, a 3.22x gain. Browse p50 stayed between 13 ms and 21 ms. In a separate topology with one coordinator wired to four peers, the coordinator answered all 7,409 queries itself and each peer received only two health checks. Size the query tier for concurrency, not for single-query latency.
File-shard fan-out is still available and engages for large scans. Keep
minReplicas >= 2 when those scans must always have peers.
The coordinator discovers Ready replicas through the query Service's SRV record and pins one membership set for the query. It assigns file shards to that set. A scale event changes the next query, after readiness and the default five-second discovery refresh. It does not move a query already in flight.
Every fan-out uses one snapshot chosen by its coordinator. Workers refuse a snapshot they cannot resolve instead of mixing generations. Separate requests can still reach replicas whose five-second metadata caches straddle a commit, so two query replicas can briefly return results from different snapshots.
Layout quality dominates query scaling. On an unconverged 1 TB layout, 32-way
mixed load collapsed to 28 QPS; on the same data converged, 115 QPS. Before you
add query replicas, check siglake_table_overlap_depth. You may have a
compaction problem rather than a capacity problem.
Autoscaling with KEDA¶
query:
replicas: 2 # starting count only; KEDA owns spec.replicas after that
keda:
enabled: true
prometheusServerAddress: http://prometheus-server.monitoring.svc.cluster.local:80
pollingIntervalSeconds: 15
cooldownPeriodSeconds: 300
scaleUpStabilizationSeconds: 0 # absorb spikes immediately
scaleDownStabilizationSeconds: 300 # anti-flap
ingester:
minReplicas: 1
maxReplicas: 10
requestsPerSecondTarget: "800"
backpressureQueueTarget: "256"
query:
minReplicas: 2
maxReplicas: 12
inFlightTarget: "8"
p95QueueWaitSecondsTarget: "0.5"
The contention trigger uses queue wait rather than request latency, because KEDA divides the signal by replica count: the signal has to fall as pods are added, and request latency does not.
keda.query.maxReplicas may exceed query.replicas. The chart renders
--query-peer-discovery-srv rather than a peer list, so every Ready replica
joins the shard membership and query.replicas is only the starting count, and
the count used when KEDA is off.
Two convergence properties come with that. A pod KEDA adds becomes eligible one
readiness probe plus one discovery refresh
(--query-peer-discovery-interval-secs, default 5 s) after it starts. A query
already running keeps the membership it pinned, so scale-out shows up on the
next query, not the one in flight. Watch
siglake_query_peer_discovery_members against the Ready replica count to see
membership follow a scale event; see Query peer
discovery.
The operator takes the same ordinary range: spec.autoscaling.query.max may
exceed min, and it scales query on in-flight queries per pod. Raising that
maximum above 1 needs a shared job store the same way the chart does. The
operator refuses the spec with
QueryJobsStoreRequiredForScaleOut when
the effective jobs-store URI is blank, which a Postgres spec.catalogUri
avoids by default.
Every operator-managed tier needs a floor of 1 or more. A min of 0 is
refused with
AutoscalingZeroFloorUnsupported,
because each tier's scaling signal comes from its own pods and nothing would
ask a stopped tier back.
For the operator's compactor, spec.autoscaling.compactor.max also fixes the
drain mode. A maximum above 1 keeps catalog claims, remote drain and
RollingUpdate while the autoscaler holds the compactor at one replica. A
maximum of 1 keeps the filesystem drain and Recreate. Scaling within the
policy does not hand retained WAL between drain protocols. Crossing the
boundary needs a manual handover.
Scale-up is immediate and scale-down is stabilised over five minutes, because log ingest is bursty and thrashing replicas costs more than briefly over-provisioning.
Only the ingester and query have KEDA ScaledObjects. The compactor does not
autoscale, because its right signal is backlog depth and its horizontal scaling
requires catalog claim.
Why scale-down is safe¶
The ingester force-seals its active WAL on SIGTERM and drains the backpressure
lane before exiting. With preStop and a long enough
terminationGracePeriodSeconds, no acknowledged event is lost when a replica
goes away.
Query pods are stateless apart from node-local spill scratch, an emptyDir
that goes away with the pod. They drain in-flight queries and exit.
Set terminationGracePeriodSeconds above preStopSleepSeconds plus your
longest expected operation.
Legacy HPAs¶
The autoscaling.* block renders CPU-based HPAs for the ingester and compactor
only, off by default. KEDA is the chart's only query autoscaler, on the
in-flight and queue-wait signals rather than CPU. CPU
works anywhere metrics-server is installed, while custom metrics need
prometheus-adapter. For the ingester, prefer KEDA where it is available.
Storage scaling¶
Size the WAL volume for your worst-case backlog, not for steady state. If the
drain stalls, sealed segments accumulate until it recovers. The chart default
is wal.size: 50Gi, which at 400 K rows/s is not much runway.
Move EFS off bursting throughput for heavy workloads. The Terraform module provisions bursting; under sustained high ingest, burst credits deplete and WAL writes slow down, which is easy to misdiagnose as an ingester problem.
Size the catalog against your commit rate. db.t4g.micro is a benchmark
default.
Size ingester, compactor, and query pods¶
Ingest rate sets the ingester and drain counts. Dataset size reaches you indirectly, as overlap depth and as query concurrency.
- Divide your target ingest rate by 50,000 EPS per CPU to get ingester CPU,
then replicas. Stop at
limits.cpu: 2per pod and add replicas instead. - Add drains until
siglake_compactor_sealed_pendingis flat under sustained load. More than one compactor replica requires catalog claim. - Size the WAL volume for backlog, not throughput.
- Confirm
siglake_table_overlap_depthconverges before you size query. - Set query replicas from concurrency. Use the measured 16-way browse plateau as a starting point, then test your workload.
For scale: the 1 TB reference fleet was four ingesters and three drain nodes at about 415,000 rows/s. See Performance.