MonPG Engineering avatar MonPG Engineering Engineering Team MariaDB 4 min read

MariaDB Galera Flow Control: Why One Slow Node Throttles Your Entire Cluster

In a 3-node MariaDB Galera cluster, an automated backup on Node 3 caused write throughput across the whole cluster to collapse from 2,800 to 12 QPS. Here is how Galera flow control works and how to design resilient replication topologies.

MariaDB

We operated a three-node MariaDB Galera cluster distributed across three availability zones, convinced that synchronous multi-master replication provided high availability and zero data loss.

Then, during an off-peak backup window, write throughput across the entire cluster suddenly collapsed from 2,800 queries per second to 12 queries per second. Nodes 1 and 2 had zero CPU load and ample disk I/O, yet client write queries were hanging indefinitely.

The cause was an I/O freeze on Node 3 triggered by a cloud disk snapshot. Galera’s Flow Control mechanism intervened and ordered Nodes 1 and 2 to halt client writes to prevent Node 3 from falling too far behind.

Certification-based replication and the receive queue

MariaDB Galera Cluster uses a certification-based replication model. When a transaction commits on Node 1, it broadcasts its write-set (the set of changed rows and primary keys) to all other nodes in the cluster via total order multicast.

Each node puts the incoming write-set into its local receive queue (wsrep_local_recv_queue) and executes a certification test to verify that the write-set does not conflict with any concurrent local transactions.

Once certification succeeds, the transaction is guaranteed to commit. However, the physical application of that write-set to the InnoDB storage engine is executed asynchronously by background applier threads (wsrep_slave_threads).

How Flow Control freezes the cluster

If Node 3 encounters a slow disk, an intensive backup, or CPU starvation, its background applier threads cannot keep up with the incoming write-sets. As a result, its wsrep_local_recv_queue begins to swell.

To prevent the lagging node from exhausting memory or falling irreversibly out of sync, Galera implements Flow Control. When the number of write-sets in any node’s receive queue exceeds wsrep_flow_control_upper_limit (default typically 16 or 30), that node broadcasts an FC_PAUSE message to the cluster.

Upon receiving FC_PAUSE, all other nodes in the cluster immediately suspend processing incoming client write transactions. The entire cluster slows down to the speed of its slowest node.

-- Checking Galera flow control activity and queue depth
SHOW STATUS LIKE 'wsrep_flow_control_paused';
SHOW STATUS LIKE 'wsrep_local_recv_queue%';
SHOW STATUS LIKE 'wsrep_slave_threads';

The wsrep_flow_control_paused status variable returns a value between 0.0 and 1.0 representing the fraction of time the cluster has been paused since the last status check. Anything above 0.05 indicates severe cluster degradation.

Architectural patterns to prevent flow control cascades

To build a resilient Galera topology that survives node slowdowns, implement these architectural boundaries:

-- Safely taking a node out of flow control before backup
SET GLOBAL wsrep_desync = ON;
-- Run backup or maintenance here...
SET GLOBAL wsrep_desync = OFF;
  • Dedicated single-writer routing: Route all application write traffic to a single designated primary node using HAProxy or ProxySQL. While Galera supports multi-master writes, multi-node writes dramatically increase certification conflict rates.
  • Tune applier concurrency: Ensure wsrep_slave_threads is configured to utilize available CPU cores (typically set to 2x to 4x the number of CPU cores) so applying write-sets does not become CPU-bound.
  • Desync nodes during maintenance: Before taking backups, running intensive DDL, or generating snapshots on a node, explicitly take it out of the flow control quorum using SET GLOBAL wsrep_desync = ON;. The node will continue replicating but will not broadcast pause signals to its peers.

The practical standard

High-level database architecture is not about drawn boxes on an infrastructure diagram. It is about how the engine manages shared resources under concurrency — memory, latches, write-ahead logs, and lock tables. When things break at 2 AM, the fix is rarely adding another replica or throwing more CPU at the host. The fix is understanding the underlying resource bottleneck, measuring the exact wait event, and applying the architectural constraint that makes the system predictable.