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 2eb060d1854 Pipe: Fix the Pipe receiver's LoadFile and 
InsertRowsStatement NPE issues (#13871)
2eb060d1854 is described below

commit 2eb060d185441e8595998eb5c53d8ffc0b646fa2
Author: Zhenyu Luo <[email protected]>
AuthorDate: Thu Oct 24 12:06:08 2024 +0800

    Pipe: Fix the Pipe receiver's LoadFile and InsertRowsStatement NPE issues 
(#13871)
---
 .../evolvable/request/PipeTransferTabletBinaryReqV2.java    | 13 +++++++++++++
 .../request/PipeTransferTabletInsertNodeReqV2.java          | 13 +++++++++++++
 .../apache/iotdb/db/protocol/session/SessionManager.java    |  4 ++--
 3 files changed, 28 insertions(+), 2 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java
index d3a241dd880..8b23c3d9392 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java
@@ -27,6 +27,8 @@ import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNod
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsNode;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertTabletNode;
 import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
 
 import org.apache.tsfile.utils.PublicBAOS;
 import org.apache.tsfile.utils.ReadWriteIOUtils;
@@ -34,6 +36,7 @@ import org.apache.tsfile.utils.ReadWriteIOUtils;
 import java.io.DataOutputStream;
 import java.io.IOException;
 import java.nio.ByteBuffer;
+import java.util.List;
 import java.util.Objects;
 
 public class PipeTransferTabletBinaryReqV2 extends PipeTransferTabletBinaryReq 
{
@@ -71,6 +74,16 @@ public class PipeTransferTabletBinaryReqV2 extends 
PipeTransferTabletBinaryReq {
 
     // Table model
     statement.setWriteToTable(true);
+    if (statement instanceof InsertRowsStatement) {
+      List<InsertRowStatement> rowStatements =
+          ((InsertRowsStatement) statement).getInsertRowStatementList();
+      if (rowStatements != null && !rowStatements.isEmpty()) {
+        for (InsertRowStatement insertRowStatement : rowStatements) {
+          insertRowStatement.setWriteToTable(true);
+          insertRowStatement.setDatabaseName(dataBaseName);
+        }
+      }
+    }
     statement.setDatabaseName(dataBaseName);
     return statement;
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/request/PipeTransferTabletInsertNodeReqV2.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/request/PipeTransferTabletInsertNodeReqV2.java
index 8a130cf4da6..25c2fcaa844 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/request/PipeTransferTabletInsertNodeReqV2.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/request/PipeTransferTabletInsertNodeReqV2.java
@@ -28,6 +28,8 @@ import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNod
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsNode;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertTabletNode;
 import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
 
 import org.apache.tsfile.utils.PublicBAOS;
 import org.apache.tsfile.utils.ReadWriteIOUtils;
@@ -35,6 +37,7 @@ import org.apache.tsfile.utils.ReadWriteIOUtils;
 import java.io.DataOutputStream;
 import java.io.IOException;
 import java.nio.ByteBuffer;
+import java.util.List;
 import java.util.Objects;
 
 public class PipeTransferTabletInsertNodeReqV2 extends 
PipeTransferTabletInsertNodeReq {
@@ -71,6 +74,16 @@ public class PipeTransferTabletInsertNodeReqV2 extends 
PipeTransferTabletInsertN
 
     // Table model
     statement.setWriteToTable(true);
+    if (statement instanceof InsertRowsStatement) {
+      List<InsertRowStatement> rowStatements =
+          ((InsertRowsStatement) statement).getInsertRowStatementList();
+      if (rowStatements != null && !rowStatements.isEmpty()) {
+        for (InsertRowStatement insertRowStatement : rowStatements) {
+          insertRowStatement.setWriteToTable(true);
+          insertRowStatement.setDatabaseName(dataBaseName);
+        }
+      }
+    }
     statement.setDatabaseName(dataBaseName);
     return statement;
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/session/SessionManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/session/SessionManager.java
index e255b76621a..5681affe8cf 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/session/SessionManager.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/session/SessionManager.java
@@ -387,7 +387,7 @@ public class SessionManager implements SessionManagerMBean {
     return new SessionInfo(
         session.getId(),
         session.getUsername(),
-        session.getZoneId(),
+        ZoneId.systemDefault(),
         session.getClientVersion(),
         session.getDatabaseName(),
         IClientSession.SqlDialect.TABLE);
@@ -397,7 +397,7 @@ public class SessionManager implements SessionManagerMBean {
     return new SessionInfo(
         session.getId(),
         session.getUsername(),
-        session.getZoneId(),
+        ZoneId.systemDefault(),
         session.getClientVersion(),
         databaseName,
         IClientSession.SqlDialect.TABLE);

Reply via email to