This is an automated email from the ASF dual-hosted git repository.
qiaojialin pushed a commit to branch rel/0.12
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rel/0.12 by this push:
new 18f0144 [To rel/0.12] Cherry pick from autoai (#3637)
18f0144 is described below
commit 18f0144d099712b4323f9c836e878a87820a80e7
Author: Jialin Qiao <[email protected]>
AuthorDate: Tue Jul 27 06:54:44 2021 -0500
[To rel/0.12] Cherry pick from autoai (#3637)
---
.../iotdb/cluster/client/DataClientProvider.java | 9 +-
.../iotdb/cluster/client/sync/SyncClientPool.java | 2 +-
.../iotdb/cluster/coordinator/Coordinator.java | 17 +-
.../apache/iotdb/cluster/metadata/CMManager.java | 71 ++++++--
.../apache/iotdb/cluster/metadata/MetaPuller.java | 18 +-
.../iotdb/cluster/query/ClusterPlanExecutor.java | 39 +++-
.../cluster/query/aggregate/ClusterAggregator.java | 9 +-
.../cluster/query/fill/ClusterPreviousFill.java | 25 ++-
.../query/groupby/RemoteGroupByExecutor.java | 21 ++-
.../query/last/ClusterLastQueryExecutor.java | 26 ++-
.../cluster/query/reader/ClusterReaderFactory.java | 8 +-
.../iotdb/cluster/query/reader/DataSourceInfo.java | 36 ++--
.../reader/RemoteSeriesReaderByTimestamp.java | 2 +
.../query/reader/RemoteSimpleSeriesReader.java | 2 +
.../query/reader/mult/MultDataSourceInfo.java | 15 +-
.../query/reader/mult/RemoteMultSeriesReader.java | 17 +-
.../apache/iotdb/cluster/server/ClientServer.java | 8 +-
.../cluster/server/heartbeat/HeartbeatThread.java | 7 +
.../cluster/server/member/DataGroupMember.java | 8 +-
.../cluster/server/member/MetaGroupMember.java | 19 +-
.../iotdb/cluster/server/member/RaftMember.java | 9 +-
.../cluster/client/DataClientProviderTest.java | 5 +-
.../org/apache/iotdb/SessionConcurrentExample.java | 198 +++++++++++++++++++++
.../db/engine/storagegroup/TsFileProcessor.java | 10 ++
.../org/apache/iotdb/db/metadata/MManager.java | 10 +-
.../iotdb/db/rescon/PrimitiveArrayManager.java | 10 +-
.../org/apache/iotdb/db/service/TSServiceImpl.java | 8 +-
.../db/writelog/node/ExclusiveWriteLogNode.java | 81 +++------
.../writelog/recover/TsFileRecoverPerformer.java | 25 ++-
.../iotdb/db/integration/IoTDBRestartIT.java | 48 +++++
.../read/expression/util/ExpressionOptimizer.java | 34 ++--
31 files changed, 601 insertions(+), 196 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/client/sync/SyncClientPool.java
b/cluster/src/main/java/org/apache/iotdb/cluster/client/sync/SyncClientPool.java
index 67dbae5..f607fa3 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/client/sync/SyncClientPool.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/client/sync/SyncClientPool.java
@@ -67,7 +67,7 @@ public class SyncClientPool {
if (clientStack.isEmpty()) {
return null;
} else {
- return clientStack.pollLast();
+ return clientStack.poll();
}
}
}
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..05d61a4 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,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) {
@@ -978,13 +984,18 @@ 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) {
+ // the connection may be broken, close it to avoid it being reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
}
}
}
@@ -1212,8 +1223,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;
+ }
}
}
@@ -1338,8 +1354,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;
@@ -1792,9 +1813,15 @@ 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) {
+ // the connection may be broken, close it to avoid it being reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
return resultBinary;
@@ -1818,9 +1845,15 @@ 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) {
+ // 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 e524772..9991c5a 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,18 @@ 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) {
+ // 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 ca16cfb..34e9904 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,13 @@ 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) {
+ // the connection may be broken, close it to avoid it being
reused
+ syncDataClient.getInputProtocol().getTransport().close();
+ throw e;
+ }
}
}
logger.debug(
@@ -363,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) {
@@ -449,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) {
@@ -558,8 +575,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 6c33f50..dc52923 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,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 33274e3..9af7082 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 (IOException e) {
+ logger.warn("{}: Cannot connect to {} during previous fill",
metaGroupMember, node);
+ } catch (TException e) {
logger.error(
"{}: Cannot perform previous fill of {} to {}",
metaGroupMember.getName(),
arguments.getPath(),
node,
e);
+ // the connection may be broken, close it to avoid it being reused
+ 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..468ec57 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,14 @@ 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();
+ throw e;
+ }
}
}
} catch (TException e) {
@@ -133,9 +139,14 @@ 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();
+ throw e;
+ }
}
}
} 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..a30db6c 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,30 @@ 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) {
+ // the connection may be broken, close it to avoid it being reused
+ 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..4ba11e4 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
@@ -153,27 +153,31 @@ public class DataSourceInfo {
}
private Long applyForReaderIdSync(Node node, boolean byTimestamp, long
timestamp)
- throws TException {
-
- Long newReaderId;
+ throws TException, IOException {
+ long newReaderId;
try (SyncDataClient client =
this.metaGroupMember
.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;
}
@@ -201,7 +205,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..e266af8 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,8 @@ public class RemoteSeriesReaderByTimestamp implements
IReaderByTimestamp {
return curSyncClient.fetchSingleSeriesByTimestamps(
sourceInfo.getHeader(), sourceInfo.getReaderId(), timestampList);
} catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ 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..f53f2bc 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,8 @@ public class RemoteSimpleSeriesReader implements
IPointReader {
curSyncClient =
sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS());
return curSyncClient.fetchSingleSeries(sourceInfo.getHeader(),
sourceInfo.getReaderId());
} catch (TException e) {
+ // the connection may be broken, close it to avoid it being reused
+ 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..a4488aa 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;
try (SyncDataClient client =
this.metaGroupMember
.getClientProvider()
@@ -189,7 +188,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;
}
}
@@ -212,7 +217,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..36513d4 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,18 @@ 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) {
+ // the connection may be broken, close it to avoid it being reused
+ 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/ClientServer.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
index c627373..722aeaf 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
@@ -317,7 +317,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 0459fef..67acc5f 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;
@@ -222,6 +223,7 @@ public class HeartbeatThread implements Runnable {
} catch (TTransportException e) {
logger.warn(
"{}: Cannot send heart beat to node {} due to network",
memberName, node, e);
+ // the connection may be broken, close it to avoid it being
reused
client.getInputProtocol().getTransport().close();
} catch (Exception e) {
logger.warn("{}: Cannot send heart beat to node {}",
memberName, node, e);
@@ -401,6 +403,11 @@ public class HeartbeatThread implements Runnable {
try {
long result = client.startElection(request);
handler.onComplete(result);
+ } catch (TException e) {
+ // the connection may be broken, close it to avoid it being
reused
+ client.getInputProtocol().getTransport().close();
+ logger.warn("{}: Cannot request a vote from {}", memberName,
node, e);
+ handler.onError(e);
} catch (Exception e) {
handler.onError(e);
} finally {
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
index 6f6936a..89d933a 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
@@ -696,6 +696,7 @@ public class DataGroupMember extends RaftMember {
*/
@Override
public TSStatus executeNonQueryPlan(PhysicalPlan plan) {
+ long startTime = System.currentTimeMillis();
if (ClusterDescriptor.getInstance().getConfig().getReplicationNum() == 1) {
try {
getLocalExecutor().processNonQuery(plan);
@@ -718,6 +719,11 @@ public class DataGroupMember extends RaftMember {
}
}
return handleLogExecutionException(plan, cause);
+ } finally {
+ long elapsed = System.currentTimeMillis() - startTime;
+ if (elapsed > 5000) {
+ logger.error("PlanExecutor execute slowly : time cost : {}ms",
elapsed);
+ }
}
} else {
TSStatus status = executeNonQueryPlanWithKnownLeader(plan);
@@ -725,7 +731,7 @@ public class DataGroupMember extends RaftMember {
return status;
}
- long startTime =
Timer.Statistic.DATA_GROUP_MEMBER_WAIT_LEADER.getOperationStartTime();
+ startTime =
Timer.Statistic.DATA_GROUP_MEMBER_WAIT_LEADER.getOperationStartTime();
waitLeader();
Timer.Statistic.DATA_GROUP_MEMBER_WAIT_LEADER.calOperationCostTimeFromStart(startTime);
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..1cc4ad9 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
@@ -169,7 +169,7 @@ public class MetaGroupMember extends RaftMember {
* every "REFRESH_CLIENT_SEC" seconds, a dataClientRefresher thread will try
to refresh one thrift
* connection for each nodes other than itself.
*/
- private static final int REFRESH_CLIENT_SEC = 5;
+ private static final int REFRESH_CLIENT_SEC = 1;
/** how many times is a data record replicated, also the number of nodes in
a data group */
private static final int REPLICATION_NUM =
@@ -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);
+ logger.info("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);
+ }
}
}
@@ -530,13 +529,13 @@ public class MetaGroupMember extends RaftMember {
try {
client.refreshConnection(new RefreshReuqest(), new
GenericHandler<>(receiver, null));
} catch (TException e) {
- logger.warn("encounter refreshing client timeout, throw broken
connection", e);
+ logger.info("encounter refreshing client timeout, throw broken
connection", e);
}
}
private void generateNodeReport() {
try {
- if (logger.isInfoEnabled()) {
+ if (logger.isDebugEnabled()) {
NodeReport report = genNodeReport();
logger.debug(report.toString());
}
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
index f2187f9..3078013 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
@@ -1317,13 +1317,13 @@ public abstract class RaftMember {
}
public TSStatus forwardPlanSync(PhysicalPlan plan, Node receiver, Node
header, Client client) {
+ long startTime = System.currentTimeMillis();
try {
ExecutNonQueryReq req = new ExecutNonQueryReq();
req.setPlanBytes(PlanSerializer.getInstance().serialize(plan));
if (header != null) {
req.setHeader(header);
}
-
TSStatus tsStatus = client.executeNonQueryPlan(req);
if (tsStatus == null) {
tsStatus = StatusUtils.TIME_OUT;
@@ -1337,7 +1337,12 @@ public abstract class RaftMember {
TSStatus status;
if (e.getCause() instanceof SocketTimeoutException) {
status = StatusUtils.TIME_OUT;
- logger.warn(MSG_FORWARD_TIMEOUT, name, plan, receiver);
+ logger.warn(
+ MSG_FORWARD_TIMEOUT + ": {}ms",
+ name,
+ plan,
+ receiver,
+ System.currentTimeMillis() - startTime);
} else {
logger.error(MSG_FORWARD_ERROR, name, plan, receiver, e);
status = StatusUtils.getStatus(StatusUtils.INTERNAL_ERROR,
e.getMessage());
diff --git
a/cluster/src/test/java/org/apache/iotdb/cluster/client/DataClientProviderTest.java
b/cluster/src/test/java/org/apache/iotdb/cluster/client/DataClientProviderTest.java
index d2ee0b7..3987450 100644
---
a/cluster/src/test/java/org/apache/iotdb/cluster/client/DataClientProviderTest.java
+++
b/cluster/src/test/java/org/apache/iotdb/cluster/client/DataClientProviderTest.java
@@ -28,7 +28,6 @@ import org.apache.iotdb.cluster.rpc.thrift.Node;
import org.apache.iotdb.cluster.utils.ClientUtils;
import org.apache.iotdb.cluster.utils.ClusterNode;
-import org.apache.thrift.TException;
import org.apache.thrift.protocol.TBinaryProtocol.Factory;
import org.junit.After;
import org.junit.Assert;
@@ -96,7 +95,7 @@ public class DataClientProviderTest {
SyncDataClient client = null;
try {
client = provider.getSyncDataClient(node, 100);
- } catch (TException e) {
+ } catch (IOException e) {
Assert.fail(e.getMessage());
} finally {
ClientUtils.putBackSyncClient(client);
@@ -135,7 +134,7 @@ public class DataClientProviderTest {
SyncDataClient client = null;
try {
client = provider.getSyncDataClient(node, 100);
- } catch (TException e) {
+ } catch (IOException e) {
Assert.fail(e.getMessage());
}
assertNotNull(client);
diff --git
a/example/session/src/main/java/org/apache/iotdb/SessionConcurrentExample.java
b/example/session/src/main/java/org/apache/iotdb/SessionConcurrentExample.java
new file mode 100644
index 0000000..fb0f793
--- /dev/null
+++
b/example/session/src/main/java/org/apache/iotdb/SessionConcurrentExample.java
@@ -0,0 +1,198 @@
+/*
+ * 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;
+
+import org.apache.iotdb.rpc.IoTDBConnectionException;
+import org.apache.iotdb.rpc.StatementExecutionException;
+import org.apache.iotdb.session.Session;
+import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
+import org.apache.iotdb.tsfile.write.record.Tablet;
+import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Random;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+public class SessionConcurrentExample {
+
+ private static final int sgNum = 20;
+ private static final int deviceNum = 100;
+ private static final int parallelDegreeForOneSG = 3;
+
+ public static void main(String[] args)
+ throws IoTDBConnectionException, StatementExecutionException {
+
+ Session session = new Session("127.0.0.1", 6667, "root", "root");
+ session.open(false);
+ createTemplate(session);
+ session.close();
+
+ CountDownLatch latch = new CountDownLatch(sgNum * parallelDegreeForOneSG);
+ ExecutorService es = Executors.newFixedThreadPool(sgNum *
parallelDegreeForOneSG);
+
+ for (int i = 0; i < sgNum * parallelDegreeForOneSG; i++) {
+ int currentIndex = i;
+ es.execute(() -> concurrentOperation(latch, currentIndex));
+ }
+
+ es.shutdown();
+
+ try {
+ latch.await();
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ }
+
+ private static void concurrentOperation(CountDownLatch latch, int
currentIndex) {
+
+ Session tmpSession = new Session("127.0.0.1", 6667, "root", "root");
+ try {
+ tmpSession.open(false);
+ } catch (IoTDBConnectionException e) {
+ e.printStackTrace();
+ }
+
+ for (int j = 0; j < deviceNum; j++) {
+ try {
+ insertTablet(
+ tmpSession, String.format("root.sg_%d.d_%d", currentIndex /
parallelDegreeForOneSG, j));
+ } catch (IoTDBConnectionException | StatementExecutionException e) {
+ e.printStackTrace();
+ }
+ }
+
+ try {
+ tmpSession.close();
+ } catch (IoTDBConnectionException e) {
+ e.printStackTrace();
+ }
+
+ latch.countDown();
+ }
+
+ private static void createTemplate(Session session)
+ throws IoTDBConnectionException, StatementExecutionException {
+ List<List<String>> measurementList = new ArrayList<>();
+ measurementList.add(Collections.singletonList("s1"));
+ measurementList.add(Collections.singletonList("s2"));
+ measurementList.add(Collections.singletonList("s3"));
+
+ List<List<TSDataType>> dataTypeList = new ArrayList<>();
+ dataTypeList.add(Collections.singletonList(TSDataType.INT64));
+ dataTypeList.add(Collections.singletonList(TSDataType.INT64));
+ dataTypeList.add(Collections.singletonList(TSDataType.INT64));
+
+ List<List<TSEncoding>> encodingList = new ArrayList<>();
+ encodingList.add(Collections.singletonList(TSEncoding.RLE));
+ encodingList.add(Collections.singletonList(TSEncoding.RLE));
+ encodingList.add(Collections.singletonList(TSEncoding.RLE));
+
+ List<CompressionType> compressionTypes = new ArrayList<>();
+ for (int i = 0; i < 3; i++) {
+ compressionTypes.add(CompressionType.SNAPPY);
+ }
+ List<String> schemaNames = new ArrayList<>();
+ schemaNames.add("s1");
+ schemaNames.add("s2");
+ schemaNames.add("s3");
+
+ session.createSchemaTemplate(
+ "template1", schemaNames, measurementList, dataTypeList, encodingList,
compressionTypes);
+ for (int i = 0; i < sgNum; i++) {
+ session.setSchemaTemplate("template1", "root.sg_" + i);
+ }
+ }
+
+ /**
+ * insert the data of a device. For each timestamp, the number of
measurements is the same.
+ *
+ * <p>Users need to control the count of Tablet and write a batch when it
reaches the maxBatchSize
+ */
+ private static void insertTablet(Session session, String deviceId)
+ throws IoTDBConnectionException, StatementExecutionException {
+ /*
+ * A Tablet example:
+ * device1
+ * time s1, s2, s3
+ * 1, 1, 1, 1
+ * 2, 2, 2, 2
+ * 3, 3, 3, 3
+ */
+ // The schema of measurements of one device
+ // only measurementId and data type in MeasurementSchema take effects in
Tablet
+ List<MeasurementSchema> schemaList = new ArrayList<>();
+ schemaList.add(new MeasurementSchema("s1", TSDataType.INT64));
+ schemaList.add(new MeasurementSchema("s2", TSDataType.INT64));
+ schemaList.add(new MeasurementSchema("s3", TSDataType.INT64));
+
+ Tablet tablet = new Tablet(deviceId, schemaList, 100);
+
+ // Method 1 to add tablet data
+ long timestamp = System.currentTimeMillis();
+
+ for (long row = 0; row < 100; row++) {
+ int rowIndex = tablet.rowSize++;
+ tablet.addTimestamp(rowIndex, timestamp);
+ for (int s = 0; s < 3; s++) {
+ long value = new Random().nextLong();
+ tablet.addValue(schemaList.get(s).getMeasurementId(), rowIndex, value);
+ }
+ if (tablet.rowSize == tablet.getMaxRowNumber()) {
+ session.insertTablet(tablet, true);
+ tablet.reset();
+ }
+ timestamp++;
+ }
+
+ if (tablet.rowSize != 0) {
+ session.insertTablet(tablet);
+ tablet.reset();
+ }
+
+ // Method 2 to add tablet data
+ long[] timestamps = tablet.timestamps;
+ Object[] values = tablet.values;
+
+ for (long time = 0; time < 100; time++) {
+ int row = tablet.rowSize++;
+ timestamps[row] = time;
+ for (int i = 0; i < 3; i++) {
+ long[] sensor = (long[]) values[i];
+ sensor[row] = i;
+ }
+ if (tablet.rowSize == tablet.getMaxRowNumber()) {
+ session.insertTablet(tablet, true);
+ tablet.reset();
+ }
+ }
+
+ if (tablet.rowSize != 0) {
+ session.insertTablet(tablet);
+ tablet.reset();
+ }
+ }
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
index 0402f4b..94c63f0 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
@@ -192,7 +192,12 @@ public class TsFileProcessor {
if (IoTDBDescriptor.getInstance().getConfig().isEnableWal()) {
try {
+ long startTime = System.currentTimeMillis();
getLogNode().write(insertRowPlan);
+ long elapsed = System.currentTimeMillis() - startTime;
+ if (elapsed > 5000) {
+ logger.error("write wal slowly : cost {}ms", elapsed);
+ }
} catch (Exception e) {
throw new WriteProcessException(
String.format(
@@ -248,11 +253,16 @@ public class TsFileProcessor {
}
try {
workMemTable.insertTablet(insertTabletPlan, start, end);
+ long startTime = System.currentTimeMillis();
if (IoTDBDescriptor.getInstance().getConfig().isEnableWal()) {
insertTabletPlan.setStart(start);
insertTabletPlan.setEnd(end);
getLogNode().write(insertTabletPlan);
}
+ long elapsed = System.currentTimeMillis() - startTime;
+ if (elapsed > 5000) {
+ logger.error("write wal slowly : cost {}ms", elapsed);
+ }
} catch (Exception e) {
for (int i = start; i < end; i++) {
results[i] = RpcUtils.getStatus(TSStatusCode.INTERNAL_SERVER_ERROR,
e.getMessage());
diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java
b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java
index 674c90e..7dad6a9 100644
--- a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java
+++ b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java
@@ -2113,7 +2113,15 @@ public class MManager {
}
private void setUsingDeviceTemplate(SetUsingDeviceTemplatePlan plan) throws
MetadataException {
- getDeviceNode(plan.getPrefixPath()).setUseTemplate(true);
+ try {
+ getDeviceNode(plan.getPrefixPath()).setUseTemplate(true);
+ } catch (PathNotExistException e) {
+ // the order of SetUsingDeviceTemplatePlan and AutoCreateDeviceMNodePlan
cannot be guaranteed
+ // during writing currently, so we need a auto-create mechanism here
+ mtree.getDeviceNodeWithAutoCreating(
+ plan.getPrefixPath(), config.getDefaultStorageGroupLevel());
+ getDeviceNode(plan.getPrefixPath()).setUseTemplate(true);
+ }
}
public long getTotalSeriesNumber() {
diff --git
a/server/src/main/java/org/apache/iotdb/db/rescon/PrimitiveArrayManager.java
b/server/src/main/java/org/apache/iotdb/db/rescon/PrimitiveArrayManager.java
index 9881e56..76d2128 100644
--- a/server/src/main/java/org/apache/iotdb/db/rescon/PrimitiveArrayManager.java
+++ b/server/src/main/java/org/apache/iotdb/db/rescon/PrimitiveArrayManager.java
@@ -40,9 +40,17 @@ public class PrimitiveArrayManager {
public static final int ARRAY_SIZE = CONFIG.getPrimitiveArraySize();
+ /**
+ * The actual used memory will be 50% larger than the statistic, so we need
to limit the size of
+ * POOLED_ARRAYS_MEMORY_THRESHOLD, make it smaller than its actual allowed
value.
+ */
+ private static final double AMPLIFICATION_FACTOR = 1.5;
+
/** threshold total size of arrays for all data types */
private static final double POOLED_ARRAYS_MEMORY_THRESHOLD =
- CONFIG.getAllocateMemoryForWrite() *
CONFIG.getBufferedArraysMemoryProportion();
+ CONFIG.getAllocateMemoryForWrite()
+ * CONFIG.getBufferedArraysMemoryProportion()
+ / AMPLIFICATION_FACTOR;
/** TSDataType#serialize() -> ArrayDeque<Array> */
private static final ArrayDeque[] POOLED_ARRAYS = new
ArrayDeque[TSDataType.values().length];
diff --git
a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
index 5be9dca..d1bf262 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
@@ -575,12 +575,12 @@ public class TSServiceImpl implements TSIService.Iface {
@Override
public TSExecuteStatementResp executeStatement(TSExecuteStatementReq req) {
+ String statement = req.getStatement();
try {
if (!checkLogin(req.getSessionId())) {
return
RpcUtils.getTSExecuteStatementResp(TSStatusCode.NOT_LOGIN_ERROR);
}
- String statement = req.getStatement();
PhysicalPlan physicalPlan =
processor.parseSQLToPhysicalPlan(
statement, sessionManager.getZoneId(req.getSessionId()),
req.fetchSize);
@@ -598,9 +598,11 @@ public class TSServiceImpl implements TSIService.Iface {
} catch (InterruptedException e) {
LOGGER.error(INFO_INTERRUPT_ERROR, req, e);
Thread.currentThread().interrupt();
- return RpcUtils.getTSExecuteStatementResp(onQueryException(e, "executing
executeStatement"));
+ return RpcUtils.getTSExecuteStatementResp(
+ onQueryException(e, "executing \"" + statement + "\""));
} catch (Exception e) {
- return RpcUtils.getTSExecuteStatementResp(onQueryException(e, "executing
executeStatement"));
+ return RpcUtils.getTSExecuteStatementResp(
+ onQueryException(e, "executing \"" + statement + "\""));
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
b/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
index c947963..e982c77 100644
---
a/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
+++
b/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
@@ -28,7 +28,6 @@ import org.apache.iotdb.db.writelog.io.ILogWriter;
import org.apache.iotdb.db.writelog.io.LogWriter;
import org.apache.iotdb.db.writelog.io.MultiFileLogReader;
-import com.google.common.util.concurrent.ThreadFactoryBuilder;
import org.apache.commons.io.FileUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -43,6 +42,7 @@ import java.util.Arrays;
import java.util.Comparator;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.ReentrantLock;
/** This WriteLogNode is used to manage insert ahead logs of a TsFile. */
@@ -51,33 +51,32 @@ public class ExclusiveWriteLogNode implements WriteLogNode,
Comparable<Exclusive
public static final String WAL_FILE_NAME = "wal";
private static final Logger logger =
LoggerFactory.getLogger(ExclusiveWriteLogNode.class);
- private String identifier;
+ private final String identifier;
- private String logDirectory;
+ private final String logDirectory;
private ILogWriter currentFileWriter;
- private IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
+ private final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
- private ByteBuffer logBufferWorking;
- private ByteBuffer logBufferIdle;
- private ByteBuffer logBufferFlushing;
+ private volatile ByteBuffer logBufferWorking;
+ private volatile ByteBuffer logBufferIdle;
+ private volatile ByteBuffer logBufferFlushing;
// used for the convenience of deletion
- private ByteBuffer[] bufferArray;
+ private volatile ByteBuffer[] bufferArray;
private final Object switchBufferCondition = new Object();
- private ReentrantLock lock = new ReentrantLock();
- private static final ExecutorService FLUSH_BUFFER_THREAD_POOL =
- Executors.newCachedThreadPool(
- new
ThreadFactoryBuilder().setNameFormat("Flush-WAL-Thread-%d").setDaemon(true).build());
+ private final ReentrantLock lock = new ReentrantLock();
+ private final ExecutorService FLUSH_BUFFER_THREAD_POOL =
+ Executors.newSingleThreadExecutor(r -> new Thread(r, "Flush-WAL-Thread-"
+ this.hashCode()));
private long fileId = 0;
private long lastFlushedId = 0;
private int bufferedLogNum = 0;
- private boolean deleted;
+ private final AtomicBoolean deleted = new AtomicBoolean(false);
/**
* constructor of ExclusiveWriteLogNode.
@@ -102,7 +101,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode,
Comparable<Exclusive
@Override
public void write(PhysicalPlan plan) throws IOException {
- if (deleted) {
+ if (deleted.get()) {
throw new IOException("WAL node deleted");
}
lock.lock();
@@ -138,7 +137,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode,
Comparable<Exclusive
lock.lock();
try {
synchronized (switchBufferCondition) {
- while (logBufferFlushing != null && !deleted) {
+ while (logBufferFlushing != null && !deleted.get()) {
switchBufferCondition.wait();
}
switchBufferCondition.notifyAll();
@@ -151,7 +150,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode,
Comparable<Exclusive
}
logger.debug("Log node {} closed successfully", identifier);
} catch (IOException e) {
- logger.error("Cannot close log node {} because:", identifier, e);
+ logger.warn("Cannot close log node {} because:", identifier, e);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
logger.warn("Waiting for current buffer being flushed interrupted");
@@ -162,7 +161,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode,
Comparable<Exclusive
@Override
public void forceSync() {
- if (deleted) {
+ if (deleted.get()) {
return;
}
sync();
@@ -208,7 +207,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode,
Comparable<Exclusive
try {
close();
FileUtils.deleteDirectory(SystemFileFactory.INSTANCE.getFile(logDirectory));
- deleted = true;
+ deleted.set(true);
return this.bufferArray;
} finally {
lock.unlock();
@@ -232,7 +231,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode,
Comparable<Exclusive
FileUtils.forceDelete(logFile);
logger.info("Log node {} cleaned old file", identifier);
} catch (IOException e) {
- logger.error("Old log file {} of {} cannot be deleted",
logFile.getName(), identifier, e);
+ logger.warn("Old log file {} of {} cannot be deleted",
logFile.getName(), identifier, e);
}
}
}
@@ -245,7 +244,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode,
Comparable<Exclusive
currentFileWriter.force();
}
} catch (IOException e) {
- logger.error("Log node {} force failed.", identifier, e);
+ logger.warn("Log node {} force failed.", identifier, e);
}
} finally {
lock.unlock();
@@ -261,8 +260,6 @@ public class ExclusiveWriteLogNode implements WriteLogNode,
Comparable<Exclusive
switchBufferWorkingToFlushing();
ILogWriter currWriter = getCurrentFileWriter();
FLUSH_BUFFER_THREAD_POOL.submit(() -> flushBuffer(currWriter));
- switchBufferIdleToWorking();
-
bufferedLogNum = 0;
logger.debug("Log node {} ends sync.", identifier);
} catch (InterruptedException e) {
@@ -281,50 +278,28 @@ public class ExclusiveWriteLogNode implements
WriteLogNode, Comparable<Exclusive
} catch (ClosedChannelException e) {
// ignore
} catch (IOException e) {
- logger.error("Log node {} sync failed, change system mode to read-only",
identifier, e);
+ logger.warn("Log node {} sync failed, change system mode to read-only",
identifier, e);
IoTDBDescriptor.getInstance().getConfig().setReadOnly(true);
return;
}
- logBufferFlushing.clear();
-
- try {
- switchBufferFlushingToIdle();
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- }
- }
- private void switchBufferWorkingToFlushing() throws InterruptedException {
+ // switch buffer flushing to idle and notify the sync thread
synchronized (switchBufferCondition) {
- while (logBufferFlushing != null && !deleted) {
- switchBufferCondition.wait();
- }
- logBufferFlushing = logBufferWorking;
- logBufferWorking = null;
+ logBufferIdle = logBufferFlushing;
+ logBufferFlushing = null;
switchBufferCondition.notifyAll();
}
}
- private void switchBufferIdleToWorking() throws InterruptedException {
+ private void switchBufferWorkingToFlushing() throws InterruptedException {
synchronized (switchBufferCondition) {
- while (logBufferIdle == null && !deleted) {
- switchBufferCondition.wait();
+ while (logBufferFlushing != null && !deleted.get()) {
+ switchBufferCondition.wait(100);
}
+ logBufferFlushing = logBufferWorking;
logBufferWorking = logBufferIdle;
+ logBufferWorking.clear();
logBufferIdle = null;
- switchBufferCondition.notifyAll();
- }
- }
-
- private void switchBufferFlushingToIdle() throws InterruptedException {
- synchronized (switchBufferCondition) {
- while (logBufferIdle != null && !deleted) {
- switchBufferCondition.wait();
- }
- logBufferIdle = logBufferFlushing;
- logBufferIdle.clear();
- logBufferFlushing = null;
- switchBufferCondition.notifyAll();
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/writelog/recover/TsFileRecoverPerformer.java
b/server/src/main/java/org/apache/iotdb/db/writelog/recover/TsFileRecoverPerformer.java
index 72ea41c..da53108 100644
---
a/server/src/main/java/org/apache/iotdb/db/writelog/recover/TsFileRecoverPerformer.java
+++
b/server/src/main/java/org/apache/iotdb/db/writelog/recover/TsFileRecoverPerformer.java
@@ -41,6 +41,8 @@ import org.slf4j.LoggerFactory;
import java.io.File;
import java.io.IOException;
import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
@@ -188,13 +190,24 @@ public class TsFileRecoverPerformer {
for (Map.Entry<String, List<ChunkMetadata>> entry :
deviceChunkMetaDataMap.entrySet()) {
String deviceId = entry.getKey();
List<ChunkMetadata> chunkMetadataList = entry.getValue();
- TSDataType dataType = entry.getValue().get(entry.getValue().size() -
1).getDataType();
- for (ChunkMetadata chunkMetaData : chunkMetadataList) {
- if (!chunkMetaData.getDataType().equals(dataType)) {
- continue;
+
+ Map<String, List<ChunkMetadata>> measurementToChunkMetadatas = new
HashMap<>();
+ for (ChunkMetadata chunkMetadata : chunkMetadataList) {
+ List<ChunkMetadata> list =
+ measurementToChunkMetadatas.computeIfAbsent(
+ chunkMetadata.getMeasurementUid(), n -> new ArrayList<>());
+ list.add(chunkMetadata);
+ }
+
+ for (List<ChunkMetadata> metadataList :
measurementToChunkMetadatas.values()) {
+ TSDataType dataType = metadataList.get(metadataList.size() -
1).getDataType();
+ for (ChunkMetadata chunkMetaData : chunkMetadataList) {
+ if (!chunkMetaData.getDataType().equals(dataType)) {
+ continue;
+ }
+ tsFileResource.updateStartTime(deviceId,
chunkMetaData.getStartTime());
+ tsFileResource.updateEndTime(deviceId, chunkMetaData.getEndTime());
}
- tsFileResource.updateStartTime(deviceId, chunkMetaData.getStartTime());
- tsFileResource.updateEndTime(deviceId, chunkMetaData.getEndTime());
}
}
tsFileResource.updatePlanIndexes(restorableTsFileIOWriter.getMinPlanIndex());
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
index d01b9d9..fa4f7f5 100644
--- a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
+++ b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
@@ -18,6 +18,8 @@
*/
package org.apache.iotdb.db.integration;
+import org.apache.iotdb.db.conf.IoTDBConfig;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.engine.StorageEngine;
import org.apache.iotdb.db.engine.compaction.CompactionMergeTaskPoolManager;
import org.apache.iotdb.db.exception.StorageEngineException;
@@ -353,6 +355,52 @@ public class IoTDBRestartIT {
}
@Test
+ public void testRecoverWALDeleteSchemaCheckResourceTime() throws Exception {
+ EnvironmentUtils.envSetUp();
+ Class.forName(Config.JDBC_DRIVER_NAME);
+ IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
+ int avgSeriesPointNumberThreshold =
config.getAvgSeriesPointNumberThreshold();
+ config.setAvgSeriesPointNumberThreshold(2);
+ long tsfileSize = config.getSeqTsFileSize();
+ config.setSeqTsFileSize(10000000);
+
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+ statement.execute("create timeseries root.turbine1.d1.s1 with
datatype=INT64");
+ statement.execute("insert into root.turbine1.d1(timestamp,s1)
values(1,1)");
+ statement.execute("insert into root.turbine1.d1(timestamp,s1)
values(2,1)");
+ statement.execute("create timeseries root.turbine1.d1.s2 with
datatype=BOOLEAN");
+ statement.execute("insert into root.turbine1.d1(timestamp,s2)
values(3,true)");
+ statement.execute("insert into root.turbine1.d1(timestamp,s2)
values(4,true)");
+ }
+
+ Thread.sleep(1000);
+ EnvironmentUtils.restartDaemon();
+
+ try (Connection connection =
+ DriverManager.getConnection(
+ Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+ Statement statement = connection.createStatement()) {
+
+ long[] result = new long[] {1L, 2L};
+ statement.execute("select s1 from root.turbine1.d1 where time < 3");
+ ResultSet resultSet = statement.getResultSet();
+ int cnt = 0;
+ while (resultSet.next()) {
+ assertEquals(resultSet.getLong(1), result[cnt]);
+ cnt++;
+ }
+ assertEquals(2, cnt);
+ }
+
+ config.setAvgSeriesPointNumberThreshold(avgSeriesPointNumberThreshold);
+ config.setSeqTsFileSize(tsfileSize);
+ EnvironmentUtils.cleanEnv();
+ }
+
+ @Test
public void testRestartCompaction()
throws SQLException, ClassNotFoundException, IOException,
StorageEngineException {
EnvironmentUtils.envSetUp();
diff --git
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/expression/util/ExpressionOptimizer.java
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/expression/util/ExpressionOptimizer.java
index ea2a24f..b056286 100644
---
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/expression/util/ExpressionOptimizer.java
+++
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/expression/util/ExpressionOptimizer.java
@@ -70,21 +70,25 @@ public class ExpressionOptimizer {
(GlobalTimeExpression) right, left, selectedSeries, relation);
} else if (left.getType() != ExpressionType.GLOBAL_TIME
&& right.getType() != ExpressionType.GLOBAL_TIME) {
- IExpression regularLeft = optimize(left, selectedSeries);
- IExpression regularRight = optimize(right, selectedSeries);
- IBinaryExpression midRet = null;
- if (relation == ExpressionType.AND) {
- midRet = BinaryExpression.and(regularLeft, regularRight);
- } else if (relation == ExpressionType.OR) {
- midRet = BinaryExpression.or(regularLeft, regularRight);
- } else {
- throw new UnsupportedOperationException("unsupported IExpression
type: " + relation);
- }
- if (midRet.getLeft().getType() == ExpressionType.GLOBAL_TIME
- || midRet.getRight().getType() == ExpressionType.GLOBAL_TIME) {
- return optimize(midRet, selectedSeries);
- } else {
- return midRet;
+ try {
+ IExpression regularLeft = optimize(left, selectedSeries);
+ IExpression regularRight = optimize(right, selectedSeries);
+ IBinaryExpression midRet = null;
+ if (relation == ExpressionType.AND) {
+ midRet = BinaryExpression.and(regularLeft, regularRight);
+ } else if (relation == ExpressionType.OR) {
+ midRet = BinaryExpression.or(regularLeft, regularRight);
+ } else {
+ throw new UnsupportedOperationException("unsupported IExpression
type: " + relation);
+ }
+ if (midRet.getLeft().getType() == ExpressionType.GLOBAL_TIME
+ || midRet.getRight().getType() == ExpressionType.GLOBAL_TIME) {
+ return optimize(midRet, selectedSeries);
+ } else {
+ return midRet;
+ }
+ } catch (StackOverflowError stackOverflowError) {
+ throw new QueryFilterOptimizationException("StackOverflowError is
encountered.");
}
}
}