Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -256,6 +256,9 @@ public class IoTDBConfig {
private String queryDir =
IoTDBConstant.DN_DEFAULT_DATA_DIR + File.separator + IoTDBConstant.QUERY_FOLDER_NAME;

/** Maximum DeviceEntry bytes kept in memory before a table-query spill. */
private long tableQueryDeviceEntryBatchSizeInBytes;

/** External lib directory, stores user-uploaded JAR files */
private String extDir = IoTDBConstant.EXT_FOLDER_NAME;

Expand Down Expand Up @@ -1789,6 +1792,14 @@ public void setQueryDir(String queryDir) {
this.queryDir = queryDir;
}

public long getTableQueryDeviceEntryBatchSizeInBytes() {
return tableQueryDeviceEntryBatchSizeInBytes;
}

public void setTableQueryDeviceEntryBatchSizeInBytes(long tableQueryDeviceEntryBatchSizeInBytes) {
this.tableQueryDeviceEntryBatchSizeInBytes = tableQueryDeviceEntryBatchSizeInBytes;
}

public String getRatisDataRegionSnapshotDir() {
return ratisDataRegionSnapshotDir;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -352,6 +352,18 @@ public void loadProperties(TrimProperties properties) throws BadNodeUrlException

conf.setQueryDir(
FilePathUtils.regularizePath(conf.getSystemDir() + IoTDBConstant.QUERY_FOLDER_NAME));
long deviceEntryBatchSize =
Long.parseLong(
properties.getProperty(
"table_query_device_entry_batch_size_in_bytes",
Long.toString(conf.getTableQueryDeviceEntryBatchSizeInBytes())));
if (deviceEntryBatchSize <= 0) {
deviceEntryBatchSize =
memoryConfig.getOperatorsMemoryManager().getTotalMemorySizeInBytes()
/ memoryConfig.getQueryThreadCount()
/ 4;
}
conf.setTableQueryDeviceEntryBatchSizeInBytes(deviceEntryBatchSize);
String[] defaultTierDirs = new String[conf.getTierDataDirs().length];
for (int i = 0; i < defaultTierDirs.length; ++i) {
defaultTierDirs[i] = String.join(",", conf.getTierDataDirs()[i]);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@
import org.apache.iotdb.db.queryengine.plan.planner.LocalExecutionPlanner;
import org.apache.iotdb.db.queryengine.plan.planner.memory.NotThreadSafeMemoryReservationManager;
import org.apache.iotdb.db.queryengine.plan.relational.function.tvf.read_tsfile.ExternalTsFileQueryResource;
import org.apache.iotdb.db.queryengine.plan.relational.metadata.spill.DeviceEntryIOContext;
import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.ExplainOutputFormat;
import org.apache.iotdb.db.queryengine.statistics.QueryPlanStatistics;

Expand Down Expand Up @@ -132,6 +133,8 @@ public enum ExplainType {

private QueryPlanStatistics queryPlanStatistics = null;

private DeviceEntryIOContext deviceEntryIOContext;

// To avoid query front-end from consuming too much memory, it needs to reserve memory when
// constructing some Expression and PlanNode.
private final MemoryReservationManager memoryReservationManager;
Expand Down Expand Up @@ -403,6 +406,13 @@ public void setStartTime(long startTime) {
this.startTime = startTime;
}

public DeviceEntryIOContext getOrCreateDeviceEntryIOContext(boolean duringFetchSchema) {
if (deviceEntryIOContext == null) {
deviceEntryIOContext = new DeviceEntryIOContext(this, duringFetchSchema);
}
return deviceEntryIOContext;
}

public void addFailedEndPoint(TEndPoint endPoint) {
this.endPointBlackList.add(endPoint);
}
Expand Down Expand Up @@ -528,6 +538,37 @@ public long getDispatchCost() {
return queryPlanStatistics.getDispatchCost();
}

public void recordDeviceEntryDiskIODuringFetchSchema(long bytes, long timeCost) {
getOrCreateQueryPlanStatistics().recordDeviceEntryDiskIODuringFetchSchema(bytes, timeCost);
}

public void recordDeviceEntryCount(long count) {
getOrCreateQueryPlanStatistics().recordDeviceEntryCount(count);
}

public long getDiskIOSizeForDeviceEntryDuringFetchSchema() {
return queryPlanStatistics == null
? 0
: queryPlanStatistics.getDiskIOSizeForDeviceEntryDuringFetchSchema();
}

public long getDiskIOTimeCostForDeviceEntryDuringFetchSchema() {
return queryPlanStatistics == null
? 0
: queryPlanStatistics.getDiskIOTimeCostForDeviceEntryDuringFetchSchema();
}

public long getDeviceEntryCount() {
return queryPlanStatistics == null ? 0 : queryPlanStatistics.getDeviceEntryCount();
}

private QueryPlanStatistics getOrCreateQueryPlanStatistics() {
if (queryPlanStatistics == null) {
queryPlanStatistics = new QueryPlanStatistics();
}
return queryPlanStatistics;
}

public void setAnalyzeCost(long analyzeCost) {
if (queryPlanStatistics == null) {
queryPlanStatistics = new QueryPlanStatistics();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,8 @@
import org.apache.iotdb.db.protocol.client.ConfigNodeClientManager;
import org.apache.iotdb.db.protocol.client.ConfigNodeInfo;
import org.apache.iotdb.db.queryengine.plan.analyze.cache.partition.PartitionCache;
import org.apache.iotdb.db.queryengine.plan.relational.metadata.spill.DeviceEntryDataSet;
import org.apache.iotdb.db.queryengine.plan.relational.metadata.spill.DeviceEntryReader;
import org.apache.iotdb.mpp.rpc.thrift.TRegionRouteReq;
import org.apache.iotdb.rpc.TSStatusCode;

Expand All @@ -57,6 +59,7 @@
import javax.annotation.Nullable;

import java.io.IOException;
import java.io.UncheckedIOException;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
Expand Down Expand Up @@ -207,28 +210,26 @@ public DataPartition getDataPartition(
final Map<String, List<DataPartitionQueryParam>> sgNameToQueryParamsMap) {
DataPartition dataPartition = partitionCache.getDataPartition(sgNameToQueryParamsMap);
if (null == dataPartition) {
try (ConfigNodeClient client =
configNodeClientManager.borrowClient(ConfigNodeInfo.CONFIG_REGION_ID)) {
final TDataPartitionTableResp dataPartitionTableResp =
client.getDataPartitionTable(constructDataPartitionReqForQuery(sgNameToQueryParamsMap));
if (dataPartitionTableResp.getStatus().getCode()
== TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
dataPartition = parseDataPartitionResp(dataPartitionTableResp);
partitionCache.updateDataPartitionCache(dataPartitionTableResp.getDataPartitionTable());
} else {
throw new StatementAnalyzeException(
String.format(
DataNodeQueryMessages
.QUERY_EXCEPTION_AN_ERROR_OCCURRED_WHEN_EXECUTING_GETDATAPARTITION_S_D21A0011,
dataPartitionTableResp.getStatus().getMessage()));
}
} catch (final ClientManagerException | TException e) {
throw new StatementAnalyzeException(
String.format(
DataNodeQueryMessages
.QUERY_EXCEPTION_AN_ERROR_OCCURRED_WHEN_EXECUTING_GETDATAPARTITION_S_D21A0011,
e.getMessage()));
}
dataPartition =
fetchDataPartition(constructDataPartitionReqForQuery(sgNameToQueryParamsMap), true);
}
return dataPartition;
}

@Override
public DataPartition getDataPartition(
final String database,
final DeviceEntryDataSet dataSet,
final List<TTimePartitionSlot> timePartitionSlots) {
final Set<TSeriesPartitionSlot> seriesPartitionSlots = collectSeriesPartitionSlots(dataSet);
DataPartition dataPartition =
partitionCache.getDataPartition(database, seriesPartitionSlots, timePartitionSlots);
if (null == dataPartition) {
dataPartition =
fetchDataPartition(
constructDataPartitionReqForQuery(
database, seriesPartitionSlots, timePartitionSlots, false, false),
true);
}
return dataPartition;
}
Expand All @@ -239,20 +240,41 @@ public DataPartition getDataPartitionWithUnclosedTimeRange(
// In this method, we must fetch from config node because it contains -oo or +oo
// and there is no need to update cache because since we will never fetch it from cache, the
// update operation will be only time waste
return fetchDataPartition(constructDataPartitionReqForQuery(sgNameToQueryParamsMap), false);
}

@Override
public DataPartition getDataPartitionWithUnclosedTimeRange(
final String database,
final DeviceEntryDataSet dataSet,
final List<TTimePartitionSlot> timePartitionSlots,
final boolean needLeftAll,
final boolean needRightAll) {
final Set<TSeriesPartitionSlot> seriesPartitionSlots = collectSeriesPartitionSlots(dataSet);
return fetchDataPartition(
constructDataPartitionReqForQuery(
database, seriesPartitionSlots, timePartitionSlots, needLeftAll, needRightAll),
false);
}

private DataPartition fetchDataPartition(
final TDataPartitionReq request, final boolean updateCache) {
try (final ConfigNodeClient client =
configNodeClientManager.borrowClient(ConfigNodeInfo.CONFIG_REGION_ID)) {
final TDataPartitionTableResp dataPartitionTableResp =
client.getDataPartitionTable(constructDataPartitionReqForQuery(sgNameToQueryParamsMap));
final TDataPartitionTableResp dataPartitionTableResp = client.getDataPartitionTable(request);
if (dataPartitionTableResp.getStatus().getCode()
== TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
return parseDataPartitionResp(dataPartitionTableResp);
} else {
throw new StatementAnalyzeException(
String.format(
DataNodeQueryMessages
.QUERY_EXCEPTION_AN_ERROR_OCCURRED_WHEN_EXECUTING_GETDATAPARTITION_S_D21A0011,
dataPartitionTableResp.getStatus().getMessage()));
final DataPartition dataPartition = parseDataPartitionResp(dataPartitionTableResp);
if (updateCache) {
partitionCache.updateDataPartitionCache(dataPartitionTableResp.getDataPartitionTable());
}
return dataPartition;
}
throw new StatementAnalyzeException(
String.format(
DataNodeQueryMessages
.QUERY_EXCEPTION_AN_ERROR_OCCURRED_WHEN_EXECUTING_GETDATAPARTITION_S_D21A0011,
dataPartitionTableResp.getStatus().getMessage()));
} catch (final ClientManagerException | TException e) {
throw new StatementAnalyzeException(
String.format(
Expand Down Expand Up @@ -543,6 +565,34 @@ private TDataPartitionReq constructDataPartitionReqForQuery(
return new TDataPartitionReq(partitionSlotsMap);
}

private TDataPartitionReq constructDataPartitionReqForQuery(
final String database,
final Set<TSeriesPartitionSlot> seriesPartitionSlots,
final List<TTimePartitionSlot> timePartitionSlots,
final boolean needLeftAll,
final boolean needRightAll) {
final TTimeSlotList sharedTimeSlotList =
new TTimeSlotList(timePartitionSlots, needLeftAll, needRightAll);
final Map<TSeriesPartitionSlot, TTimeSlotList> seriesSlotToTimeSlots = new HashMap<>();
for (final TSeriesPartitionSlot seriesPartitionSlot : seriesPartitionSlots) {
seriesSlotToTimeSlots.put(seriesPartitionSlot, sharedTimeSlotList);
}
return new TDataPartitionReq(Collections.singletonMap(database, seriesSlotToTimeSlots));
}

private Set<TSeriesPartitionSlot> collectSeriesPartitionSlots(final DeviceEntryDataSet dataSet) {
final Set<TSeriesPartitionSlot> seriesPartitionSlots = new HashSet<>();
try (final DeviceEntryReader reader = dataSet.openReader()) {
while (reader.hasNext()) {
seriesPartitionSlots.add(
partitionExecutor.getSeriesPartitionSlot(reader.next().getDeviceID()));
}
} catch (final IOException e) {
throw new UncheckedIOException(e);
}
return seriesPartitionSlots;
}

private SchemaPartition parseSchemaPartitionTableResp(
final TSchemaPartitionTableResp schemaPartitionTableResp) {
final Map<String, Map<TSeriesPartitionSlot, TRegionReplicaSet>> regionReplicaMap =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,13 @@

package org.apache.iotdb.db.queryengine.plan.analyze;

import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot;
import org.apache.iotdb.commons.partition.DataPartition;
import org.apache.iotdb.commons.partition.DataPartitionQueryParam;
import org.apache.iotdb.commons.partition.SchemaNodeManagementPartition;
import org.apache.iotdb.commons.partition.SchemaPartition;
import org.apache.iotdb.commons.path.PathPatternTree;
import org.apache.iotdb.db.queryengine.plan.relational.metadata.spill.DeviceEntryDataSet;
import org.apache.iotdb.mpp.rpc.thrift.TRegionRouteReq;

import org.apache.tsfile.file.metadata.IDeviceID;
Expand Down Expand Up @@ -59,6 +61,9 @@ default SchemaPartition getSchemaPartition(
*/
DataPartition getDataPartition(Map<String, List<DataPartitionQueryParam>> sgNameToQueryParamsMap);

DataPartition getDataPartition(
String database, DeviceEntryDataSet dataSet, List<TTimePartitionSlot> timePartitionSlots);

/**
* Get data partition, used in query scenarios which contains time filter like: time < XX or time
* > XX
Expand All @@ -68,6 +73,13 @@ default SchemaPartition getSchemaPartition(
DataPartition getDataPartitionWithUnclosedTimeRange(
Map<String, List<DataPartitionQueryParam>> sgNameToQueryParamsMap);

DataPartition getDataPartitionWithUnclosedTimeRange(
String database,
DeviceEntryDataSet dataSet,
List<TTimePartitionSlot> timePartitionSlots,
boolean needLeftAll,
boolean needRightAll);

/**
* Get or create data partition, used in standalone write scenarios. if enableAutoCreateSchema is
* true and database/series/time slots not exists, then automatically create.
Expand Down
Loading