This is an automated email from the ASF dual-hosted git repository. hui pushed a commit to branch lmh/addMetricDoc in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 74908eccce3d893ee15075b4cc623361847c5c71 Author: Minghui Liu <[email protected]> AuthorDate: Tue Jan 10 16:40:53 2023 +0800 add more metrics --- .../db/mpp/execution/exchange/LocalSinkHandle.java | 2 +- .../mpp/execution/exchange/LocalSourceHandle.java | 4 +- .../execution/exchange/MPPDataExchangeManager.java | 14 ++- .../db/mpp/execution/exchange/SinkHandle.java | 6 +- .../db/mpp/execution/exchange/SourceHandle.java | 14 ++- ...tricSet.java => DataExchangeCostMetricSet.java} | 10 +- .../db/mpp/metric/DataExchangeCountMetricSet.java | 119 +++++++++++++++++++++ .../iotdb/db/mpp/metric/QueryMetricsManager.java | 13 +-- .../db/service/metrics/DataNodeMetricsHelper.java | 6 +- 9 files changed, 155 insertions(+), 33 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/LocalSinkHandle.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/LocalSinkHandle.java index 0a1d666409..8afb07fe6e 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/LocalSinkHandle.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/LocalSinkHandle.java @@ -33,7 +33,7 @@ import java.util.List; import java.util.Optional; import static com.google.common.util.concurrent.Futures.nonCancellationPropagating; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.SINK_HANDLE_SEND_TSBLOCK_LOCAL; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.SINK_HANDLE_SEND_TSBLOCK_LOCAL; public class LocalSinkHandle implements ISinkHandle { diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/LocalSourceHandle.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/LocalSourceHandle.java index 6487a20399..f47dea04d6 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/LocalSourceHandle.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/LocalSourceHandle.java @@ -37,8 +37,8 @@ import java.nio.ByteBuffer; import static com.google.common.util.concurrent.Futures.nonCancellationPropagating; import static org.apache.iotdb.db.mpp.execution.exchange.MPPDataExchangeManager.createFullIdFrom; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.SOURCE_HANDLE_DESERIALIZE_TSBLOCK_LOCAL; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.SOURCE_HANDLE_GET_TSBLOCK_LOCAL; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.SOURCE_HANDLE_DESERIALIZE_TSBLOCK_LOCAL; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.SOURCE_HANDLE_GET_TSBLOCK_LOCAL; public class LocalSourceHandle implements ISourceHandle { diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/MPPDataExchangeManager.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/MPPDataExchangeManager.java index 885447d0c0..b3d89c0340 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/MPPDataExchangeManager.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/MPPDataExchangeManager.java @@ -54,9 +54,12 @@ import java.util.concurrent.ExecutorService; import java.util.function.Supplier; import static org.apache.iotdb.db.mpp.common.FragmentInstanceId.createFullId; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.GET_DATA_BLOCK_TASK_SERVER; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.ON_ACKNOWLEDGE_DATA_BLOCK_EVENT_TASK_SERVER; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.SEND_NEW_DATA_BLOCK_EVENT_TASK_SERVER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.GET_DATA_BLOCK_TASK_SERVER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.ON_ACKNOWLEDGE_DATA_BLOCK_EVENT_TASK_SERVER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.SEND_NEW_DATA_BLOCK_EVENT_TASK_SERVER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCountMetricSet.GET_DATA_BLOCK_NUM_SERVER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCountMetricSet.ON_ACKNOWLEDGE_DATA_BLOCK_NUM_SERVER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCountMetricSet.SEND_NEW_DATA_BLOCK_NUM_SERVER; public class MPPDataExchangeManager implements IMPPDataExchangeManager { @@ -118,6 +121,8 @@ public class MPPDataExchangeManager implements IMPPDataExchangeManager { } finally { QUERY_METRICS.recordDataExchangeCost( GET_DATA_BLOCK_TASK_SERVER, System.nanoTime() - startTime); + QUERY_METRICS.recordDataBlockNum( + GET_DATA_BLOCK_NUM_SERVER, req.getEndSequenceId() - req.getStartSequenceId()); } } @@ -150,6 +155,8 @@ public class MPPDataExchangeManager implements IMPPDataExchangeManager { } finally { QUERY_METRICS.recordDataExchangeCost( ON_ACKNOWLEDGE_DATA_BLOCK_EVENT_TASK_SERVER, System.nanoTime() - startTime); + QUERY_METRICS.recordDataBlockNum( + ON_ACKNOWLEDGE_DATA_BLOCK_NUM_SERVER, e.getEndSequenceId() - e.getStartSequenceId()); } } @@ -192,6 +199,7 @@ public class MPPDataExchangeManager implements IMPPDataExchangeManager { } finally { QUERY_METRICS.recordDataExchangeCost( SEND_NEW_DATA_BLOCK_EVENT_TASK_SERVER, System.nanoTime() - startTime); + QUERY_METRICS.recordDataBlockNum(SEND_NEW_DATA_BLOCK_NUM_SERVER, e.getBlockSizes().size()); } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SinkHandle.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SinkHandle.java index e2494af0fb..139387d510 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SinkHandle.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SinkHandle.java @@ -51,8 +51,9 @@ import java.util.concurrent.ExecutorService; import static com.google.common.util.concurrent.Futures.nonCancellationPropagating; import static org.apache.iotdb.db.mpp.common.FragmentInstanceId.createFullId; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.SEND_NEW_DATA_BLOCK_EVENT_TASK_CALLER; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.SINK_HANDLE_SEND_TSBLOCK_REMOTE; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.SEND_NEW_DATA_BLOCK_EVENT_TASK_CALLER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.SINK_HANDLE_SEND_TSBLOCK_REMOTE; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCountMetricSet.SEND_NEW_DATA_BLOCK_NUM_CALLER; import static org.apache.iotdb.tsfile.read.common.block.TsBlockBuilderStatus.DEFAULT_MAX_TSBLOCK_SIZE_IN_BYTES; public class SinkHandle implements ISinkHandle { @@ -428,6 +429,7 @@ public class SinkHandle implements ISinkHandle { } finally { QUERY_METRICS.recordDataExchangeCost( SEND_NEW_DATA_BLOCK_EVENT_TASK_CALLER, System.nanoTime() - startTime); + QUERY_METRICS.recordDataBlockNum(SEND_NEW_DATA_BLOCK_NUM_CALLER, blockSizes.size()); } } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SourceHandle.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SourceHandle.java index 2e1a4b8a6d..d164e9d1de 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SourceHandle.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SourceHandle.java @@ -51,10 +51,12 @@ import java.util.concurrent.ExecutorService; import static com.google.common.util.concurrent.Futures.nonCancellationPropagating; import static org.apache.iotdb.db.mpp.execution.exchange.MPPDataExchangeManager.createFullIdFrom; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.GET_DATA_BLOCK_TASK_CALLER; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.ON_ACKNOWLEDGE_DATA_BLOCK_EVENT_TASK_CALLER; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.SOURCE_HANDLE_DESERIALIZE_TSBLOCK_REMOTE; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.SOURCE_HANDLE_GET_TSBLOCK_REMOTE; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.GET_DATA_BLOCK_TASK_CALLER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.ON_ACKNOWLEDGE_DATA_BLOCK_EVENT_TASK_CALLER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.SOURCE_HANDLE_DESERIALIZE_TSBLOCK_REMOTE; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet.SOURCE_HANDLE_GET_TSBLOCK_REMOTE; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCountMetricSet.GET_DATA_BLOCK_NUM_CALLER; +import static org.apache.iotdb.db.mpp.metric.DataExchangeCountMetricSet.ON_ACKNOWLEDGE_DATA_BLOCK_NUM_CALLER; public class SourceHandle implements ISourceHandle { @@ -467,7 +469,7 @@ public class SourceHandle implements ISourceHandle { tsBlocks.addAll(resp.getTsBlocks()); logger.debug("[EndPullTsBlocksFromRemote] Count:{}", tsBlockNum); - QUERY_METRICS.recordDataBlockNum(tsBlockNum); + QUERY_METRICS.recordDataBlockNum(GET_DATA_BLOCK_NUM_CALLER, tsBlockNum); executorService.submit( new SendAcknowledgeDataBlockEventTask(startSequenceId, endSequenceId)); synchronized (SourceHandle.this) { @@ -581,6 +583,8 @@ public class SourceHandle implements ISourceHandle { } finally { QUERY_METRICS.recordDataExchangeCost( ON_ACKNOWLEDGE_DATA_BLOCK_EVENT_TASK_CALLER, System.nanoTime() - startTime); + QUERY_METRICS.recordDataBlockNum( + ON_ACKNOWLEDGE_DATA_BLOCK_NUM_CALLER, endSequenceId - startSequenceId); } } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/metric/DataExchangeMetricSet.java b/server/src/main/java/org/apache/iotdb/db/mpp/metric/DataExchangeCostMetricSet.java similarity index 94% rename from server/src/main/java/org/apache/iotdb/db/mpp/metric/DataExchangeMetricSet.java rename to server/src/main/java/org/apache/iotdb/db/mpp/metric/DataExchangeCostMetricSet.java index 004b640d71..aab7d1ba5d 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/metric/DataExchangeMetricSet.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/metric/DataExchangeCostMetricSet.java @@ -30,7 +30,7 @@ import org.apache.iotdb.metrics.utils.MetricType; import java.util.HashMap; import java.util.Map; -public class DataExchangeMetricSet implements IMetricSet { +public class DataExchangeCostMetricSet implements IMetricSet { private static final String metric = Metric.DATA_EXCHANGE_COST.toString(); @@ -170,19 +170,12 @@ public class DataExchangeMetricSet implements IMetricSet { "server")); } - public static final String GET_DATA_BLOCK_NUM = "get_data_block_num"; - @Override public void bindTo(AbstractMetricService metricService) { for (MetricInfo metricInfo : metricInfoMap.values()) { metricService.getOrCreateTimer( metricInfo.getName(), MetricLevel.IMPORTANT, metricInfo.getTagsInArray()); } - metricService.getOrCreateHistogram( - Metric.DATA_EXCHANGE_COUNT.toString(), - MetricLevel.IMPORTANT, - Tag.NAME.toString(), - GET_DATA_BLOCK_NUM); } @Override @@ -190,6 +183,5 @@ public class DataExchangeMetricSet implements IMetricSet { for (MetricInfo metricInfo : metricInfoMap.values()) { metricService.remove(MetricType.TIMER, metric, metricInfo.getTagsInArray()); } - metricService.remove(MetricType.HISTOGRAM, metric, Tag.NAME.toString(), GET_DATA_BLOCK_NUM); } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/metric/DataExchangeCountMetricSet.java b/server/src/main/java/org/apache/iotdb/db/mpp/metric/DataExchangeCountMetricSet.java new file mode 100644 index 0000000000..5d8ef839f4 --- /dev/null +++ b/server/src/main/java/org/apache/iotdb/db/mpp/metric/DataExchangeCountMetricSet.java @@ -0,0 +1,119 @@ +/* + * 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.db.mpp.metric; + +import org.apache.iotdb.commons.service.metric.enums.Metric; +import org.apache.iotdb.commons.service.metric.enums.Tag; +import org.apache.iotdb.metrics.AbstractMetricService; +import org.apache.iotdb.metrics.metricsets.IMetricSet; +import org.apache.iotdb.metrics.utils.MetricInfo; +import org.apache.iotdb.metrics.utils.MetricLevel; +import org.apache.iotdb.metrics.utils.MetricType; + +import java.util.HashMap; +import java.util.Map; + +public class DataExchangeCountMetricSet implements IMetricSet { + + private static final String metric = Metric.DATA_EXCHANGE_COUNT.toString(); + + public static final Map<String, MetricInfo> metricInfoMap = new HashMap<>(); + + public static final String SEND_NEW_DATA_BLOCK_NUM_CALLER = "send_new_data_block_num_caller"; + public static final String SEND_NEW_DATA_BLOCK_NUM_SERVER = "send_new_data_block_num_server"; + public static final String ON_ACKNOWLEDGE_DATA_BLOCK_NUM_CALLER = + "on_acknowledge_data_block_num_caller"; + public static final String ON_ACKNOWLEDGE_DATA_BLOCK_NUM_SERVER = + "on_acknowledge_data_block_num_server"; + public static final String GET_DATA_BLOCK_NUM_CALLER = "get_data_block_num_caller"; + public static final String GET_DATA_BLOCK_NUM_SERVER = "get_data_block_num_server"; + + static { + metricInfoMap.put( + SEND_NEW_DATA_BLOCK_NUM_CALLER, + new MetricInfo( + MetricType.HISTOGRAM, + metric, + Tag.OPERATION.toString(), + "send_new_data_block_num", + Tag.TYPE.toString(), + "caller")); + metricInfoMap.put( + SEND_NEW_DATA_BLOCK_NUM_SERVER, + new MetricInfo( + MetricType.HISTOGRAM, + metric, + Tag.OPERATION.toString(), + "send_new_data_block_num", + Tag.TYPE.toString(), + "server")); + metricInfoMap.put( + ON_ACKNOWLEDGE_DATA_BLOCK_NUM_CALLER, + new MetricInfo( + MetricType.HISTOGRAM, + metric, + Tag.OPERATION.toString(), + "on_acknowledge_data_block_num", + Tag.TYPE.toString(), + "caller")); + metricInfoMap.put( + ON_ACKNOWLEDGE_DATA_BLOCK_NUM_SERVER, + new MetricInfo( + MetricType.HISTOGRAM, + metric, + Tag.OPERATION.toString(), + "on_acknowledge_data_block_num", + Tag.TYPE.toString(), + "server")); + metricInfoMap.put( + GET_DATA_BLOCK_NUM_CALLER, + new MetricInfo( + MetricType.HISTOGRAM, + metric, + Tag.OPERATION.toString(), + "get_data_block_num", + Tag.TYPE.toString(), + "caller")); + metricInfoMap.put( + GET_DATA_BLOCK_NUM_SERVER, + new MetricInfo( + MetricType.HISTOGRAM, + metric, + Tag.OPERATION.toString(), + "get_data_block_num", + Tag.TYPE.toString(), + "server")); + } + + @Override + public void bindTo(AbstractMetricService metricService) { + for (MetricInfo metricInfo : metricInfoMap.values()) { + metricService.getOrCreateHistogram( + metricInfo.getName(), MetricLevel.IMPORTANT, metricInfo.getTagsInArray()); + } + } + + @Override + public void unbindFrom(AbstractMetricService metricService) { + for (MetricInfo metricInfo : metricInfoMap.values()) { + metricService.remove(MetricType.HISTOGRAM, metric, metricInfo.getTagsInArray()); + } + } +} diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/metric/QueryMetricsManager.java b/server/src/main/java/org/apache/iotdb/db/mpp/metric/QueryMetricsManager.java index 206125bb2e..b5bcc5be3f 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/metric/QueryMetricsManager.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/metric/QueryMetricsManager.java @@ -27,8 +27,6 @@ import org.apache.iotdb.metrics.utils.MetricLevel; import java.util.concurrent.TimeUnit; -import static org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet.GET_DATA_BLOCK_NUM; - public class QueryMetricsManager { private final MetricService metricService = MetricService.getInstance(); @@ -88,7 +86,7 @@ public class QueryMetricsManager { } public void recordDataExchangeCost(String stage, long costTimeInNanos) { - MetricInfo metricInfo = DataExchangeMetricSet.metricInfoMap.get(stage); + MetricInfo metricInfo = DataExchangeCostMetricSet.metricInfoMap.get(stage); metricService.timer( costTimeInNanos, TimeUnit.NANOSECONDS, @@ -97,13 +95,10 @@ public class QueryMetricsManager { metricInfo.getTagsInArray()); } - public void recordDataBlockNum(int num) { + public void recordDataBlockNum(String type, int num) { + MetricInfo metricInfo = DataExchangeCountMetricSet.metricInfoMap.get(type); metricService.histogram( - num, - Metric.DATA_EXCHANGE_COUNT.toString(), - MetricLevel.IMPORTANT, - Tag.NAME.toString(), - GET_DATA_BLOCK_NUM); + num, metricInfo.getName(), MetricLevel.IMPORTANT, metricInfo.getTagsInArray()); } public void recordTaskQueueTime(String name, long queueTimeInNanos) { diff --git a/server/src/main/java/org/apache/iotdb/db/service/metrics/DataNodeMetricsHelper.java b/server/src/main/java/org/apache/iotdb/db/service/metrics/DataNodeMetricsHelper.java index 34ded3bb24..074b7345a6 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/metrics/DataNodeMetricsHelper.java +++ b/server/src/main/java/org/apache/iotdb/db/service/metrics/DataNodeMetricsHelper.java @@ -20,7 +20,8 @@ package org.apache.iotdb.db.service.metrics; import org.apache.iotdb.commons.service.metric.MetricService; -import org.apache.iotdb.db.mpp.metric.DataExchangeMetricSet; +import org.apache.iotdb.db.mpp.metric.DataExchangeCostMetricSet; +import org.apache.iotdb.db.mpp.metric.DataExchangeCountMetricSet; import org.apache.iotdb.db.mpp.metric.DriverSchedulerMetricSet; import org.apache.iotdb.db.mpp.metric.QueryExecutionMetricSet; import org.apache.iotdb.db.mpp.metric.QueryPlanCostMetricSet; @@ -43,7 +44,8 @@ public class DataNodeMetricsHelper { MetricService.getInstance().addMetricSet(new SeriesScanCostMetricSet()); MetricService.getInstance().addMetricSet(new QueryExecutionMetricSet()); MetricService.getInstance().addMetricSet(new QueryResourceMetricSet()); - MetricService.getInstance().addMetricSet(new DataExchangeMetricSet()); + MetricService.getInstance().addMetricSet(new DataExchangeCostMetricSet()); + MetricService.getInstance().addMetricSet(new DataExchangeCountMetricSet()); MetricService.getInstance().addMetricSet(new DriverSchedulerMetricSet()); } }
