SGShardedCluster: query routers are registered with shouldhaveshards = true and may get shards of the distributed tables
Summary
The query routers of a citus SGShardedCluster are registered in pg_dist_node by the Patroni of
the coordinator with citus_add_node, that always registers a non coordinator node with
shouldhaveshards = true. The pg_cron job update-query-routers-flags corrects it to false up to
30 seconds later, but a create_distributed_table executed in between places shards on the query
routers. The correction does not move those placements, so after it pg_dist_node shows every
query router with shouldhaveshards = false while some of them hold shards.
The shards of a table created afterwards with the same shard count are colocated with the misplaced ones, so they land on the query routers too, even once the flag has been corrected.
A trigger on pg_dist_node can not solve it: Citus writes that table with
CatalogTupleInsert/CatalogTupleUpdate, that do not fire triggers.
Proposed solution
The query router is registered in pg_dist_node by the coordinator, as an inactive node without
shards, before its Patroni is allowed to start. Patroni of the coordinator then finds the group
already registered and only replaces the host of the node (citus_update_node), that keeps
shouldhaveshards = false. The implementation is made of generic features glued by the
SGShardedCluster controller:
SGCluster.spec.configurations.patroni.startGateLabels: an object of label keys and values. When set, the cluster-controller starts Patroni only once theSGClusterhas the same labels.SGScript.spec.scripts[].setValue: whentruethe single text value returned by the query is stored inSGCluster.status.managedSql.scripts[].scripts[].value.SGScript.spec.scripts[].cron: the script is executed following the specified schedule (parsed withcom.cronutils.parser.CronParser, Quartz definition) instead of the current rules.- The coordinator
SGScriptgenerated for a citusSGShardedCluster:- registers the missing query router groups as inactive nodes with
shouldhaveshards = false, fixes the flag of the query routers, activates the registered query routers that can be reached and replicates the reference tables to them, on a schedule; - removes, when
SGShardedCluster.spec.configurations.citus.enableNodeAutoRemovalistrue(disabled by default), the nodes of the groups removed by decreasingworkers.clustersorqueryRouterClustersthat hold no shard of a distributed table and can not be reached (see #3244); - reads, on a schedule, the list of
groupidinpg_dist_nodewithshouldhaveshards = falseinto the coordinatorSGClusterstatus (setValue); - removes the pg_cron jobs
update-query-routers-flagsandupdate-query-routers-nodes, that it replaces.
- registers the missing query router groups as inactive nodes with
- The schedule is configurable with
SGShardedCluster.spec.configurations.citus.updateNodeInterval(ISO 8601 duration, defaults toPT10S). - The
SGShardedClustersets the start gate label instartGateLabelsof the query routerSGClusters and sets the same label on the query routerSGClusterwhose group is in the list read in the coordinatorSGClusterstatus. - The managed SQL is reconciled by its own reconciliation cycle of the cluster-controller
(
ManagedSqlReconciliationCycle), that is the only one executing theSGScriptentries and updatingSGCluster.status.managedSql, so thatSGScripts are never executed in parallel and the updates of that section never race. It runs when theSGClusteror the leader of the cluster changes, each reconciliation period and whenever a scheduledSGScriptentry is due. - No check of the operator version of the coordinator is needed: the cluster-controller of the
query router and of the coordinator Pods are only upgraded when their Pods are restarted, so a
query router waiting for its start gate label starts once the coordinator is restarted with a
cluster-controller that supports
setValueandcron.
Checking the placements
SELECT count(*)
FROM pg_dist_placement p
JOIN pg_dist_shard s USING (shardid)
JOIN pg_dist_partition d ON d.logicalrelid = s.logicalrelid
WHERE p.groupid >= 1024 AND d.partmethod <> 'n'; -- must be 0Reference tables (partmethod = 'n') are replicated to every node, query routers included, and that
is correct. The misplaced shards can be moved away with rebalance_table_shards() once the flag is
false.