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 {