I've been reading through the Confluent Replicator's docs in order to evaluate its active-active multi-datacenter replication support and while the documentation is clear on the mechanism it uses to translate the consumer offset from one DC to another what remains unclear to me is how that mechanism deals with clock skew.
From what I can tell by reading the data recovery whitepaper the message timestamps used for translation are actually CreateTime timestamps coming from the producers. However, all the examples that are given for how the offset translation is performed seem to imply that the timestamps are monotonically increasing, e.g. this replicator blog post:
Why is this an issue? Because it is very hard (read: practically impossible) to guarantee monotonically increasing timestamps in distributed systems running on commodity hardware, and even with specialized hardware (e.g. AWS Time Sync Service) the best you can do is minimize the skew if you don't want to use a global time service similar to Google Spanner's TrueTime.
So, let's consider the simplest scenario where DC-1 and DC-2 have been set up with active-active replication using separate topics, e.g. topic dc1-foo in DC-1 is replicated to dc1-foo in DC2 and dc2-foo from DC2 is replicated to d2-foo in DC-1. Let's also assume that the datacenters have been operating correctly from the start but that the replication of dc-1-foo has experienced a minor hiccup which caused some records to be duplicated in DC-2, making the offset in DC-2 greater than the offset in DC-1.
Given the current comitted offset in DC-1 is O1 and DC-2 is O3, and given the contents for dc1-foo (in offset:TS format) is:
- DC-1: [O1:T3], [O2:T2], [O3:T3], [O4:T1], [O5:T3]
- DC-2: [O3:T3], [O4:T2], [O5:T3], [O6:T1], [O7:T3]
When the consumer in DC-1 commits...
- offset O2, is the DC-2 offset translated to O4?
- offset O5, is the DC-2 offset translated to O3, O5 or O7?
My guess would be that the translated offsets are O4 and O3 but just wanted to confirm.
