This is an automated email from the ASF dual-hosted git repository.

mimaison pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new eec3aad481a MINOR: Refactor downstream offset translation in 
OffsetSync (#22543)
eec3aad481a is described below

commit eec3aad481afef4e038dd564deb000f558d75d2b
Author: Russole Chen <[email protected]>
AuthorDate: Thu Jun 25 21:41:14 2026 +0800

    MINOR: Refactor downstream offset translation in OffsetSync (#22543)
    
    This is a small refactor to resolve the TODO in OffsetSyncWriter
    
    Reviewers: Mickael Maison <[email protected]>
---
 .../src/main/java/org/apache/kafka/connect/mirror/OffsetSync.java | 8 ++++++++
 .../java/org/apache/kafka/connect/mirror/OffsetSyncStore.java     | 6 +++---
 .../java/org/apache/kafka/connect/mirror/OffsetSyncWriter.java    | 5 ++---
 3 files changed, 13 insertions(+), 6 deletions(-)

diff --git 
a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSync.java 
b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSync.java
index 6e366573cfe..3241e72a132 100644
--- 
a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSync.java
+++ 
b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSync.java
@@ -45,6 +45,14 @@ public record OffsetSync(TopicPartition topicPartition, long 
upstreamOffset, lon
             topicPartition, upstreamOffset, downstreamOffset);
     }
 
+    static long downstreamOffsetAfterSync(long downstreamOffset) {
+        return downstreamOffset + 1;
+    }
+
+    long translateDownstream(long upstreamOffset) {
+        return upstreamOffset == this.upstreamOffset ? downstreamOffset : 
downstreamOffsetAfterSync(downstreamOffset);
+    }
+
     ByteBuffer serializeValue() {
         Struct struct = valueStruct();
         ByteBuffer buffer = ByteBuffer.allocate(VALUE_SCHEMA.sizeOf(struct));
diff --git 
a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSyncStore.java
 
b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSyncStore.java
index 635ab732773..0829a9adbac 100644
--- 
a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSyncStore.java
+++ 
b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSyncStore.java
@@ -150,13 +150,13 @@ public class OffsetSyncStore implements AutoCloseable {
             //          | /
             //          vv
             // target |-sg----r-----|
-            long upstreamStep = upstreamOffset == 
offsetSync.get().upstreamOffset() ? 0 : 1;
+            long translatedDownstreamOffset = 
offsetSync.get().translateDownstream(upstreamOffset);
             log.debug("translateDownstream({},{},{}): Translated {} (relative 
to {})",
                     group, sourceTopicPartition, upstreamOffset,
-                    offsetSync.get().downstreamOffset() + upstreamStep,
+                    translatedDownstreamOffset,
                     offsetSync.get()
             );
-            return OptionalLong.of(offsetSync.get().downstreamOffset() + 
upstreamStep);
+            return OptionalLong.of(translatedDownstreamOffset);
         } else {
             log.debug("translateDownstream({},{},{}): Skipped (offset sync not 
found)",
                     group, sourceTopicPartition, upstreamOffset);
diff --git 
a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSyncWriter.java
 
b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSyncWriter.java
index 1a5ef6cc458..c696333ec82 100644
--- 
a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSyncWriter.java
+++ 
b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/OffsetSyncWriter.java
@@ -169,9 +169,8 @@ class OffsetSyncWriter implements AutoCloseable {
         boolean update(long upstreamOffset, long downstreamOffset) {
             // Emit an offset sync if any of the following conditions are true
             boolean noPreviousSyncThisLifetime = lastSyncDownstreamOffset == 
-1L;
-            // the OffsetSync::translateDownstream method will translate this 
offset 1 past the last sync, so add 1.
-            // TODO: share common implementation to enforce this relationship
-            boolean translatedOffsetTooStale = downstreamOffset - 
(lastSyncDownstreamOffset + 1) >= maxOffsetLag;
+            // OffsetSync translates offsets after the last sync to one 
downstream offset past the sync.
+            boolean translatedOffsetTooStale = downstreamOffset - 
OffsetSync.downstreamOffsetAfterSync(lastSyncDownstreamOffset) >= maxOffsetLag;
             boolean skippedUpstreamRecord = upstreamOffset - 
previousUpstreamOffset != 1L;
             boolean truncatedDownstreamTopic = downstreamOffset < 
previousDownstreamOffset;
             if (noPreviousSyncThisLifetime || translatedOffsetTooStale || 
skippedUpstreamRecord || truncatedDownstreamTopic) {

Reply via email to