This is an automated email from the ASF dual-hosted git repository. tanxinyu pushed a commit to branch fix_thrift_out_of_sequence in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 1252a68e872e8e2b8e3ce4a8bea0e78db9785870 Author: LebronAl <[email protected]> AuthorDate: Thu Jul 15 19:03:39 2021 +0800 fix --- .../iotdb/cluster/client/DataClientProvider.java | 9 ++- .../iotdb/cluster/coordinator/Coordinator.java | 17 ++---- .../apache/iotdb/cluster/metadata/CMManager.java | 65 +++++++++++++++------- .../apache/iotdb/cluster/metadata/MetaPuller.java | 17 ++++-- .../iotdb/cluster/query/ClusterPlanExecutor.java | 35 +++++++++--- .../cluster/query/aggregate/ClusterAggregator.java | 8 ++- .../cluster/query/fill/ClusterPreviousFill.java | 25 ++++++--- .../query/groupby/RemoteGroupByExecutor.java | 21 +++++-- .../query/last/ClusterLastQueryExecutor.java | 25 ++++++--- .../cluster/query/reader/ClusterReaderFactory.java | 8 ++- .../iotdb/cluster/query/reader/DataSourceInfo.java | 15 +++-- .../reader/RemoteSeriesReaderByTimestamp.java | 1 + .../query/reader/RemoteSimpleSeriesReader.java | 1 + .../query/reader/mult/MultDataSourceInfo.java | 14 +++-- .../query/reader/mult/RemoteMultSeriesReader.java | 16 +++--- .../cluster/server/member/MetaGroupMember.java | 11 ++-- 16 files changed, 193 insertions(+), 95 deletions(-) diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/client/DataClientProvider.java b/cluster/src/main/java/org/apache/iotdb/cluster/client/DataClientProvider.java index 106705f..0950958 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/client/DataClientProvider.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/client/DataClientProvider.java @@ -28,7 +28,6 @@ import org.apache.iotdb.cluster.config.ClusterDescriptor; import org.apache.iotdb.cluster.rpc.thrift.Node; import org.apache.iotdb.cluster.rpc.thrift.RaftService.Client; -import org.apache.thrift.TException; import org.apache.thrift.protocol.TProtocolFactory; import java.io.IOException; @@ -102,10 +101,10 @@ public class DataClientProvider { * @param node the node to be connected * @param timeout timeout threshold of connection */ - public SyncDataClient getSyncDataClient(Node node, int timeout) throws TException { + public SyncDataClient getSyncDataClient(Node node, int timeout) throws IOException { SyncDataClient client = (SyncDataClient) getDataSyncClientPool().getClient(node); if (client == null) { - throw new TException(GET_CLIENT_FAILED_MSG + node); + throw new IOException(GET_CLIENT_FAILED_MSG + node); } client.setTimeout(timeout); return client; @@ -121,10 +120,10 @@ public class DataClientProvider { * @param node the node to be connected * @param timeout timeout threshold of connection */ - public SyncDataClient getSyncDataClientForRefresh(Node node, int timeout) throws TException { + public SyncDataClient getSyncDataClientForRefresh(Node node, int timeout) throws IOException { SyncDataClient client = (SyncDataClient) getDataSyncClientPool().getClientForRefresh(node); if (client == null) { - throw new TException(GET_CLIENT_FAILED_MSG + node); + throw new IOException(GET_CLIENT_FAILED_MSG + node); } client.setTimeout(timeout); return client; diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/coordinator/Coordinator.java b/cluster/src/main/java/org/apache/iotdb/cluster/coordinator/Coordinator.java index 11f99e8..db0fce3 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/coordinator/Coordinator.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/coordinator/Coordinator.java @@ -58,7 +58,6 @@ import org.apache.iotdb.rpc.TSStatusCode; import org.apache.iotdb.service.rpc.thrift.EndPoint; import org.apache.iotdb.service.rpc.thrift.TSStatus; -import org.apache.thrift.TException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -703,15 +702,11 @@ public class Coordinator { private TSStatus forwardDataPlanSync(PhysicalPlan plan, Node receiver, Node header) throws IOException { - RaftService.Client client = null; - try { - client = - metaGroupMember - .getClientProvider() - .getSyncDataClient(receiver, RaftServer.getWriteOperationTimeoutMS()); - } catch (TException e) { - throw new IOException(e); - } + RaftService.Client client = + metaGroupMember + .getClientProvider() + .getSyncDataClient(receiver, RaftServer.getWriteOperationTimeoutMS()); + return this.metaGroupMember.forwardPlanSync(plan, receiver, header, client); } @@ -735,7 +730,7 @@ public class Coordinator { * @param node the node to be connected * @param timeout timeout threshold of connection */ - public SyncDataClient getSyncDataClient(Node node, int timeout) throws TException { + public SyncDataClient getSyncDataClient(Node node, int timeout) throws IOException { return metaGroupMember.getClientProvider().getSyncDataClient(node, timeout); } } diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java index 319159a..9a3d69d 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java @@ -794,8 +794,13 @@ public class CMManager extends MManager { metaGroupMember .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - result = - syncDataClient.getUnregisteredTimeseries(partitionGroup.getHeader(), seriesList); + try { + result = + syncDataClient.getUnregisteredTimeseries(partitionGroup.getHeader(), seriesList); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } if (result != null) { @@ -978,13 +983,17 @@ public class CMManager extends MManager { metaGroupMember .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - - PullSchemaResp pullSchemaResp = syncDataClient.pullTimeSeriesSchema(request); - ByteBuffer buffer = pullSchemaResp.schemaBytes; - int size = buffer.getInt(); - schemas = new ArrayList<>(size); - for (int i = 0; i < size; i++) { - schemas.add(TimeseriesSchema.deserializeFrom(buffer)); + try { + PullSchemaResp pullSchemaResp = syncDataClient.pullTimeSeriesSchema(request); + ByteBuffer buffer = pullSchemaResp.schemaBytes; + int size = buffer.getInt(); + schemas = new ArrayList<>(size); + for (int i = 0; i < size; i++) { + schemas.add(TimeseriesSchema.deserializeFrom(buffer)); + } + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; } } } @@ -1212,8 +1221,12 @@ public class CMManager extends MManager { metaGroupMember .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - - result = syncDataClient.getAllPaths(header, pathsToQuery, withAlias); + try { + result = syncDataClient.getAllPaths(header, pathsToQuery, withAlias); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } @@ -1338,8 +1351,12 @@ public class CMManager extends MManager { metaGroupMember .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - - paths = syncDataClient.getAllDevices(header, pathsToQuery); + try { + paths = syncDataClient.getAllDevices(header, pathsToQuery); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } return paths; @@ -1792,9 +1809,14 @@ public class CMManager extends MManager { .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { plan.serialize(dataOutputStream); - resultBinary = - syncDataClient.getAllMeasurementSchema( - group.getHeader(), ByteBuffer.wrap(byteArrayOutputStream.toByteArray())); + try { + resultBinary = + syncDataClient.getAllMeasurementSchema( + group.getHeader(), ByteBuffer.wrap(byteArrayOutputStream.toByteArray())); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } return resultBinary; @@ -1818,9 +1840,14 @@ public class CMManager extends MManager { .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { plan.serialize(dataOutputStream); - resultBinary = - syncDataClient.getDevices( - group.getHeader(), ByteBuffer.wrap(byteArrayOutputStream.toByteArray())); + try { + resultBinary = + syncDataClient.getDevices( + group.getHeader(), ByteBuffer.wrap(byteArrayOutputStream.toByteArray())); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } return resultBinary; diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/MetaPuller.java b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/MetaPuller.java index e524772..a772259 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/MetaPuller.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/MetaPuller.java @@ -231,12 +231,17 @@ public class MetaPuller { .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { // only need measurement name - PullSchemaResp pullSchemaResp = syncDataClient.pullMeasurementSchema(request); - ByteBuffer buffer = pullSchemaResp.schemaBytes; - int size = buffer.getInt(); - schemas = new ArrayList<>(size); - for (int i = 0; i < size; i++) { - schemas.add(MeasurementSchema.deserializeFrom(buffer)); + try { + PullSchemaResp pullSchemaResp = syncDataClient.pullMeasurementSchema(request); + ByteBuffer buffer = pullSchemaResp.schemaBytes; + int size = buffer.getInt(); + schemas = new ArrayList<>(size); + for (int i = 0; i < size; i++) { + schemas.add(MeasurementSchema.deserializeFrom(buffer)); + } + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; } } } diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterPlanExecutor.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterPlanExecutor.java index ca16cfb..980ba38 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterPlanExecutor.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterPlanExecutor.java @@ -259,8 +259,12 @@ public class ClusterPlanExecutor extends PlanExecutor { metaGroupMember .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - syncDataClient.setTimeout(RaftServer.getReadOperationTimeoutMS()); - count = syncDataClient.getPathCount(partitionGroup.getHeader(), pathsToQuery, level); + try { + count = syncDataClient.getPathCount(partitionGroup.getHeader(), pathsToQuery, level); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } logger.debug( @@ -363,8 +367,13 @@ public class ClusterPlanExecutor extends PlanExecutor { metaGroupMember .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - paths = - syncDataClient.getNodeList(group.getHeader(), schemaPattern.getFullPath(), level); + try { + paths = + syncDataClient.getNodeList(group.getHeader(), schemaPattern.getFullPath(), level); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } if (paths != null) { @@ -449,8 +458,13 @@ public class ClusterPlanExecutor extends PlanExecutor { metaGroupMember .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - nextChildrenNodes = - syncDataClient.getChildNodeInNextLevel(group.getHeader(), path.getFullPath()); + try { + nextChildrenNodes = + syncDataClient.getChildNodeInNextLevel(group.getHeader(), path.getFullPath()); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } if (nextChildrenNodes != null) { @@ -558,8 +572,13 @@ public class ClusterPlanExecutor extends PlanExecutor { metaGroupMember .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - nextChildren = - syncDataClient.getChildNodePathInNextLevel(group.getHeader(), path.getFullPath()); + try { + nextChildren = + syncDataClient.getChildNodePathInNextLevel(group.getHeader(), path.getFullPath()); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } if (nextChildren != null) { diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java index 6c33f50..e8ec8e2 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java @@ -274,8 +274,12 @@ public class ClusterAggregator { metaGroupMember .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - - resultBuffers = syncDataClient.getAggrResult(request); + try { + resultBuffers = syncDataClient.getAggrResult(request); + } catch (TException e) { + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } return resultBuffers; diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterPreviousFill.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterPreviousFill.java index 33274e3..9c46845 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterPreviousFill.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterPreviousFill.java @@ -32,6 +32,7 @@ import org.apache.iotdb.cluster.server.RaftServer; import org.apache.iotdb.cluster.server.handlers.caller.PreviousFillHandler; import org.apache.iotdb.cluster.server.member.DataGroupMember; import org.apache.iotdb.cluster.server.member.MetaGroupMember; +import org.apache.iotdb.cluster.utils.ClientUtils; import org.apache.iotdb.cluster.utils.PartitionUtils.Intervals; import org.apache.iotdb.db.exception.StorageEngineException; import org.apache.iotdb.db.exception.query.QueryProcessException; @@ -41,6 +42,7 @@ import org.apache.iotdb.db.query.executor.fill.PreviousFill; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; import org.apache.iotdb.tsfile.read.TimeValuePair; +import org.apache.thrift.TException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -240,19 +242,28 @@ public class ClusterPreviousFill extends PreviousFill { private ByteBuffer remoteSyncPreviousFill( Node node, PreviousFillRequest request, PreviousFillArguments arguments) { ByteBuffer byteBuffer = null; - try (SyncDataClient syncDataClient = - metaGroupMember - .getClientProvider() - .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - - byteBuffer = syncDataClient.previousFill(request); - } catch (Exception e) { + SyncDataClient client = null; + try { + client = + metaGroupMember + .getClientProvider() + .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS()); + byteBuffer = client.previousFill(request); + } catch (TException | IOException e) { logger.error( "{}: Cannot perform previous fill of {} to {}", metaGroupMember.getName(), arguments.getPath(), node, e); + if (e instanceof TException) { + client.getInputProtocol().getTransport().close(); + } + + } finally { + if (client != null) { + ClientUtils.putBackSyncClient(client); + } } return byteBuffer; } diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/groupby/RemoteGroupByExecutor.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/groupby/RemoteGroupByExecutor.java index 02df747..30c115c 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/groupby/RemoteGroupByExecutor.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/groupby/RemoteGroupByExecutor.java @@ -89,8 +89,13 @@ public class RemoteGroupByExecutor implements GroupByExecutor { .getClientProvider() .getSyncDataClient(source, RaftServer.getReadOperationTimeoutMS())) { - aggrBuffers = - syncDataClient.getGroupByResult(header, executorId, curStartTime, curEndTime); + try { + aggrBuffers = + syncDataClient.getGroupByResult(header, executorId, curStartTime, curEndTime); + } catch (TException e) { + // the connection may be broken, close it to avoid it being reused + syncDataClient.getInputProtocol().getTransport().close(); + } } } } catch (TException e) { @@ -118,7 +123,7 @@ public class RemoteGroupByExecutor implements GroupByExecutor { @Override public Pair<Long, Object> peekNextNotNullValue(long nextStartTime, long nextEndTime) throws IOException { - ByteBuffer aggrBuffer; + ByteBuffer aggrBuffer = null; try { if (ClusterDescriptor.getInstance().getConfig().isUseAsyncServer()) { AsyncDataClient client = @@ -133,9 +138,13 @@ public class RemoteGroupByExecutor implements GroupByExecutor { metaGroupMember .getClientProvider() .getSyncDataClient(source, RaftServer.getReadOperationTimeoutMS())) { - - aggrBuffer = - syncDataClient.peekNextNotNullValue(header, executorId, nextStartTime, nextEndTime); + try { + aggrBuffer = + syncDataClient.peekNextNotNullValue(header, executorId, nextStartTime, nextEndTime); + } catch (TException e) { + // the connection may be broken, close it to avoid it being reused + syncDataClient.getInputProtocol().getTransport().close(); + } } } } catch (TException e) { diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java index d5ec324..bea7c1d 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java @@ -30,6 +30,7 @@ import org.apache.iotdb.cluster.rpc.thrift.Node; import org.apache.iotdb.cluster.server.RaftServer; import org.apache.iotdb.cluster.server.member.DataGroupMember; import org.apache.iotdb.cluster.server.member.MetaGroupMember; +import org.apache.iotdb.cluster.utils.ClientUtils; import org.apache.iotdb.cluster.utils.ClusterQueryUtils; import org.apache.iotdb.db.exception.StorageEngineException; import org.apache.iotdb.db.exception.query.QueryProcessException; @@ -258,19 +259,29 @@ public class ClusterLastQueryExecutor extends LastQueryExecutor { } private ByteBuffer lastSync(Node node, QueryContext context) throws TException { - try (SyncDataClient syncDataClient = - metaGroupMember - .getClientProvider() - .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - - return syncDataClient.last( + SyncDataClient client = null; + try { + client = + metaGroupMember + .getClientProvider() + .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS()); + return client.last( new LastQueryRequest( PartialPath.toStringList(seriesPaths), dataTypeOrdinals, context.getQueryId(), queryPlan.getDeviceToMeasurements(), group.getHeader(), - syncDataClient.getNode())); + client.getNode())); + } catch (IOException e) { + return null; + } catch (TException e) { + client.getInputProtocol().getTransport().close(); + throw e; + } finally { + if (client != null) { + ClientUtils.putBackSyncClient(client); + } } } } diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java index 64a2e3b..e3e0d53 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java @@ -896,7 +896,13 @@ public class ClusterReaderFactory { .getClientProvider() .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { - executorId = syncDataClient.getGroupByExecutor(request); + try { + executorId = syncDataClient.getGroupByExecutor(request); + } catch (TException e) { + // the connection may be broken, close it to avoid it being reused + syncDataClient.getInputProtocol().getTransport().close(); + throw e; + } } } return executorId; diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/DataSourceInfo.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/DataSourceInfo.java index 8889535..44767b7 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/DataSourceInfo.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/DataSourceInfo.java @@ -156,10 +156,15 @@ public class DataSourceInfo { throws TException { Long newReaderId; - try (SyncDataClient client = - this.metaGroupMember - .getClientProvider() - .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) { + try { + SyncDataClient client = + this.metaGroupMember + .getClientProvider() + .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS()) + }catch (IOException e){ + return null; + } + try () { if (byTimestamp) { newReaderId = client.querySingleSeriesByTimestamp(request); @@ -201,7 +206,7 @@ public class DataSourceInfo { : metaGroupMember.getClientProvider().getAsyncDataClient(this.curSource, timeout); } - SyncDataClient getCurSyncClient(int timeout) throws TException { + SyncDataClient getCurSyncClient(int timeout) throws IOException { return isNoClient ? null : metaGroupMember.getClientProvider().getSyncDataClient(this.curSource, timeout); diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSeriesReaderByTimestamp.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSeriesReaderByTimestamp.java index d077f02..2431951 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSeriesReaderByTimestamp.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSeriesReaderByTimestamp.java @@ -108,6 +108,7 @@ public class RemoteSeriesReaderByTimestamp implements IReaderByTimestamp { return curSyncClient.fetchSingleSeriesByTimestamps( sourceInfo.getHeader(), sourceInfo.getReaderId(), timestampList); } catch (TException e) { + curSyncClient.getInputProtocol().getTransport().close(); // try other node if (!sourceInfo.switchNode(true, timestamps[0])) { return null; diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSimpleSeriesReader.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSimpleSeriesReader.java index 2dcc1b7..70a0641 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSimpleSeriesReader.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSimpleSeriesReader.java @@ -149,6 +149,7 @@ public class RemoteSimpleSeriesReader implements IPointReader { curSyncClient = sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS()); return curSyncClient.fetchSingleSeries(sourceInfo.getHeader(), sourceInfo.getReaderId()); } catch (TException e) { + curSyncClient.getInputProtocol().getTransport().close(); // try other node if (!sourceInfo.switchNode(false, lastTimestamp)) { return null; diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/MultDataSourceInfo.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/MultDataSourceInfo.java index 27c9f59..8ffe4ce2 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/MultDataSourceInfo.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/MultDataSourceInfo.java @@ -172,9 +172,8 @@ public class MultDataSourceInfo { return result.get(); } - private Long applyForReaderIdSync(Node node, long timestamp) throws TException { - - Long newReaderId; + private Long applyForReaderIdSync(Node node, long timestamp) throws TException, IOException { + Long newReaderId = null; try (SyncDataClient client = this.metaGroupMember .getClientProvider() @@ -189,7 +188,12 @@ public class MultDataSourceInfo { newFilter = TimeFilter.gt(timestamp); } request.setTimeFilterBytes(SerializeUtils.serializeFilter(newFilter)); - newReaderId = client.queryMultSeries(request); + try { + newReaderId = client.queryMultSeries(request); + } catch (TException e) { + // the connection may be broken, close it to avoid it being reused + client.getInputProtocol().getTransport().close(); + } return newReaderId; } } @@ -212,7 +216,7 @@ public class MultDataSourceInfo { : metaGroupMember.getClientProvider().getAsyncDataClient(this.curSource, timeout); } - SyncDataClient getCurSyncClient(int timeout) throws TException { + SyncDataClient getCurSyncClient(int timeout) throws IOException { return isNoClient ? null : metaGroupMember.getClientProvider().getSyncDataClient(this.curSource, timeout); diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/RemoteMultSeriesReader.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/RemoteMultSeriesReader.java index bf20b35..31f7caa 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/RemoteMultSeriesReader.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/RemoteMultSeriesReader.java @@ -181,15 +181,17 @@ public class RemoteMultSeriesReader extends AbstractMultPointReader { } private Map<String, ByteBuffer> fetchResultSync(List<String> paths) throws IOException { - try (SyncDataClient curSyncClient = - sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS()); ) { - - return curSyncClient.fetchMultSeries(sourceInfo.getHeader(), sourceInfo.getReaderId(), paths); - } catch (TException e) { - logger.error("Failed to fetch result sync, connect to {}", sourceInfo, e); - return null; + sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS())) { + try { + return curSyncClient.fetchMultSeries( + sourceInfo.getHeader(), sourceInfo.getReaderId(), paths); + } catch (TException e) { + curSyncClient.getInputProtocol().getTransport().close(); + logger.error("Failed to fetch result sync, connect to {}", sourceInfo, e); + } } + return null; } /** select path, which could batch-fetch result */ diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java index cede9a3..e5f0a0f 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java @@ -498,23 +498,22 @@ public class MetaGroupMember extends RaftMember { } private void refreshClientOnceSync(Node receiver) { - RaftService.Client client; + RaftService.Client client = null; try { client = getClientProvider() .getSyncDataClientForRefresh(receiver, RaftServer.getWriteOperationTimeoutMS()); - } catch (TException e) { - return; - } - try { RefreshReuqest req = new RefreshReuqest(); client.refreshConnection(req); + } catch (IOException ignored) { } catch (TException e) { logger.warn("encounter refreshing client timeout, throw broken connection", e); // the connection may be broken, close it to avoid it being reused client.getInputProtocol().getTransport().close(); } finally { - ClientUtils.putBackSyncClient(client); + if (client != null) { + ClientUtils.putBackSyncClient(client); + } } }
