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

rong pushed a commit to branch pipe-parallel-connector
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/pipe-parallel-connector by 
this push:
     new 57bdcc9e202 fix return client bug
57bdcc9e202 is described below

commit 57bdcc9e202bda529e353016272672bd9daad5f3
Author: Steve Yurong Su <[email protected]>
AuthorDate: Fri Jun 16 02:39:04 2023 +0800

    fix return client bug
---
 .../apache/iotdb/pipe/api/access/RowIterator.java  | 75 ----------------------
 .../async/AsyncPipeDataTransferServiceClient.java  |  2 +-
 .../pipe/agent/receiver/IoTDBThriftReceiver.java   |  4 +-
 .../db/pipe/agent/receiver/PipeReceiverAgent.java  |  6 +-
 ...ava => IoTDBThriftConnectorRequestVersion.java} |  4 +-
 .../pipe/connector/v1/IoTDBThriftReceiverV1.java   |  6 +-
 .../v1/request/PipeTransferFilePieceReq.java       |  4 +-
 .../v1/request/PipeTransferFileSealReq.java        |  4 +-
 .../v1/request/PipeTransferHandshakeReq.java       |  4 +-
 .../v1/request/PipeTransferInsertNodeReq.java      |  4 +-
 .../v1/request/PipeTransferTabletReq.java          |  4 +-
 .../PipeTransferTsFileInsertionEventHandler.java   |  2 +
 12 files changed, 23 insertions(+), 96 deletions(-)

diff --git 
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/access/RowIterator.java
 
b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/access/RowIterator.java
deleted file mode 100644
index ae9849eb3cc..00000000000
--- 
a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/access/RowIterator.java
+++ /dev/null
@@ -1,75 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing,
- * software distributed under the License is distributed on an
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
- * KIND, either express or implied.  See the License for the
- * specific language governing permissions and limitations
- * under the License.
- */
-
-package org.apache.iotdb.pipe.api.access;
-
-import org.apache.iotdb.pipe.api.exception.PipeParameterNotValidException;
-import org.apache.iotdb.pipe.api.type.Type;
-import org.apache.iotdb.tsfile.read.common.Path;
-
-import java.io.IOException;
-import java.util.List;
-
-public interface RowIterator {
-
-  /**
-   * Returns {@code true} if the iteration has more rows.
-   *
-   * @return {@code true} if the iteration has more rows
-   */
-  boolean hasNextRow();
-
-  /**
-   * Returns the next row in the iteration.
-   *
-   * <p>Note that the Row instance returned by this method each time is the 
same instance. In other
-   * words, calling {@code next()} will only change the member variables 
inside the Row instance,
-   * but will not generate a new Row instance.
-   *
-   * @return the next element in the iteration
-   * @throws IOException if any I/O errors occur
-   */
-  Row next() throws IOException;
-
-  /** Resets the iteration. */
-  void reset();
-
-  /**
-   * Returns the actual column index of the given column name.
-   *
-   * @param columnName the column name in Path form
-   * @throws PipeParameterNotValidException if the given column name is not 
existed
-   * @return the actual column index of the given column name
-   */
-  int getColumnIndex(Path columnName) throws PipeParameterNotValidException;
-
-  /**
-   * Returns the column names
-   *
-   * @return the column names
-   */
-  List<Path> getColumnNames();
-
-  /**
-   * Returns the column data types
-   *
-   * @return the column data types
-   */
-  List<Type> getColumnTypes();
-}
diff --git 
a/node-commons/src/main/java/org/apache/iotdb/commons/client/async/AsyncPipeDataTransferServiceClient.java
 
b/node-commons/src/main/java/org/apache/iotdb/commons/client/async/AsyncPipeDataTransferServiceClient.java
index fcb4b90321a..534349460f1 100644
--- 
a/node-commons/src/main/java/org/apache/iotdb/commons/client/async/AsyncPipeDataTransferServiceClient.java
+++ 
b/node-commons/src/main/java/org/apache/iotdb/commons/client/async/AsyncPipeDataTransferServiceClient.java
@@ -99,7 +99,7 @@ public class AsyncPipeDataTransferServiceClient extends 
IClientRPCService.AsyncC
    * return self, the method doesn't need to be called by the user and will be 
