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 @@ -3818,5 +3818,8 @@ private DataNodeQueryMessages() {}
public static final String EXCEPTION_VISIBLEALIASES_IS_NULL_630B27F1 = "visibleAliases is null";
public static final String EXCEPTION_HAS_NO_PERMISSION_TO_EXECUTE_ARG_BECAUSE_ONLY_THE_SUPERUSER_CAN_ALTER_HIM_HERSELF_C5902893 =
"Has no permission to execute %s, because only the superuser can alter him/herself.";
public static final String
LOG_FAILED_TO_CLEAN_DEVICEENTRY_DATA_SET_ASYNCHRONOUSLY_QUERYID_ARG_PLANNODEID_ARG_9106C4C5 =
"Failed to clean DeviceEntry data set asynchronously: queryId=%s, planNodeId=%s";

}
Original file line number Diff line number Diff line change
Expand Up @@ -4575,5 +4575,8 @@ private DataNodeQueryMessages() {}
public static final String EXCEPTION_VISIBLEALIASES_IS_NULL_630B27F1 = "visibleAliases 不能为空";
public static final String EXCEPTION_HAS_NO_PERMISSION_TO_EXECUTE_ARG_BECAUSE_ONLY_THE_SUPERUSER_CAN_ALTER_HIM_HERSELF_C5902893 =
"无权执行 %s,因为只有超级用户可以修改其自身。";
public static final String
LOG_FAILED_TO_CLEAN_DEVICEENTRY_DATA_SET_ASYNCHRONOUSLY_QUERYID_ARG_PLANNODEID_ARG_9106C4C5 =
"异步清理 DeviceEntry 数据集失败:queryId=%s,planNodeId=%s";

}
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 @@ -43,11 +43,14 @@
import org.apache.iotdb.db.queryengine.execution.memory.LocalMemoryManager;
import org.apache.iotdb.db.queryengine.metric.DataExchangeCostMetricSet;
import org.apache.iotdb.db.queryengine.metric.DataExchangeCountMetricSet;
import org.apache.iotdb.db.queryengine.plan.relational.metadata.spill.DeviceEntrySpillManager;
import org.apache.iotdb.db.utils.SetThreadName;
import org.apache.iotdb.mpp.rpc.thrift.MPPDataExchangeService;
import org.apache.iotdb.mpp.rpc.thrift.TAcknowledgeDataBlockEvent;
import org.apache.iotdb.mpp.rpc.thrift.TCloseSinkChannelEvent;
import org.apache.iotdb.mpp.rpc.thrift.TEndOfDataBlockEvent;
import org.apache.iotdb.mpp.rpc.thrift.TFetchDeviceEntrySegmentReq;
import org.apache.iotdb.mpp.rpc.thrift.TFetchDeviceEntrySegmentResp;
import org.apache.iotdb.mpp.rpc.thrift.TFragmentInstanceId;
import org.apache.iotdb.mpp.rpc.thrift.TGetDataBlockRequest;
import org.apache.iotdb.mpp.rpc.thrift.TGetDataBlockResponse;
Expand Down Expand Up @@ -96,6 +99,46 @@ class MPPDataExchangeServiceImpl implements MPPDataExchangeService.Iface {
private final DataExchangeCountMetricSet DATA_EXCHANGE_COUNT_METRICS =
DataExchangeCountMetricSet.getInstance();

@Override
public TFetchDeviceEntrySegmentResp fetchDeviceEntrySegment(
TFetchDeviceEntrySegmentReq request) {
try {
DeviceEntrySpillManager spillManager = DeviceEntrySpillManager.getInstance();
byte[] payload =
spillManager.readSegment(
request.getQueryId(), request.getPlanNodeId(), request.getSegmentId());
if (request.getSegmentId() > 0) {
spillManager.deleteSegment(
request.getQueryId(), request.getPlanNodeId(), request.getSegmentId() - 1);
}
return new TFetchDeviceEntrySegmentResp(
new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()))
.setPayload(payload);
} catch (IOException | RuntimeException e) {
return new TFetchDeviceEntrySegmentResp(
new TSStatus(TSStatusCode.INTERNAL_SERVER_ERROR.getStatusCode())
.setMessage(e.getMessage()));
}
}

@Override
public TSStatus finishDeviceEntrySegment(String queryId, String planNodeId) {
executorService.submit(
() -> {
try {
DeviceEntrySpillManager.getInstance().finishSegmentDataSet(queryId, planNodeId);
} catch (IOException | RuntimeException e) {
LOGGER.warn(
DataNodeQueryMessages
.LOG_FAILED_TO_CLEAN_DEVICEENTRY_DATA_SET_ASYNCHRONOUSLY_QUERYID_ARG_PLANNODEID_ARG_9106C4C5,
queryId,
planNodeId,
e);
}
});
return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
}

@Override
public TGetDataBlockResponse getDataBlock(TGetDataBlockRequest req) throws TException {
long startTime = System.nanoTime();
Expand Down Expand Up @@ -623,6 +666,11 @@ public MPPDataExchangeServiceImpl getOrCreateMPPDataExchangeServiceImpl() {
return mppDataExchangeService;
}

public IClientManager<TEndPoint, SyncDataNodeMPPDataExchangeServiceClient>
getMppDataExchangeServiceClientManager() {
return mppDataExchangeServiceClientManager;
}

public void deRegisterFragmentInstanceFromMemoryPool(
String queryId, String fragmentInstanceId, boolean forceDeregister) {
localMemoryManager
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