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) {