This is an automated email from the ASF dual-hosted git repository. rong pushed a commit to branch pipe-cache-leader in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 7463e365461d570e9bef31bb6f50e0ffbf0224ff Author: Steve Yurong Su <[email protected]> AuthorDate: Fri Dec 29 00:55:18 2023 +0800 refactor iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java --- .../async/IoTDBThriftAsyncClientManager.java | 226 +++++++++++++++++++++ .../thrift/async/IoTDBThriftAsyncConnector.java | 209 ++++--------------- .../thrift/sync/IoTDBThriftSyncConnector.java | 2 +- .../async/AsyncPipeDataTransferServiceClient.java | 8 + 4 files changed, 280 insertions(+), 165 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncClientManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncClientManager.java new file mode 100644 index 00000000000..06a5e336b76 --- /dev/null +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncClientManager.java @@ -0,0 +1,226 @@ +/* + * 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.db.pipe.connector.protocol.thrift.async; + +import org.apache.iotdb.common.rpc.thrift.TEndPoint; +import org.apache.iotdb.commons.client.ClientPoolFactory; +import org.apache.iotdb.commons.client.IClientManager; +import org.apache.iotdb.commons.client.async.AsyncPipeDataTransferServiceClient; +import org.apache.iotdb.commons.conf.CommonDescriptor; +import org.apache.iotdb.commons.pipe.config.PipeConfig; +import org.apache.iotdb.commons.pipe.connector.client.IoTDBThriftSyncConnectorClient; +import org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferHandshakeReq; +import org.apache.iotdb.db.pipe.connector.protocol.thrift.LeaderCacheManager; +import org.apache.iotdb.pipe.api.exception.PipeConnectionException; +import org.apache.iotdb.pipe.api.exception.PipeException; +import org.apache.iotdb.rpc.TSStatusCode; +import org.apache.iotdb.service.rpc.thrift.TPipeTransferResp; +import org.apache.iotdb.tsfile.utils.Pair; + +import org.apache.thrift.async.AsyncMethodCallback; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.Closeable; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; + +public class IoTDBThriftAsyncClientManager implements Closeable { + + private static final Logger LOGGER = LoggerFactory.getLogger(IoTDBThriftAsyncClientManager.class); + + private static final PipeConfig PIPE_CONFIG = PipeConfig.getInstance(); + + private final boolean useLeaderCache; + + private final List<TEndPoint> endPoints; + + private final LeaderCacheManager leaderCacheManager = new LeaderCacheManager(); + + private static final AtomicReference< + IClientManager<TEndPoint, AsyncPipeDataTransferServiceClient>> + ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER = new AtomicReference<>(); + private final IClientManager<TEndPoint, AsyncPipeDataTransferServiceClient> + asyncPipeDataTransferClientManager; + + private long currentClientIndex = 0; + + public IoTDBThriftAsyncClientManager(List<TEndPoint> endPoints, boolean useLeaderCache) { + this.useLeaderCache = useLeaderCache; + this.endPoints = endPoints; + + if (ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER.get() == null) { + synchronized (IoTDBThriftAsyncConnector.class) { + if (ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER.get() == null) { + ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER.set( + new IClientManager.Factory<TEndPoint, AsyncPipeDataTransferServiceClient>() + .createClientManager( + new ClientPoolFactory.AsyncPipeDataTransferServiceClientPoolFactory())); + } + } + } + asyncPipeDataTransferClientManager = ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER.get(); + } + + public AsyncPipeDataTransferServiceClient borrowClient() throws Exception { + final int clientSize = endPoints.size(); + while (true) { + final TEndPoint targetNodeUrl = endPoints.get((int) (currentClientIndex++ % clientSize)); + final AsyncPipeDataTransferServiceClient client = + asyncPipeDataTransferClientManager.borrowClient(targetNodeUrl); + if (handshakeIfNecessary(targetNodeUrl, client)) { + return client; + } + } + } + + public AsyncPipeDataTransferServiceClient borrowClient(String deviceId) { + return null; + } + + /** + * Handshake with the target if necessary. + * + * @param client client to handshake + * @return true if the handshake is already finished, false if the handshake is not finished yet + * and finished in this method + * @throws Exception if an error occurs. + */ + private boolean handshakeIfNecessary( + TEndPoint targetNodeUrl, AsyncPipeDataTransferServiceClient client) throws Exception { + if (client.isHandshakeFinished()) { + return true; + } + + final AtomicBoolean isHandshakeFinished = new AtomicBoolean(false); + final AtomicReference<Exception> exception = new AtomicReference<>(); + + client.pipeTransfer( + PipeTransferHandshakeReq.toTPipeTransferReq( + CommonDescriptor.getInstance().getConfig().getTimestampPrecision()), + new AsyncMethodCallback<TPipeTransferResp>() { + @Override + public void onComplete(TPipeTransferResp response) { + if (response.getStatus().getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) { + LOGGER.warn( + "Handshake error with receiver {}:{}, code: {}, message: {}.", + targetNodeUrl.getIp(), + targetNodeUrl.getPort(), + response.getStatus().getCode(), + response.getStatus().getMessage()); + exception.set( + new PipeConnectionException( + String.format( + "Handshake error with receiver %s:%s, code: %d, message: %s.", + targetNodeUrl.getIp(), + targetNodeUrl.getPort(), + response.getStatus().getCode(), + response.getStatus().getMessage()))); + } else { + LOGGER.info( + "Handshake successfully with receiver {}:{}.", + targetNodeUrl.getIp(), + targetNodeUrl.getPort()); + client.markHandshakeFinished(); + } + + isHandshakeFinished.set(true); + } + + @Override + public void onError(Exception e) { + LOGGER.warn( + "Handshake error with receiver {}:{}.", + targetNodeUrl.getIp(), + targetNodeUrl.getPort(), + e); + exception.set(e); + + isHandshakeFinished.set(true); + } + }); + + try { + while (!isHandshakeFinished.get()) { + Thread.sleep(10); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new PipeException("Interrupted while waiting for handshake response.", e); + } + + if (exception.get() != null) { + throw new PipeConnectionException("Failed to handshake.", exception.get()); + } + + return false; + } + + public void updateLeaderCache(String deviceId, TEndPoint endPoint) { + if (!useLeaderCache) { + return; + } + + try { + if (!endPoint2ClientAndStatus.containsKey(endPoint)) { + endPoints.add(endPoint); + endPoint2ClientAndStatus.put(endPoint, new Pair<>(null, false)); + reconstructClient(endPoint); + } + + leaderCacheManager.updateLeaderEndPoint(deviceId, endPoint); + } catch (Exception e) { + LOGGER.warn( + "Failed to update leader cache for device {} with endpoint {}:{}.", + deviceId, + endPoint.getIp(), + endPoint.getPort(), + e); + } + } + + @Override + public void close() { + for (final Map.Entry<TEndPoint, Pair<IoTDBThriftSyncConnectorClient, Boolean>> entry : + endPoint2ClientAndStatus.entrySet()) { + final TEndPoint endPoint = entry.getKey(); + final Pair<IoTDBThriftSyncConnectorClient, Boolean> clientAndStatus = entry.getValue(); + + if (clientAndStatus == null) { + continue; + } + + try { + if (clientAndStatus.getLeft() != null) { + clientAndStatus.getLeft().close(); + clientAndStatus.setLeft(null); + } + LOGGER.info("Client {}:{} closed.", endPoint.getIp(), endPoint.getPort()); + } catch (Exception e) { + LOGGER.warn( + "Failed to close client {}:{}, because: {}.", endPoint.getIp(), endPoint.getPort(), e); + } finally { + clientAndStatus.setRight(false); + } + } + } +} diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java index e608b39609c..6bea990bb31 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java @@ -19,14 +19,9 @@ package org.apache.iotdb.db.pipe.connector.protocol.thrift.async; -import org.apache.iotdb.common.rpc.thrift.TEndPoint; -import org.apache.iotdb.commons.client.ClientPoolFactory; -import org.apache.iotdb.commons.client.IClientManager; import org.apache.iotdb.commons.client.async.AsyncPipeDataTransferServiceClient; -import org.apache.iotdb.commons.conf.CommonDescriptor; import org.apache.iotdb.commons.pipe.plugin.builtin.connector.iotdb.IoTDBConnector; import org.apache.iotdb.db.pipe.connector.payload.evolvable.builder.IoTDBThriftAsyncPipeTransferBatchReqBuilder; -import org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferHandshakeReq; import org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferTabletBinaryReq; import org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferTabletInsertNodeReq; import org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferTabletRawReq; @@ -47,41 +42,37 @@ import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters; import org.apache.iotdb.pipe.api.event.Event; import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent; import org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent; -import org.apache.iotdb.pipe.api.exception.PipeConnectionException; -import org.apache.iotdb.pipe.api.exception.PipeException; -import org.apache.iotdb.rpc.TSStatusCode; import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq; -import org.apache.iotdb.service.rpc.thrift.TPipeTransferResp; -import org.apache.thrift.async.AsyncMethodCallback; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.FileNotFoundException; import java.io.IOException; +import java.util.Arrays; import java.util.Comparator; import java.util.HashMap; import java.util.concurrent.PriorityBlockingQueue; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicReference; import static org.apache.iotdb.commons.pipe.config.constant.PipeConnectorConstant.CONNECTOR_IOTDB_BATCH_MODE_ENABLE_KEY; +import static org.apache.iotdb.commons.pipe.config.constant.PipeConnectorConstant.CONNECTOR_LEADER_CACHE_ENABLE_DEFAULT_VALUE; +import static org.apache.iotdb.commons.pipe.config.constant.PipeConnectorConstant.CONNECTOR_LEADER_CACHE_ENABLE_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeConnectorConstant.SINK_IOTDB_BATCH_MODE_ENABLE_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeConnectorConstant.SINK_IOTDB_SSL_ENABLE_KEY; +import static org.apache.iotdb.commons.pipe.config.constant.PipeConnectorConstant.SINK_LEADER_CACHE_ENABLE_KEY; public class IoTDBThriftAsyncConnector extends IoTDBConnector { private static final Logger LOGGER = LoggerFactory.getLogger(IoTDBThriftAsyncConnector.class); - private static final String THRIFT_ERROR_FORMATTER = + private static final String THRIFT_ERROR_FORMATTER_WITHOUT_ENDPOINT = + "Failed to borrow client from client pool or exception occurred " + + "when sending to receiver."; + private static final String THRIFT_ERROR_FORMATTER_WITH_ENDPOINT = "Failed to borrow client from client pool or exception occurred " + "when sending to receiver %s:%s."; - private static final AtomicReference< - IClientManager<TEndPoint, AsyncPipeDataTransferServiceClient>> - ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER = new AtomicReference<>(); - private final IClientManager<TEndPoint, AsyncPipeDataTransferServiceClient> - asyncPipeDataTransferClientManager; + private IoTDBThriftAsyncClientManager clientManager; private final IoTDBThriftSyncConnector retryConnector = new IoTDBThriftSyncConnector(); private final PriorityBlockingQueue<Event> retryEventQueue = @@ -95,20 +86,6 @@ public class IoTDBThriftAsyncConnector extends IoTDBConnector { private IoTDBThriftAsyncPipeTransferBatchReqBuilder tabletBatchBuilder; - public IoTDBThriftAsyncConnector() { - if (ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER.get() == null) { - synchronized (IoTDBThriftAsyncConnector.class) { - if (ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER.get() == null) { - ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER.set( - new IClientManager.Factory<TEndPoint, AsyncPipeDataTransferServiceClient>() - .createClientManager( - new ClientPoolFactory.AsyncPipeDataTransferServiceClientPoolFactory())); - } - } - } - asyncPipeDataTransferClientManager = ASYNC_PIPE_DATA_TRANSFER_CLIENT_MANAGER_HOLDER.get(); - } - @Override public void validate(PipeParameterValidator validator) throws Exception { super.validate(validator); @@ -132,6 +109,13 @@ public class IoTDBThriftAsyncConnector extends IoTDBConnector { retryParameters.getAttribute().put(CONNECTOR_IOTDB_BATCH_MODE_ENABLE_KEY, "false"); retryConnector.customize(retryParameters, configuration); + clientManager = + new IoTDBThriftAsyncClientManager( + nodeUrls, + parameters.getBooleanOrDefault( + Arrays.asList(SINK_LEADER_CACHE_ENABLE_KEY, CONNECTOR_LEADER_CACHE_ENABLE_KEY), + CONNECTOR_LEADER_CACHE_ENABLE_DEFAULT_VALUE)); + if (isTabletBatchModeEnabled) { tabletBatchBuilder = new IoTDBThriftAsyncPipeTransferBatchReqBuilder(parameters); } @@ -173,14 +157,12 @@ public class IoTDBThriftAsyncConnector extends IoTDBConnector { return; } - final long commitId = ((EnrichedEvent) tabletInsertionEvent).getCommitId(); - if (isTabletBatchModeEnabled) { if (tabletBatchBuilder.onEvent(tabletInsertionEvent)) { final PipeTransferTabletBatchEventHandler pipeTransferTabletBatchEventHandler = new PipeTransferTabletBatchEventHandler(tabletBatchBuilder, this); - transfer(commitId, pipeTransferTabletBatchEventHandler); + transfer(pipeTransferTabletBatchEventHandler); tabletBatchBuilder.onSuccess(); } @@ -198,7 +180,7 @@ public class IoTDBThriftAsyncConnector extends IoTDBConnector { new PipeTransferTabletInsertNodeEventHandler( pipeInsertNodeTabletInsertionEvent, pipeTransferReq, this); - transfer(commitId, pipeTransferInsertNodeReqHandler); + transfer(pipeTransferInsertNodeReqHandler); } else { // tabletInsertionEvent instanceof PipeRawTabletInsertionEvent final PipeRawTabletInsertionEvent pipeRawTabletInsertionEvent = (PipeRawTabletInsertionEvent) tabletInsertionEvent; @@ -210,54 +192,40 @@ public class IoTDBThriftAsyncConnector extends IoTDBConnector { new PipeTransferTabletRawEventHandler( pipeRawTabletInsertionEvent, pipeTransferTabletRawReq, this); - transfer(commitId, pipeTransferTabletReqHandler); + transfer(pipeTransferTabletReqHandler); } } } - private void transfer( - long requestCommitId, - PipeTransferTabletBatchEventHandler pipeTransferTabletBatchEventHandler) { - final TEndPoint targetNodeUrl = nodeUrls.get((int) (requestCommitId % nodeUrls.size())); - + private void transfer(PipeTransferTabletBatchEventHandler pipeTransferTabletBatchEventHandler) { + AsyncPipeDataTransferServiceClient client = null; try { - final AsyncPipeDataTransferServiceClient client = borrowClient(targetNodeUrl); + client = clientManager.borrowClient(); pipeTransferTabletBatchEventHandler.transfer(client); } catch (Exception ex) { - LOGGER.warn( - String.format(THRIFT_ERROR_FORMATTER, targetNodeUrl.getIp(), targetNodeUrl.getPort()), - ex); + logOnClientException(client, ex); pipeTransferTabletBatchEventHandler.onError(ex); } } - private void transfer( - long requestCommitId, - PipeTransferTabletInsertNodeEventHandler pipeTransferInsertNodeReqHandler) { - final TEndPoint targetNodeUrl = nodeUrls.get((int) (requestCommitId % nodeUrls.size())); - + private void transfer(PipeTransferTabletInsertNodeEventHandler pipeTransferInsertNodeReqHandler) { + AsyncPipeDataTransferServiceClient client = null; try { - final AsyncPipeDataTransferServiceClient client = borrowClient(targetNodeUrl); + client = clientManager.borrowClient(); pipeTransferInsertNodeReqHandler.transfer(client); } catch (Exception ex) { - LOGGER.warn( - String.format(THRIFT_ERROR_FORMATTER, targetNodeUrl.getIp(), targetNodeUrl.getPort()), - ex); + logOnClientException(client, ex); pipeTransferInsertNodeReqHandler.onError(ex); } } - private void transfer( - long requestCommitId, PipeTransferTabletRawEventHandler pipeTransferTabletReqHandler) { - final TEndPoint targetNodeUrl = nodeUrls.get((int) (requestCommitId % nodeUrls.size())); - + private void transfer(PipeTransferTabletRawEventHandler pipeTransferTabletReqHandler) { + AsyncPipeDataTransferServiceClient client = null; try { - final AsyncPipeDataTransferServiceClient client = borrowClient(targetNodeUrl); + client = clientManager.borrowClient(); pipeTransferTabletReqHandler.transfer(client); } catch (Exception ex) { - LOGGER.warn( - String.format(THRIFT_ERROR_FORMATTER, targetNodeUrl.getIp(), targetNodeUrl.getPort()), - ex); + logOnClientException(client, ex); pipeTransferTabletReqHandler.onError(ex); } } @@ -303,21 +271,17 @@ public class IoTDBThriftAsyncConnector extends IoTDBConnector { final PipeTransferTsFileInsertionEventHandler pipeTransferTsFileInsertionEventHandler = new PipeTransferTsFileInsertionEventHandler(pipeTsFileInsertionEvent, this); - transfer(pipeTsFileInsertionEvent.getCommitId(), pipeTransferTsFileInsertionEventHandler); + transfer(pipeTransferTsFileInsertionEventHandler); } private void transfer( - long requestCommitId, PipeTransferTsFileInsertionEventHandler pipeTransferTsFileInsertionEventHandler) { - final TEndPoint targetNodeUrl = nodeUrls.get((int) (requestCommitId % nodeUrls.size())); - + AsyncPipeDataTransferServiceClient client = null; try { - final AsyncPipeDataTransferServiceClient client = borrowClient(targetNodeUrl); + client = clientManager.borrowClient(); pipeTransferTsFileInsertionEventHandler.transfer(client); } catch (Exception ex) { - LOGGER.warn( - String.format(THRIFT_ERROR_FORMATTER, targetNodeUrl.getIp(), targetNodeUrl.getPort()), - ex); + logOnClientException(client, ex); pipeTransferTsFileInsertionEventHandler.onError(ex); } } @@ -333,93 +297,15 @@ public class IoTDBThriftAsyncConnector extends IoTDBConnector { } } - private AsyncPipeDataTransferServiceClient borrowClient(TEndPoint targetNodeUrl) - throws Exception { - while (true) { - final AsyncPipeDataTransferServiceClient client = - asyncPipeDataTransferClientManager.borrowClient(targetNodeUrl); - if (handshakeIfNecessary(targetNodeUrl, client)) { - return client; - } - } - } - - /** - * Handshake with the target if necessary. - * - * @param client client to handshake - * @return true if the handshake is already finished, false if the handshake is not finished yet - * and finished in this method - * @throws Exception if an error occurs. - */ - private boolean handshakeIfNecessary( - TEndPoint targetNodeUrl, AsyncPipeDataTransferServiceClient client) throws Exception { - if (client.isHandshakeFinished()) { - return true; - } - - final AtomicBoolean isHandshakeFinished = new AtomicBoolean(false); - final AtomicReference<Exception> exception = new AtomicReference<>(); - - client.pipeTransfer( - PipeTransferHandshakeReq.toTPipeTransferReq( - CommonDescriptor.getInstance().getConfig().getTimestampPrecision()), - new AsyncMethodCallback<TPipeTransferResp>() { - @Override - public void onComplete(TPipeTransferResp response) { - if (response.getStatus().getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) { - LOGGER.warn( - "Handshake error with receiver {}:{}, code: {}, message: {}.", - targetNodeUrl.getIp(), - targetNodeUrl.getPort(), - response.getStatus().getCode(), - response.getStatus().getMessage()); - exception.set( - new PipeConnectionException( - String.format( - "Handshake error with receiver %s:%s, code: %d, message: %s.", - targetNodeUrl.getIp(), - targetNodeUrl.getPort(), - response.getStatus().getCode(), - response.getStatus().getMessage()))); - } else { - LOGGER.info( - "Handshake successfully with receiver {}:{}.", - targetNodeUrl.getIp(), - targetNodeUrl.getPort()); - client.markHandshakeFinished(); - } - - isHandshakeFinished.set(true); - } - - @Override - public void onError(Exception e) { - LOGGER.warn( - "Handshake error with receiver {}:{}.", - targetNodeUrl.getIp(), - targetNodeUrl.getPort(), - e); - exception.set(e); - - isHandshakeFinished.set(true); - } - }); + //////////////////////////// Exception handlers //////////////////////////// - try { - while (!isHandshakeFinished.get()) { - Thread.sleep(10); - } - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw new PipeException("Interrupted while waiting for handshake response.", e); - } - - if (exception.get() != null) { - throw new PipeConnectionException("Failed to handshake.", exception.get()); + private void logOnClientException(AsyncPipeDataTransferServiceClient client, Exception e) { + if (client == null) { + LOGGER.warn(THRIFT_ERROR_FORMATTER_WITHOUT_ENDPOINT, e); + } else { + LOGGER.warn( + String.format(THRIFT_ERROR_FORMATTER_WITH_ENDPOINT, client.getIp(), client.getPort()), e); } - - return false; } /** @@ -462,14 +348,7 @@ public class IoTDBThriftAsyncConnector extends IoTDBConnector { return; } - // requestCommitId can not be generated by commitIdGenerator because the commit id must - // be bind to a specific InsertTabletEvent or TsFileInsertionEvent, otherwise the commit - // process will stuck. - final long requestCommitId = tabletBatchBuilder.getLastCommitId(); - final PipeTransferTabletBatchEventHandler pipeTransferTabletBatchEventHandler = - new PipeTransferTabletBatchEventHandler(tabletBatchBuilder, this); - - transfer(requestCommitId, pipeTransferTabletBatchEventHandler); + transfer(new PipeTransferTabletBatchEventHandler(tabletBatchBuilder, this)); tabletBatchBuilder.onSuccess(); } @@ -483,6 +362,8 @@ public class IoTDBThriftAsyncConnector extends IoTDBConnector { retryEventQueue.offer(event); } + //////////////////////////// Operations for close //////////////////////////// + /** * When a pipe is dropped, the connector maybe reused and will not be closed. So we just discard * its queued events in the output pipe connector. diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java index c87d00ce7f0..0d4ef07d539 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java @@ -283,7 +283,7 @@ public class IoTDBThriftSyncConnector extends IoTDBConnector { } } - private void doTransfer() throws IOException { + private void doTransfer() { Pair<IoTDBThriftSyncConnectorClient, Boolean> clientAndStatus = clientManager.getClient(); final TPipeTransferResp resp; try { diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/async/AsyncPipeDataTransferServiceClient.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/async/AsyncPipeDataTransferServiceClient.java index 83f47956812..2b8110bc7bc 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/async/AsyncPipeDataTransferServiceClient.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/async/AsyncPipeDataTransferServiceClient.java @@ -146,6 +146,14 @@ public class AsyncPipeDataTransferServiceClient extends IClientRPCService.AsyncC LOGGER.info("Handshake finished for client {}", this); } + public String getIp() { + return endpoint.getIp(); + } + + public int getPort() { + return endpoint.getPort(); + } + @Override public String toString() { return String.format("AsyncPipeDataTransferServiceClient{%s}, id = {%d}", endpoint, id);
