From 3dcc142730f5a4aecde8254fe2abe5f52270a967 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Thu, 6 Aug 2026 14:56:48 +0800 Subject: [PATCH 1/2] [To dev/1.3] Fix write-back sink tree target database case (#18398) --- .../it/manual/IoTDBPipeWriteBackSinkIT.java | 24 +++++++++++++++---- .../protocol/writeback/WriteBackSink.java | 3 +-- 2 files changed, 20 insertions(+), 7 deletions(-) diff --git a/integration-test/src/test/java/org/apache/iotdb/pipe/it/manual/IoTDBPipeWriteBackSinkIT.java b/integration-test/src/test/java/org/apache/iotdb/pipe/it/manual/IoTDBPipeWriteBackSinkIT.java index dda55ad4871d4..257035238b775 100644 --- a/integration-test/src/test/java/org/apache/iotdb/pipe/it/manual/IoTDBPipeWriteBackSinkIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/pipe/it/manual/IoTDBPipeWriteBackSinkIT.java @@ -44,13 +44,27 @@ public class IoTDBPipeWriteBackSinkIT extends AbstractPipeDualManualIT { @Test public void testWriteBackSinkWithTargetDatabaseForTreeModel() throws Exception { + testWriteBackSinkWithTargetDatabaseForTreeModel("root.target.db"); + } + + @Test + public void testWriteBackSinkPreservesTreeModelTargetDatabaseCase() throws Exception { + testWriteBackSinkWithTargetDatabaseForTreeModel("TargetDB"); + } + + private void testWriteBackSinkWithTargetDatabaseForTreeModel(final String targetDatabase) + throws Exception { + final String qualifiedTargetDatabase = + targetDatabase.startsWith("root.") ? targetDatabase : "root." + targetDatabase; TestUtils.executeNonQueries( senderEnv, Arrays.asList( "create database root.source", "create timeseries root.source.d1.s1 with datatype=INT32,encoding=PLAIN", - "create database root.target.db", - "create timeseries root.target.db.d1.s1 with datatype=INT32,encoding=PLAIN"), + "create database " + qualifiedTargetDatabase, + "create timeseries " + + qualifiedTargetDatabase + + ".d1.s1 with datatype=INT32,encoding=PLAIN"), null); try (final SyncConfigNodeIServiceClient client = @@ -65,7 +79,7 @@ public void testWriteBackSinkWithTargetDatabaseForTreeModel() throws Exception { sourceAttributes.put("user", "root"); sinkAttributes.put("sink", "write-back-sink"); - sinkAttributes.put("sink.database", "root.target.db"); + sinkAttributes.put("sink.database", targetDatabase); sinkAttributes.put("user", "root"); final TSStatus status = @@ -89,8 +103,8 @@ public void testWriteBackSinkWithTargetDatabaseForTreeModel() throws Exception { TestUtils.assertDataEventuallyOnEnv( senderEnv, - "select * from root.target.db.**", - "Time,root.target.db.d1.s1,", + "select * from " + qualifiedTargetDatabase + ".**", + "Time," + qualifiedTargetDatabase + ".d1.s1,", Collections.unmodifiableSet(new HashSet<>(Arrays.asList("1,1,", "2,2,")))); } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/writeback/WriteBackSink.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/writeback/WriteBackSink.java index f5a7ad2ff31a3..e3ebd879eda67 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/writeback/WriteBackSink.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/writeback/WriteBackSink.java @@ -64,7 +64,6 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; -import java.util.Locale; import java.util.Objects; import static org.apache.iotdb.commons.conf.IoTDBConstant.MAX_DATABASE_NAME_LENGTH; @@ -115,7 +114,7 @@ private static String validateTargetDatabase(final String targetDatabase) { return validateAndNormalizeTreeModelDatabaseName( IoTDBConstant.PATH_ROOT + IoTDBConstant.PATH_SEPARATOR - + trimmedTargetDatabase.toLowerCase(Locale.ENGLISH)); + + trimmedTargetDatabase); } catch (final Exception e) { throw new PipeException( String.format("The target database %s is invalid.", targetDatabase), e); From ba4176c8808556c0c21135250e53bee79a2325a8 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Thu, 6 Aug 2026 15:04:02 +0800 Subject: [PATCH 2/2] Update WriteBackSink.java --- .../iotdb/db/pipe/sink/protocol/writeback/WriteBackSink.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/writeback/WriteBackSink.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/writeback/WriteBackSink.java index e3ebd879eda67..abf6ccda00843 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/writeback/WriteBackSink.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/writeback/WriteBackSink.java @@ -112,9 +112,7 @@ private static String validateTargetDatabase(final String targetDatabase) { try { PathUtils.checkAndReturnSingleMeasurement(trimmedTargetDatabase); return validateAndNormalizeTreeModelDatabaseName( - IoTDBConstant.PATH_ROOT - + IoTDBConstant.PATH_SEPARATOR - + trimmedTargetDatabase); + IoTDBConstant.PATH_ROOT + IoTDBConstant.PATH_SEPARATOR + trimmedTargetDatabase); } catch (final Exception e) { throw new PipeException( String.format("The target database %s is invalid.", targetDatabase), e);