triggered after the
    * RPC is finished.
    */
-  private void returnSelf() {
+  public void returnSelf() {
     if (shouldReturnSelf.get()) {
       clientManager.returnClient(endpoint, this);
     }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/agent/receiver/IoTDBThriftReceiver.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/agent/receiver/IoTDBThriftReceiver.java
index 91c2b980da3..80b18288c3b 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/agent/receiver/IoTDBThriftReceiver.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/agent/receiver/IoTDBThriftReceiver.java
@@ -21,13 +21,13 @@ package org.apache.iotdb.db.pipe.agent.receiver;
 
 import org.apache.iotdb.db.mpp.plan.analyze.IPartitionFetcher;
 import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaFetcher;
-import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorVersion;
+import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorRequestVersion;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferResp;
 
 public interface IoTDBThriftReceiver {
 
-  IoTDBThriftConnectorVersion getVersion();
+  IoTDBThriftConnectorRequestVersion getVersion();
 
   TPipeTransferResp receive(
       TPipeTransferReq req, IPartitionFetcher partitionFetcher, ISchemaFetcher 
schemaFetcher);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/agent/receiver/PipeReceiverAgent.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/agent/receiver/PipeReceiverAgent.java
index 97489d7e433..663e169c0bf 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/agent/receiver/PipeReceiverAgent.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/agent/receiver/PipeReceiverAgent.java
@@ -21,7 +21,7 @@ package org.apache.iotdb.db.pipe.agent.receiver;
 
 import org.apache.iotdb.db.mpp.plan.analyze.IPartitionFetcher;
 import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaFetcher;
-import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorVersion;
+import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorRequestVersion;
 import org.apache.iotdb.db.pipe.connector.v1.IoTDBThriftReceiverV1;
 import org.apache.iotdb.rpc.RpcUtils;
 import org.apache.iotdb.rpc.TSStatusCode;
@@ -40,7 +40,7 @@ public class PipeReceiverAgent {
   public TPipeTransferResp receive(
       TPipeTransferReq req, IPartitionFetcher partitionFetcher, ISchemaFetcher 
schemaFetcher) {
     final byte reqVersion = req.getVersion();
-    if (reqVersion == IoTDBThriftConnectorVersion.VERSION_1.getVersion()) {
+    if (reqVersion == 
IoTDBThriftConnectorRequestVersion.VERSION_1.getVersion()) {
       return getReceiver(reqVersion).receive(req, partitionFetcher, 
schemaFetcher);
     } else {
       return new TPipeTransferResp(
@@ -71,7 +71,7 @@ public class PipeReceiverAgent {
   }
 
   private IoTDBThriftReceiver setAndGetReceiver(byte reqVersion) {
-    if (reqVersion == IoTDBThriftConnectorVersion.VERSION_1.getVersion()) {
+    if (reqVersion == 
IoTDBThriftConnectorRequestVersion.VERSION_1.getVersion()) {
       receiverThreadLocal.set(new IoTDBThriftReceiverV1());
     } else {
       throw new UnsupportedOperationException(
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/IoTDBThriftConnectorVersion.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/IoTDBThriftConnectorRequestVersion.java
similarity index 91%
rename from 
server/src/main/java/org/apache/iotdb/db/pipe/connector/IoTDBThriftConnectorVersion.java
rename to 
server/src/main/java/org/apache/iotdb/db/pipe/connector/IoTDBThriftConnectorRequestVersion.java
index 8b7328a0425..670dd9eb6e0 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/IoTDBThriftConnectorVersion.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/IoTDBThriftConnectorRequestVersion.java
@@ -19,14 +19,14 @@
 
 package org.apache.iotdb.db.pipe.connector;
 
-public enum IoTDBThriftConnectorVersion {
+public enum IoTDBThriftConnectorRequestVersion {
   VERSION_1((byte) 1),
   VERSION_2((byte) 1),
   ;
 
   private final byte version;
 
-  IoTDBThriftConnectorVersion(byte type) {
+  IoTDBThriftConnectorRequestVersion(byte type) {
     this.version = type;
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/IoTDBThriftReceiverV1.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/IoTDBThriftReceiverV1.java
index 9c4d689d20b..5ab2b19c5a7 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/IoTDBThriftReceiverV1.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/IoTDBThriftReceiverV1.java
@@ -30,7 +30,7 @@ import org.apache.iotdb.db.mpp.plan.statement.Statement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertTabletStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.LoadTsFileStatement;
 import org.apache.iotdb.db.pipe.agent.receiver.IoTDBThriftReceiver;
-import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorVersion;
+import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorRequestVersion;
 import org.apache.iotdb.db.pipe.connector.v1.reponse.PipeTransferFilePieceResp;
 import org.apache.iotdb.db.pipe.connector.v1.request.PipeTransferFilePieceReq;
 import org.apache.iotdb.db.pipe.connector.v1.request.PipeTransferFileSealReq;
@@ -297,7 +297,7 @@ public class IoTDBThriftReceiverV1 implements 
IoTDBThriftReceiver {
   }
 
   @Override
-  public IoTDBThriftConnectorVersion getVersion() {
-    return IoTDBThriftConnectorVersion.VERSION_1;
+  public IoTDBThriftConnectorRequestVersion getVersion() {
+    return IoTDBThriftConnectorRequestVersion.VERSION_1;
   }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferFilePieceReq.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferFilePieceReq.java
index a40e6195fdd..d26939ec54f 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferFilePieceReq.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferFilePieceReq.java
@@ -19,7 +19,7 @@
 
 package org.apache.iotdb.db.pipe.connector.v1.request;
 
-import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorVersion;
+import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorRequestVersion;
 import org.apache.iotdb.db.pipe.connector.v1.PipeRequestType;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 import org.apache.iotdb.tsfile.utils.Binary;
@@ -58,7 +58,7 @@ public class PipeTransferFilePieceReq extends 
TPipeTransferReq {
     filePieceReq.startWritingOffset = startWritingOffset;
     filePieceReq.filePiece = filePiece;
 
-    filePieceReq.version = IoTDBThriftConnectorVersion.VERSION_1.getVersion();
+    filePieceReq.version = 
IoTDBThriftConnectorRequestVersion.VERSION_1.getVersion();
     filePieceReq.type = PipeRequestType.TRANSFER_FILE_PIECE.getType();
     try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferFileSealReq.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferFileSealReq.java
index 8f6e7186988..1314e5e5ae0 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferFileSealReq.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferFileSealReq.java
@@ -19,7 +19,7 @@
 
 package org.apache.iotdb.db.pipe.connector.v1.request;
 
-import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorVersion;
+import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorRequestVersion;
 import org.apache.iotdb.db.pipe.connector.v1.PipeRequestType;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 import org.apache.iotdb.tsfile.utils.PublicBAOS;
@@ -51,7 +51,7 @@ public class PipeTransferFileSealReq extends TPipeTransferReq 
{
     fileSealReq.fileName = fileName;
     fileSealReq.fileLength = fileLength;
 
-    fileSealReq.version = IoTDBThriftConnectorVersion.VERSION_1.getVersion();
+    fileSealReq.version = 
IoTDBThriftConnectorRequestVersion.VERSION_1.getVersion();
     fileSealReq.type = PipeRequestType.TRANSFER_FILE_SEAL.getType();
     try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferHandshakeReq.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferHandshakeReq.java
index b8b9b9f38ad..aceb9a5877b 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferHandshakeReq.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferHandshakeReq.java
@@ -19,7 +19,7 @@
 
 package org.apache.iotdb.db.pipe.connector.v1.request;
 
-import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorVersion;
+import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorRequestVersion;
 import org.apache.iotdb.db.pipe.connector.v1.PipeRequestType;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 import org.apache.iotdb.tsfile.utils.PublicBAOS;
@@ -45,7 +45,7 @@ public class PipeTransferHandshakeReq extends 
TPipeTransferReq {
 
     handshakeReq.timestampPrecision = timestampPrecision;
 
-    handshakeReq.version = IoTDBThriftConnectorVersion.VERSION_1.getVersion();
+    handshakeReq.version = 
IoTDBThriftConnectorRequestVersion.VERSION_1.getVersion();
     handshakeReq.type = PipeRequestType.HANDSHAKE.getType();
     try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferInsertNodeReq.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferInsertNodeReq.java
index f852cf936d4..157a4200fd6 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferInsertNodeReq.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferInsertNodeReq.java
@@ -26,7 +26,7 @@ import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertTabletNode;
 import org.apache.iotdb.db.mpp.plan.statement.Statement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertRowStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertTabletStatement;
-import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorVersion;
+import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorRequestVersion;
 import org.apache.iotdb.db.pipe.connector.v1.PipeRequestType;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 
@@ -83,7 +83,7 @@ public class PipeTransferInsertNodeReq extends 
TPipeTransferReq {
 
     req.insertNode = insertNode;
 
-    req.version = IoTDBThriftConnectorVersion.VERSION_1.getVersion();
+    req.version = IoTDBThriftConnectorRequestVersion.VERSION_1.getVersion();
     req.type = PipeRequestType.TRANSFER_INSERT_NODE.getType();
     req.body = insertNode.serializeToByteBuffer();
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferTabletReq.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferTabletReq.java
index 7c5716be77b..ead686a6acf 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferTabletReq.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v1/request/PipeTransferTabletReq.java
@@ -23,7 +23,7 @@ import org.apache.iotdb.commons.exception.MetadataException;
 import org.apache.iotdb.commons.utils.PathUtils;
 import org.apache.iotdb.db.mpp.plan.parser.StatementGenerator;
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertTabletStatement;
-import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorVersion;
+import org.apache.iotdb.db.pipe.connector.IoTDBThriftConnectorRequestVersion;
 import org.apache.iotdb.db.pipe.connector.v1.PipeRequestType;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 import org.apache.iotdb.service.rpc.thrift.TSInsertTabletReq;
@@ -52,7 +52,7 @@ public class PipeTransferTabletReq extends TPipeTransferReq {
 
     tabletReq.tablet = tablet;
 
-    tabletReq.version = IoTDBThriftConnectorVersion.VERSION_1.getVersion();
+    tabletReq.version = 
IoTDBThriftConnectorRequestVersion.VERSION_1.getVersion();
     tabletReq.type = PipeRequestType.TRANSFER_TABLET.getType();
     tabletReq.body = tablet.serialize();
     return tabletReq;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v2/handler/PipeTransferTsFileInsertionEventHandler.java
 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v2/handler/PipeTransferTsFileInsertionEventHandler.java
index 3a017f6863b..b597cce2408 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/pipe/connector/v2/handler/PipeTransferTsFileInsertionEventHandler.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/pipe/connector/v2/handler/PipeTransferTsFileInsertionEventHandler.java
@@ -130,6 +130,7 @@ public class PipeTransferTsFileInsertionEventHandler
       } finally {
         if (client != null) {
           client.setShouldReturnSelf(true);
+          client.returnSelf();
         }
 
         connector.commit(requestCommitId, event);
@@ -172,6 +173,7 @@ public class PipeTransferTsFileInsertionEventHandler
     } finally {
       if (client != null) {
         client.setShouldReturnSelf(true);
+        client.returnSelf();
       }
     }
 

Reply via email to