This is an automated email from the ASF dual-hosted git repository.
rong 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 46fc1e0a6fa Pipe: Fixed the tsFile parsing & write-back-sink auto
create db bug (#15240)
46fc1e0a6fa is described below
commit 46fc1e0a6fad9614d871dbe67b3ce137be4ea33a
Author: Caideyipi <[email protected]>
AuthorDate: Tue Apr 1 00:06:14 2025 +0800
Pipe: Fixed the tsFile parsing & write-back-sink auto create db bug (#15240)
---
.../pipe/it/dual/tablemodel/TableModelUtils.java | 8 ++--
.../pipe/it/single/IoTDBPipePermissionIT.java | 43 ++++++++++++++++++++++
.../agent/task/connection/PipeEventCollector.java | 4 +-
.../subtask/processor/PipeProcessorSubtask.java | 15 +++++++-
.../protocol/writeback/WriteBackConnector.java | 10 +++++
.../common/tsfile/PipeTsFileInsertionEvent.java | 4 +-
6 files changed, 75 insertions(+), 9 deletions(-)
diff --git
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/tablemodel/TableModelUtils.java
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/tablemodel/TableModelUtils.java
index 17b470296cb..cc154a9d943 100644
---
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/tablemodel/TableModelUtils.java
+++
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/tablemodel/TableModelUtils.java
@@ -100,11 +100,11 @@ public class TableModelUtils {
public static boolean insertData(
final String dataBaseName,
final String tableName,
- final int start,
- final int end,
+ final int startInclusive,
+ final int endExclusive,
final BaseEnv baseEnv) {
- List<String> list = new ArrayList<>(end - start + 1);
- for (int i = start; i < end; ++i) {
+ List<String> list = new ArrayList<>(endExclusive - startInclusive + 1);
+ for (int i = startInclusive; i < endExclusive; ++i) {
list.add(
String.format(
"insert into %s (s0, s3, s2, s1, s4, s5, s6, s7, s8, s9, s10,
s11, time) values ('t%s','t%s','t%s','t%s','%s', %s.0, %s, %s, %d, %d.0, '%s',
'%s', %s)",
diff --git
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
index 005eb49afef..64afeeb594d 100644
---
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
+++
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
@@ -32,6 +32,7 @@ import org.junit.runner.RunWith;
import java.sql.Connection;
import java.sql.SQLException;
import java.sql.Statement;
+import java.util.Arrays;
import static org.junit.Assert.fail;
@@ -154,4 +155,46 @@ public class IoTDBPipePermissionIT extends
AbstractPipeSingleIT {
TableModelUtils.assertCountData("test", "test", 100, env);
}
+
+ @Test
+ public void testSinkPermissionWithHistoricalDataAndTablePattern() {
+ TableModelUtils.createDataBaseAndTable(env, "test", "test1");
+ TableModelUtils.createDataBaseAndTable(env, "test1", "test1");
+ TableModelUtils.createDataBaseAndTable(env, "test", "test");
+ TableModelUtils.createDataBaseAndTable(env, "test1", "test");
+
+ if (!TestUtils.tryExecuteNonQueriesWithRetry(
+ "test",
+ BaseEnv.TABLE_SQL_DIALECT,
+ env,
+ Arrays.asList(
+ "create user thulab 'passwd'", "grant INSERT on test.test1 to user
thulab"))) {
+ return;
+ }
+
+ // Write some data
+ if (!TableModelUtils.insertData("test1", "test", 0, 100, env)) {
+ return;
+ }
+
+ if (!TableModelUtils.insertData("test1", "test1", 0, 100, env)) {
+ return;
+ }
+
+ // Use current session, user is root
+ try (final Connection connection =
env.getConnection(BaseEnv.TABLE_SQL_DIALECT);
+ final Statement statement = connection.createStatement()) {
+ statement.execute(
+ "create pipe a2b "
+ + "with source ('database'='test1', 'table'='test1') "
+ + "with processor('processor'='rename-database-processor',
'processor.new-db-name'='test') "
+ + "with sink ('sink'='write-back-sink', 'username'='thulab',
'password'='passwd')");
+ } catch (final SQLException e) {
+ e.printStackTrace();
+ fail("Create pipe without user shall succeed if use the current
session");
+ }
+
+ TableModelUtils.assertCountData("test", "test", 0, env);
+ TableModelUtils.assertCountData("test", "test1", 100, env);
+ }
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
index 6e265711ad0..3bc4553c852 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
@@ -135,9 +135,7 @@ public class PipeEventCollector implements EventCollector {
return;
}
- if (!forceTabletFormat
- && !sourceEvent.shouldParse4Privilege()
- && canSkipParsing4TsFileEvent(sourceEvent)) {
+ if (!forceTabletFormat && canSkipParsing4TsFileEvent(sourceEvent)) {
collectEvent(sourceEvent);
return;
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
index f611b18ad61..40352766630 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
@@ -31,6 +31,7 @@ import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
import org.apache.iotdb.db.pipe.agent.task.connection.PipeEventCollector;
import org.apache.iotdb.db.pipe.event.UserDefinedEnrichedEvent;
import org.apache.iotdb.db.pipe.event.common.heartbeat.PipeHeartbeatEvent;
+import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
import
org.apache.iotdb.db.pipe.metric.overview.PipeDataNodeRemainingEventAndTimeMetrics;
import org.apache.iotdb.db.pipe.metric.processor.PipeProcessorMetrics;
import org.apache.iotdb.db.pipe.processor.pipeconsensus.PipeConsensusProcessor;
@@ -141,7 +142,19 @@ public class PipeProcessorSubtask extends
PipeReportableSubtask {
pipeProcessor.process((TabletInsertionEvent) event,
outputEventCollector);
PipeProcessorMetrics.getInstance().markTabletEvent(taskID);
} else if (event instanceof TsFileInsertionEvent) {
- pipeProcessor.process((TsFileInsertionEvent) event,
outputEventCollector);
+ // We have to parse the privilege first, to avoid passing
no-privilege data to processor
+ if (event instanceof PipeTsFileInsertionEvent
+ && ((PipeTsFileInsertionEvent) event).shouldParse4Privilege()) {
+ try (final PipeTsFileInsertionEvent tsFileInsertionEvent =
+ (PipeTsFileInsertionEvent) event) {
+ for (final TabletInsertionEvent tabletInsertionEvent :
+ tsFileInsertionEvent.toTabletInsertionEvents()) {
+ pipeProcessor.process(tabletInsertionEvent,
outputEventCollector);
+ }
+ }
+ } else {
+ pipeProcessor.process((TsFileInsertionEvent) event,
outputEventCollector);
+ }
PipeProcessorMetrics.getInstance().markTsFileEvent(taskID);
PipeDataNodeRemainingEventAndTimeMetrics.getInstance()
.markTsFileCollectInvocationCount(
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/writeback/WriteBackConnector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/writeback/WriteBackConnector.java
index 201f776b9e6..0bc4f76d25c 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/writeback/WriteBackConnector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/writeback/WriteBackConnector.java
@@ -374,6 +374,16 @@ public class WriteBackConnector implements PipeConnector {
return;
}
+ try {
+ Coordinator.getInstance()
+ .getAccessControl()
+ .checkCanCreateDatabase(session.getUsername(), database);
+ } catch (final AccessDeniedException e) {
+ // Auto create failed, we still check if there are existing databases
+ // If there are not, this will be removed by catching database not
exists exception
+ ALREADY_CREATED_DATABASES.add(database);
+ return;
+ }
final TDatabaseSchema schema = new TDatabaseSchema(new
TDatabaseSchema(database));
schema.setIsTableModel(true);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
index 04cf1d58e59..3077b05237e 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
@@ -672,7 +672,9 @@ public class PipeTsFileInsertionEvent extends
PipeInsertionEvent
startTime,
endTime,
pipeTaskMeta,
- userName,
+ // Do not parse privilege if it should not be parsed
+ // To avoid renaming of the tsFile database
+ shouldParse4Privilege ? userName : null,
this)
.provide());
return eventParser.get();