This is an automated email from the ASF dual-hosted git repository.
haonan 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 72799f3 [IOTDB-1503] 1 node crash causes whole cluster cannot work
(#3579)
72799f3 is described below
commit 72799f33215d5f6291544c354caf58c23a3dd04a
Author: Jamber <[email protected]>
AuthorDate: Fri Jul 23 10:21:58 2021 +0800
[IOTDB-1503] 1 node crash causes whole cluster cannot work (#3579)
Co-authored-by: haiyi.zb <[email protected]>
---
.../apache/iotdb/cluster/metadata/CMManager.java | 57 ++++++++++++++++------
.../apache/iotdb/cluster/metadata/MetaPuller.java | 27 +++++-----
.../iotdb/cluster/query/ClusterPlanExecutor.java | 40 ++++++++++++---
.../cluster/query/aggregate/ClusterAggregator.java | 9 +++-
.../cluster/query/fill/ClusterPreviousFill.java | 10 +++-
.../query/groupby/RemoteGroupByExecutor.java | 27 +++++++---
.../query/last/ClusterLastQueryExecutor.java | 26 ++++++----
.../cluster/query/reader/ClusterReaderFactory.java | 9 +++-
.../iotdb/cluster/query/reader/DataSourceInfo.java | 28 ++++++-----
.../reader/RemoteSeriesReaderByTimestamp.java | 3 ++
.../query/reader/RemoteSimpleSeriesReader.java | 3 ++
.../query/reader/mult/MultDataSourceInfo.java | 8 ++-
.../query/reader/mult/RemoteMultSeriesReader.java | 10 +++-
.../apache/iotdb/cluster/server/ClientServer.java | 8 ++-
.../cluster/server/heartbeat/HeartbeatThread.java | 6 +++
15 files changed, 199 insertions(+), 72 deletions(-)
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 f01870c..d21659a 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
@@ -913,8 +913,14 @@ 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) {
+ // the connection may be broken, close it to avoid it being
reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
if (result != null) {
@@ -1163,8 +1169,13 @@ 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) {
+ // the connection may be broken, close it to avoid it being reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
@@ -1291,8 +1302,13 @@ 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) {
+ // the connection may be broken, close it to avoid it being reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
return paths;
@@ -1752,10 +1768,16 @@ public class CMManager extends MManager {
metaGroupMember
.getClientProvider()
.getSyncDataClient(node,
RaftServer.getReadOperationTimeoutMS())) {
- plan.serialize(dataOutputStream);
- resultBinary =
- syncDataClient.getAllMeasurementSchema(
- group.getHeader(),
ByteBuffer.wrap(byteArrayOutputStream.toByteArray()));
+ try {
+ plan.serialize(dataOutputStream);
+ resultBinary =
+ syncDataClient.getAllMeasurementSchema(
+ group.getHeader(),
ByteBuffer.wrap(byteArrayOutputStream.toByteArray()));
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
return resultBinary;
@@ -1777,11 +1799,16 @@ public class CMManager extends MManager {
metaGroupMember
.getClientProvider()
.getSyncDataClient(node,
RaftServer.getReadOperationTimeoutMS())) {
-
- plan.serialize(dataOutputStream);
- resultBinary =
- syncDataClient.getDevices(
- group.getHeader(),
ByteBuffer.wrap(byteArrayOutputStream.toByteArray()));
+ try {
+ plan.serialize(dataOutputStream);
+ resultBinary =
+ syncDataClient.getDevices(
+ group.getHeader(),
ByteBuffer.wrap(byteArrayOutputStream.toByteArray()));
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ 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 ba85bd1..b212487 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
@@ -236,17 +236,22 @@ public class MetaPuller {
metaGroupMember
.getClientProvider()
.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(
- buffer.get() == 0
- ? MeasurementSchema.partialDeserializeFrom(buffer)
- : VectorMeasurementSchema.partialDeserializeFrom(buffer));
+ try {
+ // 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(
+ buffer.get() == 0
+ ? MeasurementSchema.partialDeserializeFrom(buffer)
+ : VectorMeasurementSchema.partialDeserializeFrom(buffer));
+ }
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ 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 bebe5af..849b802 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
@@ -258,8 +258,14 @@ public class ClusterPlanExecutor extends PlanExecutor {
metaGroupMember
.getClientProvider()
.getSyncDataClient(node,
RaftServer.getReadOperationTimeoutMS())) {
- syncDataClient.setTimeout(RaftServer.getReadOperationTimeoutMS());
- count = syncDataClient.getPathCount(partitionGroup.getHeader(),
pathsToQuery, level);
+ try {
+
syncDataClient.setTimeout(RaftServer.getReadOperationTimeoutMS());
+ count = syncDataClient.getPathCount(partitionGroup.getHeader(),
pathsToQuery, level);
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being
reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
logger.debug(
@@ -362,8 +368,14 @@ 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) {
+ // the connection may be broken, close it to avoid it being
reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
if (paths != null) {
@@ -448,8 +460,14 @@ 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) {
+ // the connection may be broken, close it to avoid it being
reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
if (nextChildrenNodes != null) {
@@ -556,8 +574,14 @@ 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) {
+ // the connection may be broken, close it to avoid it being
reused
+ 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 80aa81e..a27b8f3 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
@@ -269,8 +269,13 @@ public class ClusterAggregator {
metaGroupMember
.getClientProvider()
.getSyncDataClient(node,
RaftServer.getReadOperationTimeoutMS())) {
-
- resultBuffers = syncDataClient.getAggrResult(request);
+ try {
+ resultBuffers = syncDataClient.getAggrResult(request);
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ 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 e0f6053..548af4b 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
@@ -41,6 +41,7 @@ import org.apache.iotdb.db.utils.TimeValuePairUtils.Intervals;
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;
@@ -245,8 +246,13 @@ public class ClusterPreviousFill extends PreviousFill {
metaGroupMember
.getClientProvider()
.getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) {
-
- byteBuffer = syncDataClient.previousFill(request);
+ try {
+ byteBuffer = syncDataClient.previousFill(request);
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
} catch (Exception e) {
logger.error(
"{}: Cannot perform previous fill of {} to {}",
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 b2d52ac..99dfec6 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
@@ -27,7 +27,6 @@ import org.apache.iotdb.cluster.rpc.thrift.Node;
import org.apache.iotdb.cluster.rpc.thrift.RaftNode;
import org.apache.iotdb.cluster.server.RaftServer;
import org.apache.iotdb.cluster.server.member.MetaGroupMember;
-import org.apache.iotdb.cluster.utils.ClientUtils;
import org.apache.iotdb.db.query.aggregation.AggregateResult;
import org.apache.iotdb.db.query.dataset.groupby.GroupByExecutor;
import org.apache.iotdb.db.utils.SerializeUtils;
@@ -90,8 +89,14 @@ public class RemoteGroupByExecutor implements
GroupByExecutor {
metaGroupMember
.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();
+ throw e;
+ }
}
}
} catch (TException e) {
@@ -130,13 +135,19 @@ public class RemoteGroupByExecutor implements
GroupByExecutor {
SyncClientAdaptor.peekNextNotNullValue(
client, header, executorId, nextStartTime, nextEndTime);
} else {
- SyncDataClient syncDataClient =
+ try (SyncDataClient syncDataClient =
metaGroupMember
.getClientProvider()
- .getSyncDataClient(source,
RaftServer.getReadOperationTimeoutMS());
- aggrBuffer =
- syncDataClient.peekNextNotNullValue(header, executorId,
nextStartTime, nextEndTime);
- ClientUtils.putBackSyncClient(syncDataClient);
+ .getSyncDataClient(source,
RaftServer.getReadOperationTimeoutMS())) {
+ 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();
+ throw e;
+ }
+ }
}
} catch (TException e) {
throw new IOException(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 5ef7327..b42df93 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
@@ -259,20 +259,28 @@ public class ClusterLastQueryExecutor extends
LastQueryExecutor {
}
private ByteBuffer lastSync(Node node, QueryContext context) throws
TException {
+ ByteBuffer res;
try (SyncDataClient syncDataClient =
metaGroupMember
.getClientProvider()
.getSyncDataClient(node,
RaftServer.getReadOperationTimeoutMS())) {
-
- return syncDataClient.last(
- new LastQueryRequest(
- PartialPath.toStringList(seriesPaths),
- dataTypeOrdinals,
- context.getQueryId(),
- queryPlan.getDeviceToMeasurements(),
- group.getHeader(),
- syncDataClient.getNode()));
+ try {
+ res =
+ syncDataClient.last(
+ new LastQueryRequest(
+ PartialPath.toStringList(seriesPaths),
+ dataTypeOrdinals,
+ context.getQueryId(),
+ queryPlan.getDeviceToMeasurements(),
+ group.getHeader(),
+ syncDataClient.getNode()));
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
+ return res;
}
}
}
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 f379da5..f4ed148 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
@@ -952,8 +952,13 @@ public class ClusterReaderFactory {
metaGroupMember
.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 b9c4f00..685f0fe 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
@@ -162,19 +162,25 @@ public class DataSourceInfo {
.getClientProvider()
.getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) {
- if (byTimestamp) {
- newReaderId = client.querySingleSeriesByTimestamp(request);
- } else {
- Filter newFilter;
- // add timestamp to as a timeFilter to skip the data which has been
read
- if (request.isSetTimeFilterBytes()) {
- Filter timeFilter =
FilterFactory.deserialize(request.timeFilterBytes);
- newFilter = new AndFilter(timeFilter, TimeFilter.gt(timestamp));
+ try {
+ if (byTimestamp) {
+ newReaderId = client.querySingleSeriesByTimestamp(request);
} else {
- newFilter = TimeFilter.gt(timestamp);
+ Filter newFilter;
+ // add timestamp to as a timeFilter to skip the data which has been
read
+ if (request.isSetTimeFilterBytes()) {
+ Filter timeFilter =
FilterFactory.deserialize(request.timeFilterBytes);
+ newFilter = new AndFilter(timeFilter, TimeFilter.gt(timestamp));
+ } else {
+ newFilter = TimeFilter.gt(timestamp);
+ }
+
request.setTimeFilterBytes(SerializeUtils.serializeFilter(newFilter));
+ newReaderId = client.querySingleSeries(request);
}
- request.setTimeFilterBytes(SerializeUtils.serializeFilter(newFilter));
- newReaderId = client.querySingleSeries(request);
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ client.getInputProtocol().getTransport().close();
+ throw e;
}
return newReaderId;
}
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..e8b8d0e 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,9 @@ public class RemoteSeriesReaderByTimestamp implements
IReaderByTimestamp {
return curSyncClient.fetchSingleSeriesByTimestamps(
sourceInfo.getHeader(), sourceInfo.getReaderId(), timestampList);
} catch (TException e) {
+ if (curSyncClient != null) {
+ 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..b85b53b 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,9 @@ public class RemoteSimpleSeriesReader implements
IPointReader {
curSyncClient =
sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS());
return curSyncClient.fetchSingleSeries(sourceInfo.getHeader(),
sourceInfo.getReaderId());
} catch (TException e) {
+ if (curSyncClient != null) {
+ 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 14ca954..73dee46 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
@@ -190,7 +190,13 @@ 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();
+ throw e;
+ }
return newReaderId;
}
}
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 5e6bfba..a9bbe22 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
@@ -185,8 +185,14 @@ public class RemoteMultSeriesReader extends
AbstractMultPointReader {
try (SyncDataClient curSyncClient =
sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS()); )
{
-
- return curSyncClient.fetchMultSeries(sourceInfo.getHeader(),
sourceInfo.getReaderId(), paths);
+ try {
+ return curSyncClient.fetchMultSeries(
+ sourceInfo.getHeader(), sourceInfo.getReaderId(), paths);
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ curSyncClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
} catch (TException e) {
logger.error("Failed to fetch result sync, connect to {}", sourceInfo,
e);
return null;
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
index 1200c2d..d107635 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
@@ -319,7 +319,13 @@ public class ClientServer extends TSServiceImpl {
try (SyncDataClient syncDataClient =
coordinator.getSyncDataClient(
queriedNode, RaftServer.getReadOperationTimeoutMS())) {
- syncDataClient.endQuery(header, coordinator.getThisNode(),
queryId);
+ try {
+ syncDataClient.endQuery(header, coordinator.getThisNode(),
queryId);
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being
reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
} catch (IOException | TException e) {
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
index c9919d3..98b0fa7 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
@@ -34,6 +34,7 @@ import
org.apache.iotdb.cluster.server.handlers.caller.HeartbeatHandler;
import org.apache.iotdb.cluster.server.member.RaftMember;
import org.apache.iotdb.cluster.utils.ClientUtils;
+import org.apache.thrift.TException;
import org.apache.thrift.transport.TTransportException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -409,6 +410,11 @@ public class HeartbeatThread implements Runnable {
try {
long result = client.startElection(request);
handler.onComplete(result);
+ } catch (TException e) {
+ client.getInputProtocol().getTransport().close();
+ logger.warn(
+ "{}: Cannot request a vote from {} due to network",
memberName, node, e);
+ handler.onError(e);
} catch (Exception e) {
handler.onError(e);
} finally {