Routing, cluster groups, and clusters
QueryFlux separates where a query should go logically (cluster group) from which physical backend instance serves it (cluster / adapter). This document explains that split, how routers work, and how the cluster manager balances load and enforces limits.
Vocabulary
| Term | Meaning |
|---|---|
| Cluster | A named backend instance in config (clusters.<name>): one engine (Trino, DuckDB, StarRocks, …), endpoint or DB path, auth, etc. At runtime it has an engine adapter and a ClusterState (running count, health, limits). |
| Cluster group | A named pool (clusterGroups.<name>) listing member cluster names, a per-group maxRunningQueries, optional maxQueuedQueries, optional selection strategy, and enabled. Routing returns a group; the cluster manager then picks a member cluster. |
| Router | A rule that inspects the SQL string, session context, and frontend protocol and optionally returns a target group name. |
| Router chain | All routers in config order, plus routingFallback when every router returns “no match.” |
Two-stage placement
-
Routing (group selection)
RouterChainevaluates eachRouterTraitimplementation in order. The first router that returnsSome(ClusterGroupName)wins. If all returnNone,routingFallbackis used.
Implementation:queryflux_routing::chain::RouterChain(route,route_with_trace). -
Cluster selection (member selection)
ClusterGroupManager::acquire_cluster(group)considers only clusters in that group that are enabled, healthy, and undermax_running_queries. It then uses the group’s strategy to pick one member and increments that cluster’s running count.
If no member is eligible (e.g. all at capacity or unhealthy),acquire_clusterreturnsNone→ the query is queued (Trino HTTP async path) or retried with backoff (syncexecute_to_sinkpath).
Implementation:queryflux_cluster_manager::simple::SimpleClusterGroupManager.
When the query finishes (success, failure, or cancel), release_cluster decrements the running count on that cluster.
Router types (config → code)
Configured under routers: in YAML (queryflux_core::config::RouterConfig). Wired in queryflux/src/main.rs.
type | Behavior |
|---|---|
protocolBased | Maps the active frontend (trinoHttp, postgresWire, mysqlWire, flightSql, clickhouseHttp) to a group name. |
header | Matches a header value to a group (useful for Trino HTTP and similar). |
queryRegex | Ordered rules: first regex match on the SQL text wins. |
tags | Routes based on query tags attached to the session. Each rule specifies one or more tag key/value conditions (AND logic); first matching rule wins. See Tags router below. |
pythonScript | Embedded or file-backed Python route(query, ctx) returning a group name or None. See Python script router below. |
compound | Multiple conditions combined with all (AND) or any (OR). Supported condition types: protocol, header (name + value), user, clientTag, queryRegex. |
All six router types are implemented. Unknown type values in config are skipped at startup with a warning.
Tags router (tags)
Routes queries based on key/value tags that clients attach to their sessions. Rules are evaluated in order; the first rule where all specified tags match wins (AND logic within a rule).
routers:
- type: tags
rules:
- tags:
team: eng
env: prod
targetGroup: prod-eng-cluster
- tags:
team: analytics
targetGroup: analytics-cluster
- tags:
batch: null # key-only match — any value (or no value) accepted
targetGroup: batch-cluster
Matching semantics per tag entry:
| Config value | Matches when |
|---|---|
"some-value" | Tag key is present and its value equals "some-value" |
null | Tag key is present, regardless of its value (key-only match) |
How clients attach tags depends on the frontend protocol — see Query tags for full details on how to set tags from each frontend (Trino, MySQL wire, Postgres wire, ClickHouse HTTP).
Cached routing config and DB reload (Postgres)
When persistence.type is postgres, routing rules and cluster/group definitions loaded from the database are held in memory inside LiveConfig (including the compiled RouterChain). Each request reads the current chain from that shared snapshot (Arc<tokio::sync::RwLock<LiveConfig>> in queryflux-frontend).
- Periodic refresh:
queryflux.configReloadIntervalSecsin YAML (default 30 when omitted) controls how often a background task re-reads Postgres and replacesLiveConfigin one atomic swap. Implementation:crates/queryflux/src/main.rs(reload task) andreload_live_config→load_routing_config. 0disables polling only: WithconfigReloadIntervalSecs: 0, there is no timer-driven refresh; the in-memory config stays as loaded at startup until an immediate refresh runs (below).- Immediate refresh: After Studio/admin API writes to routing, clusters, or groups, the proxy notifies the same task so a reload runs without waiting for the interval (
config_reload_notifyinadmin.rs). - YAML-only mode: With
inMemorypersistence there is no DB reload loop; routing comes from the process config at startup until restart.
Python script router (pythonScript)
The script must define:
def route(query: str, ctx: dict) -> str | None:
...
query: SQL text (the same string the router chain sees).ctx: plaindictbuilt by QueryFlux (string keys).protocolis always set; all other keys are protocol-agnostic.
| Key | Type | Meaning |
|---|---|---|
protocol | str | One of trinoHttp, postgresWire, mysqlWire, clickhouseHttp, flightSql, snowflakeHttp, snowflakeSqlApi (camelCase). |
user | str | None | Resolved user identity. Extracted from X-Trino-User, Postgres startup message, MySQL handshake, etc. |
database | str | None | Target database/catalog hint. Extracted from X-Trino-Catalog, Postgres database, MySQL USE, etc. |
extra | dict[str, str] | Protocol-specific key-value bag. For Trino/ClickHouse HTTP: lowercased header names (e.g. x-trino-source). For Postgres: startup parameters. For MySQL: session variables. Empty for other protocols. |
auth | dict | None | {"user": str, "groups": [str, …], "roles": [str, …]} when the request was authenticated. Raw JWT / bearer tokens are not passed into Python. |
Query tags are not injected into the Python context dict. Use the tags router type for tag-based routing instead.
Example (route by user):
def route(query: str, ctx: dict) -> str | None:
if ctx.get("user") == "batch":
return "heavy-trino"
return None
Example (inspect a Trino HTTP header via extra):
def route(query: str, ctx: dict) -> str | None:
if ctx.get("protocol") != "trinoHttp":
return None
source = (ctx.get("extra") or {}).get("x-trino-source")
if source == "dbt":
return "dbt-group"
return None
Routing trace
route_with_trace records each router’s decision (matched, optional result) and whether the fallback group was used. This supports debugging and future UI/metrics (see RoutingTrace in queryflux_routing::chain).
Cluster group configuration (actual shape)
Clusters are defined once at the top level; groups reference them by name:
clusters:
trino-1:
engine: trino
endpoint: http://trino:8080
enabled: true
duckdb-1:
engine: duckDb
enabled: true
clusterGroups:
trino-default:
enabled: true
maxRunningQueries: 100
members: [trino-1]
strategy:
type: leastLoaded
duckdb-local:
enabled: true
maxRunningQueries: 4
members: [duckdb-1]
Notes:
maxRunningQuerieson the group applies to each member cluster’sClusterStatewhen those states are built (seemain.rspass 2). It is the cap used for acquire / capacity checks.memberscan list multiple clusters, including mixed engines (e.g. Trino and DuckDB in one group). For that,engineAffinity(or another strategy) helps express preference order across engine types (queryflux_cluster_manager::strategy).
Selection strategies
Configured as strategy: { type: ... } on a group. Implemented in strategy.rs:
| Strategy | Behavior |
|---|---|
roundRobin | Rotates among eligible members (default when strategy omitted). |
leastLoaded | Picks the member with the smallest running_queries. |
failover | First eligible member in member list order (priority ordering in YAML). |
engineAffinity | Ordered engine preference; within each engine, least loaded. |
weighted | Distributes by configured weights (deterministic pseudo-random from load). |
pythonScript | Embedded or file-backed Python select_cluster(candidates) returning a member cluster name or None. See Python script strategy below. |
Eligible candidates are always healthy, enabled, and not at capacity before the strategy runs.
Python script strategy (pythonScript)
For groups where none of the built-in strategies fit — custom cost/tiering rules, time-of-day shaping, or any placement logic specific to your deployment — a group's strategy can run operator-supplied Python instead. This is the same pattern as the Python script router above, one stage down: that picks a cluster group, this picks a member cluster within a group already chosen by routing.
clusterGroups:
trino-default:
enabled: true
maxRunningQueries: 100
members: [trino-a, trino-b, duckdb-1]
strategy:
type: pythonScript
script: |
def select_cluster(candidates: list[dict]) -> str | None:
trino = [c for c in candidates if c["engineType"] == "trino"]
if trino:
return min(trino, key=lambda c: c["runningQueries"])["name"]
return None
The script must define:
def select_cluster(candidates: list[dict]) -> str | None:
...
candidates: the same eligible (healthy, enabled, under capacity) member list every strategy receives, as a list of dicts:
| Key | Type | Meaning |
|---|---|---|
name | str | Cluster name (matches clusterGroups.<group>.members). |
engineType | str | e.g. trino, duckDb, starRocks (camelCase, matches EngineType). |
runningQueries | int | Current in-flight query count on this cluster. |
maxRunningQueries | int | Per-cluster capacity limit. |
candidatespreservesclusterGroups.<group>.membersorder —acquire_clusterfilters that list down to eligible members without reordering it. Return anamefromcandidatesto pick that cluster, orNoneto fall back to the first eligible candidate in that member order. Returning an unrecognized name also falls back and logs a warning naming the bad value; an in-script exception falls back the same way (also logged) — a broken script degrades tofailover-like behavior rather than dropping the query.- Load from a file instead of inline with
scriptFile: /path/to/select.py. If bothscriptandscriptFileare set,scriptFilewins (same precedence as the router-levelpythonScriptconfig). - Unlike
script/scriptFileat the router level, a missing/unreadablescriptFile(or apythonScriptstrategy with neither field set) is a hard failure at startup — the process won't boot with a misconfigured strategy. On a config reload, the same failure is non-fatal: that group keeps its last-known-good strategy (the last one that built successfully) and a warning is logged; only a group whose strategy has never successfully built falls back toroundRobin. Either way, one bad script edit can't take down an already-running proxy. - Unlike the router-level script,
select_clusterdoes not receive the query text or session context —ClusterGroupManager::acquire_clusteronly knows the group name at this stage. Placement logic here is necessarily load/topology-based, not query-content-based. - Runs off the async runtime (
spawn_blocking) so a slow or GIL-bound script can't block a Tokio worker thread. This doesn't make it free:spawn_blockinguses a shared, bounded thread pool, and the Python GIL serializes concurrent Python calls process-wide (shared with translation fixups) — a script that hangs or loops forever still ties up a blocking-pool thread and contends for the GIL indefinitely. Treat a runaway script as a process-level operational risk. - The script is compiled and
select_clusterlooked up once, on the first pick, then reused for the lifetime of that strategy instance (module-level state persists across calls, same as a normal Python module). A config reload rebuilds a fresh instance, so a script edit is always picked up on the next reload.
Health and runtime updates
Each ClusterState tracks health (is_healthy) and a running-query counter (running_queries). Background tasks in the main binary update both every 30 seconds:
- Health checks mark clusters unhealthy when probes fail; unhealthy clusters are excluded from acquisition.
- Reconciliation aligns
running_querieswith backend ground truth (engine introspection or optional custom SQL). In distributed mode, one replica publishes counts to Postgres (cluster_capacity_counters.running); other replicas read that value instead of querying backends.
For SaaS ADBC backends (Snowflake, Databricks, BigQuery, Redshift), QueryFlux uses built-in introspection that avoids waking auto-suspending warehouses (REST APIs or metadata queries instead of naive SELECT 1). Operators can override behavior with optional healthCheckQuery and reconcileQuery on the cluster config.
Distributed admission (Postgres persistence, multiple replicas) is separate from reconcile: maxRunningQueries is enforced fleet-wide via capacity leases (cluster_capacity_leases), not via the running counter. See Cluster variants, health checks & reconciliation.
Cluster variants expand one persisted config into multiple runtime clusters (base::warehouse-name), each with its own adapter and health/reconcile targets. Reference expanded names in group members.
See Cluster variants, health checks & reconciliation for YAML/API examples, resolution order, distributed CapacityStore behavior, and Studio UI.
The ClusterGroupManager trait also supports update_cluster (enable/disable, change max_running_queries) for admin-driven changes.
Frontend dispatch: async vs sync
After a group is chosen, the Trino HTTP handler (post_statement) branches:
- If the group is considered async-capable (e.g. Trino-style polling), it uses
dispatch_query: acquire cluster → translate →submit_query→ persist executing state → rewritenextUrito point back at QueryFlux. - Otherwise it uses
execute_to_sink: wait for capacity (backoff loop), translate, stream Arrow batches, and synthesize a Trino-compatible JSON response.
So routing and cluster selection are shared concepts; the result delivery shape depends on engine and frontend capabilities.
Mental model
- Routers answer: which pool (group) should handle this query?
- Cluster manager + strategy answer: which replica/instance in that pool?
- Translation (separate doc) then aligns SQL with that instance’s engine.
See system-map.md for the high-level query flow, persistence model, and component status, and query-translation.md for dialect conversion details.