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();
}
}