Tune Replicator for Confluent Platform
Learn how to get maximum replication throughput with a minimal number of Connect workers and machines.
Assume you already have a Connect cluster with Replicator running. This documentation also assumes you’re running a dedicated Connect cluster for Replicator, that is, no other connector is running on the same cluster, and all resources of the cluster are available for Replicator.
Sizing Replicator cluster
Sizing a Replicator cluster comes down to two questions:
How many nodes does the cluster need, and how large should these nodes be?
How many Replicator tasks does the cluster need to run?
This section describes a method to determine the number of tasks, and then uses this information to determine the number of nodes you need.
The first step is to find out the throughput per task. You can do this by running just one task for a bit and checking the throughput. One way to do it is to fill a topic at the origin cluster with large amounts of data, start Replicator with a single task (configuration parameter tasks.max=1), and measure the rate at which event records are written to the destination cluster. This shows you what you can achieve with a single task.
Then, take the desired throughput - how many MBps do you need to replicate between the clusters - and divide by the throughput per task to determine the number of tasks you need.
For example, suppose that you ran Replicator with a single task and saw that it can replicate 30 MBps. You know that between all your topics, you need to replicate 100 MBps. This means you need at least four Replicator tasks.
There are two caveats to this formula:
You can’t have more tasks than partitions. If you find out that your throughput requires more tasks than you have partitions, you’ll need to add partitions. Having more partitions than tasks is fine; one task can consume from multiple partitions.
Having too many tasks on a single machine causes you to saturate the machine resources. Those are either the network or the CPU. If you are adding more tasks and the throughput doesn’t increase, you’ve probably saturated one of these resources. The easy solution is to add more nodes, so you’ll have more resources to use.
Recommended guidelines:
No more than 150 partitions per task
At most 20 tasks per worker
Getting more throughput from Replicator tasks
One way to increase throughput for Replicator is to run more tasks across more Connect worker nodes. The other is to try to squeeze more throughput from each Replicator task. In reality, your tuning effort is likely to be a combination of both: start by tuning a single task and then achieving the desired throughput with multiple tasks running in parallel.
When tuning performance of Connect tasks, you can try to make each task use less CPU or to use the network more efficiently. As a best practice, first check which of these is the bottleneck, and tune for the right resource. One way to know which one is the bottleneck is to add tasks to Replicator running on a single node until adding tasks no longer increases the throughput. If you see high CPU utilization, improve CPU utilization of Connect tasks. If CPU utilization is low, but adding tasks doesn’t improve throughput, tune the network utilization instead.
Improving CPU utilization of a Connect task
Make sure you are not seeing excessive garbage collection pauses by enabling and checking Java’s garbage collection log. Use G1GC and a reasonable heap size (4 GB is a good default to start with).
Make sure Replicator is configured with
key.converterandvalue.converterset toio.confluent.connect.replicator.util.ByteArrayConverter. If you havesrc.key.converterandsrc.value.converterconfigured, they should also be set toio.confluent.connect.replicator.util.ByteArrayConverter(the default value). This eliminates costly conversion of messages to Connect’s internal data format and back to bytes.Disable CRC checks. This isn’t recommended because it can lead to data corruption, but CRC checks for data integrity use CPU and can be disabled for improved performance by setting
src.consumer.check.crcs=falsein Replicator configuration.
Improving network utilization of a Connect task
Replicator typically reads or writes messages between two datacenters, which means latency is usually very high. High network latency can lead to reduced throughput. The suggestions below all aim at configuring the TCP stack to improve throughput in this environment. Note that this is different from configuring applications that are running within the same datacenter as the Apache Kafka® cluster. Inside a datacenter, you are usually trading off latency vs throughput and you can decide which of these to tune for. When replicating messages between two datacenters, high latency is typically a fact you must deal with.
Use a TCP throughput calculator to compute TCP buffer sizes based on the network link properties. Increase TCP buffer sizes to handle networks with high bandwidth and high latency (this includes most cross-datacenter links). You need to do this computation in two places:
Increase the application-level buffer size requested for sending and receiving requests. On producers and consumers, use
send.buffer.bytesandreceive.buffer.bytes. On brokers, usesocket.send.buffer.bytesandsocket.receive.buffer.bytes. If you don’t follow the second step, this level may still be silently overridden by the operating system.Increase operating system level TCP buffer size (
net.core.rmem_default,net.core.rmem_max,net.core.wmem_default,net.core.wmem_max,net.core.optmem_max). You can usesysctlfor testing in the current session, but you’ll need the values in/etc/sysctl.confto make them permanent.Important: Enable logging to double-check this actually took effect. In some instances, the operating system silently overrides or ignores settings.
log4j.logger.org.apache.kafka.common.network.Selector=DEBUG.
In addition, you can try two additional network optimizations:
Enable automatic window scaling (
sysctl -w net.ipv4.tcp_window_scaling=1or addnet.ipv4.tcp_window_scaling=1to/etc/sysctl.conf). If the latency justifies it, this allows the TCP buffer to grow beyond its usual maximum of 64 KB.Reduce the TCP slow start time (set
/proc/sys/net/ipv4/tcp_slow_start_after_idleto 0) to make the network connection reach its maximum capacity sooner.Increase producer
batch.size,linger.msand consumerfetch.max.bytes,fetch.min.bytes, andfetch.max.waitto improve throughput by reading and writing bigger batches.
Depending on the method you use to run Replicator, where and how you configure these settings varies:
When running Replicator as a connector, configure consumer settings in the Replicator configuration file and producer settings in the Connect worker configuration. For example, in this type of deployment, to tune consumer
fetch.max.bytes, setsrc.consumer.fetch.max.bytesin the Replicator configuration. To tune producerbatch.size, setproducer.batch.sizein the Connect worker configuration. In this scenario, these configuration files are used because Replicator is consuming from the origin cluster, and passing the records to the Connect worker. The worker is producing to the destination cluster.When running Replicator as an executable, configure consumer settings in
consumer.propertiesand producer settings inproducer.properties. In this scenario, these configuration files are used because theconsumer.propertiesfile controls the configurations for the embedded source consumer within Replicator, whileproducer.propertiescontrols the settings for the producer within the Connect framework.
Setting compression to improve performance with increased data loads
Increased data loads can impact Replicator performance with slow processing. You can address this by using data compression on both the source topic and the destination topic. Setting compression on the source topic alone is typically not sufficient. You must also explicitly set compression on the destination Replicator connector or worker configuration as follows:
To set compression on the Replicator connector, use
producer.override.compression.type(for example,"producer.override.compression.type":"lz4"). To learn more, see Override the Worker Configuration, and examples in Tutorial: Configure and Run Replicator for Confluent Platform as an Executable or Connector.To set compression in the worker, use
compression.type(for example,"compression.type":"lz4"). To learn more, see Kafka Producer Configuration Reference for Confluent Platform, and examples in Tutorial: Configure and Run Replicator for Confluent Platform as an Executable or Connector.