Tuning Percona ClusterSync for MongoDB Performance

Share this Post:

Percona ClusterSync for MongoDB (PCSM) synchronizes data from a source MongoDB cluster (including Atlas) to a target MongoDB cluster in two stages. It first creates a consistent starting point and clones the existing collections and indexes. At the same time, it captures source changes through MongoDB change streams. Once the initial copy is complete, PCSM applies those changes to the target until the target catches up with the source. 

PCSM’s latest architectural changes introduce several new parameters to improve initial clone speed and continuous replication throughput. This post aims to equip you with some ideas to help tune PCSM for maximum performance.

As usual, tune gradually and monitor source and target CPU, disk, network usage, connection counts, and replication lag after each change. Use Percona Monitoring and Management (PMM) to help you identify any bottlenecks. Don’t forget to keep an eye on the same metrics for the host running PCSM as well. 

Initial Clone Phase

The Initial cloning phase is mainly affected by collection parallelism, read and insert workers, segment size, and read batch size. You configure all of this while starting the pcsm process. 

Collection Parallelism

Parameter: clone-num-parallel-collections

The default is 2, which is quite conservative. This controls how many collections are cloned at the same time.

Increase it when:

  • The source has many collections.
  • The source and target have spare CPU and disk capacity.

It could be a good idea to start with 4 parallel collections for each CPU core and test higher values as long as the target is not saturated. Excessive parallelism will increase disk contention and connection usage.

Read and Insert Workers

Parameters:

  • clone-num-read-workers
  • clone-num-insert-workers

Defaults are:

  • Read workers: max(NumCPU/4, 1)
  • Insert workers: NumCPU × 2

Read workers fetch documents from the source. Insert workers write batches to the target.

Increase the number of read workers when the source reads limit throughput. Increase insert workers when the target has available write capacity and insert queues are building up. Target indexes are maintained during clone inserts, so additional insert workers can increase write contention.

Do not increase both aggressively at the same time. If target CPU or disk latency increases without a corresponding increase in throughput, reduce insert concurrency.

For continuous replication, PCSM does expose queue-depth metrics you can check in PMM:

Continuous Replication Phase

Continuous replication is mainly affected by worker count, change-stream batch size, queue sizes, bulk size, flush interval, and bulk-write mode.

Replication Workers

Parameter: repl-num-workers

The default is the number of CPUs.

Increase this setting when worker queues grow and the target still has available CPU, storage, and connections. A useful test value is twice the number of CPUs.

More workers will not improve throughput for high contention workloads (e.g. operations targeting the same set of documents). Events for the same document remain processed in order on a single worker at a time.

Change-Stream Batch Size

Parameter: repl-change-stream-batch-size

The default is 10,000 events.

Larger values reduce change-stream round trips during sustained workloads. You may experiment with values such as 20,000 or 50,000.

Larger batches use more memory and may increase burst latency.

Event and Worker Queues

Parameters:

  • repl-event-queue-size
  • repl-worker-queue-size

Both default to 5,000 events.

These queues absorb temporary bursts between change-stream reading, event routing, and target writes.

Test values such as 10,000 or 25,000 when queues fill briefly during write stalls. Larger queues do not improve sustained throughput if the target is permanently slower than the source.

The worker queue is allocated per worker, so its memory impact increases with repl-num-workers.

Bulk Size and Flush Interval

Parameters:

  • repl-bulk-ops-size
  • repl-worker-flush-interval

The default bulk size is 5,000 operations. The default flush interval is 1 second.

For higher throughput, test:

  • repl-bulk-ops-size=10000
  • repl-worker-flush-interval=2s

Larger bulks reduce MongoDB command overhead but may increase replication latency.

For lower latency, use a shorter flush interval, such as 100ms to 500ms.

Actual bulks remain limited by MongoDB’s wire-message size. Large documents may reach the byte limit before reaching the configured operation count.

Asynchronous Bulk Queue

Parameter: repl-worker-bulk-queue-size

The default is 3. This controls how many completed bulks can wait behind an in-flight write.

Test values of 5 or 10 when target writes are slower than event processing. Increasing this value can keep workers busy, but it also increases memory usage and may allow replication lag to grow when the target cannot keep up.

Collection-Level Bulk Writes

Parameter: use-collection-bulk-write

Benchmark this option against the default (false).

By default, a single bulk written can contain operations for multiple databases and collections. Collection-level bulk writes=true groups operations by namespace instead, and may improve workloads spread across many collections.

PCSM may select collection-level writes automatically for target versions without client-level bulk-write support.

Connection Pools

PCSM supports the MongoDB URI option maxPoolSize. If it is omitted, the driver default is 100 connections per client.

Example:

Connection-pool sizing becomes important when increasing clone or replication workers.

Start with a moderate explicit value and increase it only when workers are waiting for connections. Avoid unlimited pools unless there is a specific reason to remove the limit.

Compression

PCSM supports separate compressor settings for source and target connections.

Parameters:

  • source-client-compressors
  • target-client-compressors

The default preference is:

snappy, zstd, zlib

PCSM matches the order the compressors are negotiated by MongoDB by default. Typically  snappy is used since it has lowest overhead but the compression is not as good as zstd or zlib. You may try the other compressors if optimizing for least possible network bandwidth usage.

Recommended Tuning Order 

Now that we know the most important parameters, it makes sense to run some tests and keep iterating to find the optimal values. 

The recommended workflow is:

  1. Establish a baseline using the default settings.
  2. Set an appropriate maxPoolSize on the source and target URIs.
  3. Find the ideal clone collection parallelism.
  4. Play with clone read and insert workers independently.
  5. Tune the replication worker count. 

In the end, the fastest configuration is the one that keeps PCSM and MongoDB busy without saturating storage, CPU, memory, or connections on either source or target. 

0 0 votes
Article Rating
Subscribe
Notify of
guest

0 Comments
Oldest
Newest Most Voted

Far
Enough.

Said no pioneer ever.
MySQL, PostgreSQL, InnoDB, MariaDB, MongoDB and Kubernetes are trademarks for their respective owners.
© 2026 Percona All Rights Reserved