10.2 The realities of cross-datacenter communication
"The solutions we'll discuss may seem overly complicated without understanding that they represent TRADE-OFFS IN THE FACE OF SPECIFIC NETWORK CONDITIONS."
- High latencies. “increases as the Distance and the Number of network hops between the two clusters increase”
- Limited bandwidth. “WANs typically have Far lower available bandwidth than what you’ll see inside a single datacenter, and the available bandwidth Can vary from minute to minute. In addition, Higher latencies make it more challenging to utilize all the available bandwidth.”
- Higher costs. “partly because the bandwidth is limited and adding bandwidth can be Prohibitively expensive, and also because of The prices vendors charge for transferring data among datacenters, regions, and clouds”
⚠️ Why you can't just stretch a normal cluster over a WAN
"Apache Kafka's brokers and clients were designed, developed, tested, and tuned, ALL WITHIN A SINGLE DATACENTER. We assumed low latency and high bandwidth between brokers and clients. This is apparent in THE DEFAULT TIMEOUTS AND SIZING OF VARIOUS BUFFERS. For this reason, it is NOT RECOMMENDED (except in specific cases) to install some Kafka brokers in one datacenter and others in another datacenter."
💡 The central design insight: consume remotely, don't produce remotely
Three possible cross-DC communication patterns, eliminated in order:
Why ③ is safe — this is the argument to internalize:
"in the event of a network partition that prevents a consumer from reading data, THE RECORDS REMAIN SAFE INSIDE THE KAFKA BROKERS until communications resume and consumers can read them. THERE IS NO RISK OF ACCIDENTAL DATA LOSS due to network partitions."
(Ch. 3 §12’s cascade, over a WAN.)
If you do have to produce remotely: "you need to account for higher latency and the potential for more network errors. You can handle the errors by increasing the number of producer retries, and handle the higher latency by increasing the size of the buffers that hold records between attempts to send them."
And one more efficiency argument:
"because bandwidth is limited, if there are MULTIPLE applications in one datacenter that need to read data from Kafka brokers in another datacenter, we prefer to install a Kafka cluster in each datacenter and MIRROR THE NECESSARY DATA BETWEEN THEM ONCE rather than have MULTIPLE applications consume the same data across the WAN."
The three guiding principles
- NO LESS THAN ONE CLUSTER PER DATACENTER.
- REPLICATE EACH EVENT EXACTLY ONCE (barring retries due to errors) between each pair of datacenters.
- WHEN POSSIBLE, CONSUME FROM A REMOTE DATACENTER RATHER THAN PRODUCE TO A REMOTE DATACENTER.