Skip to content

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.

query:
  replicas: 4
  distributed:
    enabled: true

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.

The chart renders KEDA ScaledObjects for the ingester and query only. It does not autoscale the compactor, because the right compactor signal is backlog depth and its horizontal scaling requires catalog claim. The operator does scale the compactor, on that backlog depth.

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.

  1. Divide your target ingest rate by 50,000 EPS per CPU to get ingester CPU, then replicas. Stop at limits.cpu: 2 per pod and add replicas instead.
  2. Add drains until siglake_compactor_sealed_pending is flat under sustained load. More than one compactor replica requires catalog claim.
  3. Size the WAL volume for backlog, not throughput.
  4. Confirm siglake_table_overlap_depth converges before you size query.
  5. 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.