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

Caideyipi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new ada62cb1048 [Pipe] Fix conversion task ID collision after leader 
switch (#18494)
ada62cb1048 is described below

commit ada62cb1048b85507cbc478edd4ebd50af037ac5
Author: Caideyipi <[email protected]>
AuthorDate: Wed Aug 19 15:14:25 2026 +0800

    [Pipe] Fix conversion task ID collision after leader switch (#18494)
---
 .../request/PipeTransferTsFileSealWithModReq.java  | 11 ++++----
 .../pipe/sink/PipeDataNodeThriftRequestTest.java   | 30 ++++++++++++++++++++++
 2 files changed, 35 insertions(+), 6 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTsFileSealWithModReq.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTsFileSealWithModReq.java
index 1328ae59ba4..efd6db23a74 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTsFileSealWithModReq.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTsFileSealWithModReq.java
@@ -162,12 +162,11 @@ public class PipeTransferTsFileSealWithModReq extends 
PipeTransferFileSealReqV2
         } catch (final UnsupportedOperationException ignored) {
           appendStablePart(eventIdentity, UNSUPPORTED_REPLICATE_INDEX);
         }
-        if (event.getCommitterKey() == null) {
-          try {
-            appendStablePart(eventIdentity, 
String.valueOf(event.getProgressIndex()));
-          } catch (final UnsupportedOperationException ignored) {
-            appendStablePart(eventIdentity, UNSUPPORTED_PROGRESS_INDEX);
-          }
+        // Commit ids are local to a DataNode and may collide after a leader 
change.
+        try {
+          appendStablePart(eventIdentity, 
String.valueOf(event.getProgressIndex()));
+        } catch (final UnsupportedOperationException ignored) {
+          appendStablePart(eventIdentity, UNSUPPORTED_PROGRESS_INDEX);
         }
         eventIdentities.add(eventIdentity.toString());
       }
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
index 915860fd239..143d195d284 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
@@ -19,7 +19,10 @@
 
 package org.apache.iotdb.db.pipe.sink;
 
+import org.apache.iotdb.commons.consensus.index.impl.IoTProgressIndex;
 import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
+import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
 import 
org.apache.iotdb.commons.pipe.sink.payload.thrift.common.PipeTransferHandshakeConstant;
 import 
org.apache.iotdb.commons.pipe.sink.payload.thrift.request.IoTDBSinkRequestVersion;
 import 
org.apache.iotdb.commons.pipe.sink.payload.thrift.request.PipeRequestType;
@@ -71,6 +74,7 @@ import org.apache.tsfile.write.schema.IMeasurementSchema;
 import org.apache.tsfile.write.schema.MeasurementSchema;
 import org.junit.Assert;
 import org.junit.Test;
+import org.mockito.Mockito;
 
 import java.io.DataOutputStream;
 import java.io.IOException;
@@ -1215,6 +1219,32 @@ public class PipeDataNodeThriftRequestTest {
     Assert.assertFalse(deserialized.shouldAsyncLoadOnTypeMismatch());
   }
 
+  @Test
+  public void 
testPipeTransferTsFileSealConversionTaskIdDistinguishesProgressIndexes() {
+    final CommitterKey committerKey = new CommitterKey("pipe", 1L, 1, 0);
+    final EnrichedEvent firstEvent = Mockito.mock(EnrichedEvent.class);
+    Mockito.when(firstEvent.getCommitterKey()).thenReturn(committerKey);
+    
Mockito.when(firstEvent.getCommitIds()).thenReturn(Collections.singletonList(1L));
+    Mockito.when(firstEvent.getProgressIndex()).thenReturn(new 
IoTProgressIndex(1, 1L));
+
+    final EnrichedEvent secondEvent = Mockito.mock(EnrichedEvent.class);
+    Mockito.when(secondEvent.getCommitterKey()).thenReturn(committerKey);
+    
Mockito.when(secondEvent.getCommitIds()).thenReturn(Collections.singletonList(1L));
+    Mockito.when(secondEvent.getProgressIndex()).thenReturn(new 
IoTProgressIndex(1, 100L));
+
+    final String firstTaskId =
+        PipeTransferTsFileSealWithModReq.generateConversionTaskId(
+            "sink-task", Collections.singletonList(firstEvent), "root.db", 0);
+    Assert.assertEquals(
+        firstTaskId,
+        PipeTransferTsFileSealWithModReq.generateConversionTaskId(
+            "sink-task", Collections.singletonList(firstEvent), "root.db", 0));
+    Assert.assertNotEquals(
+        firstTaskId,
+        PipeTransferTsFileSealWithModReq.generateConversionTaskId(
+            "sink-task", Collections.singletonList(secondEvent), "root.db", 
0));
+  }
+
   @Test
   public void 
testPipeTransferTsFileSealWithModReqFromLegacyV13BodyWithoutDatabaseName()
       throws IOException {

Reply via email to