diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/CheckAndMutate.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/CheckAndMutate.java index 7f8d0677b6be..954d03d18ca6 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/CheckAndMutate.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/CheckAndMutate.java @@ -70,6 +70,7 @@ public static final class Builder { private byte[] value; private Filter filter; private TimeRange timeRange; + private boolean queryMetricsEnabled = false; private Builder(byte[] row) { this.row = Preconditions.checkNotNull(row, "row is null"); @@ -133,6 +134,21 @@ public Builder timeRange(TimeRange timeRange) { return this; } + /** + * Enables the return of {@link QueryMetrics} alongside the corresponding result for this query + *

+ * This is intended for advanced users who need result-granular, server-side metrics + *

+ * Does not work + * @param queryMetricsEnabled {@code true} to enable collection of per-result query metrics + * {@code false} to disable metrics collection (resulting in + * {@code null} metrics) + */ + public Builder queryMetricsEnabled(boolean queryMetricsEnabled) { + this.queryMetricsEnabled = queryMetricsEnabled; + return this; + } + private void preCheck(Row action) { Preconditions.checkNotNull(action, "action is null"); if (!Bytes.equals(row, action.getRow())) { @@ -154,9 +170,10 @@ private void preCheck(Row action) { public CheckAndMutate build(Put put) { preCheck(put); if (filter != null) { - return new CheckAndMutate(row, filter, timeRange, put); + return new CheckAndMutate(row, filter, timeRange, put, queryMetricsEnabled); } else { - return new CheckAndMutate(row, family, qualifier, op, value, timeRange, put); + return new CheckAndMutate(row, family, qualifier, op, value, timeRange, put, + queryMetricsEnabled); } } @@ -168,9 +185,10 @@ public CheckAndMutate build(Put put) { public CheckAndMutate build(Delete delete) { preCheck(delete); if (filter != null) { - return new CheckAndMutate(row, filter, timeRange, delete); + return new CheckAndMutate(row, filter, timeRange, delete, queryMetricsEnabled); } else { - return new CheckAndMutate(row, family, qualifier, op, value, timeRange, delete); + return new CheckAndMutate(row, family, qualifier, op, value, timeRange, delete, + queryMetricsEnabled); } } @@ -182,9 +200,10 @@ public CheckAndMutate build(Delete delete) { public CheckAndMutate build(Increment increment) { preCheck(increment); if (filter != null) { - return new CheckAndMutate(row, filter, timeRange, increment); + return new CheckAndMutate(row, filter, timeRange, increment, queryMetricsEnabled); } else { - return new CheckAndMutate(row, family, qualifier, op, value, timeRange, increment); + return new CheckAndMutate(row, family, qualifier, op, value, timeRange, increment, + queryMetricsEnabled); } } @@ -196,9 +215,10 @@ public CheckAndMutate build(Increment increment) { public CheckAndMutate build(Append append) { preCheck(append); if (filter != null) { - return new CheckAndMutate(row, filter, timeRange, append); + return new CheckAndMutate(row, filter, timeRange, append, queryMetricsEnabled); } else { - return new CheckAndMutate(row, family, qualifier, op, value, timeRange, append); + return new CheckAndMutate(row, family, qualifier, op, value, timeRange, append, + queryMetricsEnabled); } } @@ -210,9 +230,10 @@ public CheckAndMutate build(Append append) { public CheckAndMutate build(RowMutations mutations) { preCheck(mutations); if (filter != null) { - return new CheckAndMutate(row, filter, timeRange, mutations); + return new CheckAndMutate(row, filter, timeRange, mutations, queryMetricsEnabled); } else { - return new CheckAndMutate(row, family, qualifier, op, value, timeRange, mutations); + return new CheckAndMutate(row, family, qualifier, op, value, timeRange, mutations, + queryMetricsEnabled); } } } @@ -234,9 +255,10 @@ public static Builder newBuilder(byte[] row) { private final Filter filter; private final TimeRange timeRange; private final Row action; + private final boolean queryMetricsEnabled; private CheckAndMutate(byte[] row, byte[] family, byte[] qualifier, final CompareOperator op, - byte[] value, TimeRange timeRange, Row action) { + byte[] value, TimeRange timeRange, Row action, boolean queryMetricsEnabled) { this.row = row; this.family = family; this.qualifier = qualifier; @@ -245,9 +267,11 @@ private CheckAndMutate(byte[] row, byte[] family, byte[] qualifier, final Compar this.filter = null; this.timeRange = timeRange != null ? timeRange : TimeRange.allTime(); this.action = action; + this.queryMetricsEnabled = queryMetricsEnabled; } - private CheckAndMutate(byte[] row, Filter filter, TimeRange timeRange, Row action) { + private CheckAndMutate(byte[] row, Filter filter, TimeRange timeRange, Row action, + boolean queryMetricsEnabled) { this.row = row; this.family = null; this.qualifier = null; @@ -256,6 +280,7 @@ private CheckAndMutate(byte[] row, Filter filter, TimeRange timeRange, Row actio this.filter = filter; this.timeRange = timeRange != null ? timeRange : TimeRange.allTime(); this.action = action; + this.queryMetricsEnabled = queryMetricsEnabled; } /** Returns the row */ @@ -326,4 +351,9 @@ public TimeRange getTimeRange() { public Row getAction() { return action; } + + /** Returns whether query metrics are enabled */ + public boolean isQueryMetricsEnabled() { + return queryMetricsEnabled; + } } diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/CheckAndMutateResult.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/CheckAndMutateResult.java index 8ecb49e3d5f1..6ed6f7e26a9b 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/CheckAndMutateResult.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/CheckAndMutateResult.java @@ -27,6 +27,8 @@ public class CheckAndMutateResult { private final boolean success; private final Result result; + private QueryMetrics metrics = null; + public CheckAndMutateResult(boolean success, Result result) { this.success = success; this.result = result; @@ -41,4 +43,13 @@ public boolean isSuccess() { public Result getResult() { return result; } + + public CheckAndMutateResult setMetrics(QueryMetrics metrics) { + this.metrics = metrics; + return this; + } + + public QueryMetrics getMetrics() { + return metrics; + } } diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Get.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Get.java index d11045f89b03..da4754227999 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Get.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Get.java @@ -96,6 +96,7 @@ public Get(Get get) { this.setFilter(get.getFilter()); this.setReplicaId(get.getReplicaId()); this.setConsistency(get.getConsistency()); + this.setQueryMetricsEnabled(get.isQueryMetricsEnabled()); // from Get this.cacheBlocks = get.getCacheBlocks(); this.maxVersions = get.getMaxVersions(); @@ -511,6 +512,7 @@ public Map toMap(int maxCols) { map.put("colFamTimeRangeMap", colFamTimeRangeMapStr); } map.put("priority", getPriority()); + map.put("queryMetricsEnabled", queryMetricsEnabled); return map; } diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/HTable.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/HTable.java index fd3de615cb2f..8767257bf5b7 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/HTable.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/HTable.java @@ -744,10 +744,8 @@ public boolean checkAndPut(final byte[] row, final byte[] family, final byte[] q .setTableName(tableName).setOperation(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE) .setContainerOperations(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE, HBaseSemanticAttributes.Operation.PUT); - return TraceUtil.trace( - () -> doCheckAndMutate(row, family, qualifier, CompareOperator.EQUAL, value, null, null, put) - .isSuccess(), - supplier); + return TraceUtil.trace(() -> doCheckAndMutate(row, family, qualifier, CompareOperator.EQUAL, + value, null, null, put, false).isSuccess(), supplier); } @Override @@ -759,7 +757,7 @@ public boolean checkAndPut(final byte[] row, final byte[] family, final byte[] q .setContainerOperations(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE, HBaseSemanticAttributes.Operation.PUT); return TraceUtil.trace(() -> doCheckAndMutate(row, family, qualifier, - toCompareOperator(compareOp), value, null, null, put).isSuccess(), supplier); + toCompareOperator(compareOp), value, null, null, put, false).isSuccess(), supplier); } @Override @@ -771,7 +769,7 @@ public boolean checkAndPut(final byte[] row, final byte[] family, final byte[] q .setContainerOperations(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE, HBaseSemanticAttributes.Operation.PUT); return TraceUtil.trace( - () -> doCheckAndMutate(row, family, qualifier, op, value, null, null, put).isSuccess(), + () -> doCheckAndMutate(row, family, qualifier, op, value, null, null, put, false).isSuccess(), supplier); } @@ -784,7 +782,7 @@ public boolean checkAndDelete(final byte[] row, final byte[] family, final byte[ .setContainerOperations(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE, HBaseSemanticAttributes.Operation.DELETE); return TraceUtil.trace(() -> doCheckAndMutate(row, family, qualifier, CompareOperator.EQUAL, - value, null, null, delete).isSuccess(), supplier); + value, null, null, delete, false).isSuccess(), supplier); } @Override @@ -796,7 +794,7 @@ public boolean checkAndDelete(final byte[] row, final byte[] family, final byte[ .setContainerOperations(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE, HBaseSemanticAttributes.Operation.DELETE); return TraceUtil.trace(() -> doCheckAndMutate(row, family, qualifier, - toCompareOperator(compareOp), value, null, null, delete).isSuccess(), supplier); + toCompareOperator(compareOp), value, null, null, delete, false).isSuccess(), supplier); } @Override @@ -807,9 +805,9 @@ public boolean checkAndDelete(final byte[] row, final byte[] family, final byte[ .setTableName(tableName).setOperation(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE) .setContainerOperations(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE, HBaseSemanticAttributes.Operation.DELETE); - return TraceUtil.trace( - () -> doCheckAndMutate(row, family, qualifier, op, value, null, null, delete).isSuccess(), - supplier); + return TraceUtil + .trace(() -> doCheckAndMutate(row, family, qualifier, op, value, null, null, delete, false) + .isSuccess(), supplier); } @Override @@ -826,7 +824,8 @@ public CheckAndMutateWithFilterBuilder checkAndMutate(byte[] row, Filter filter) private CheckAndMutateResult doCheckAndMutate(final byte[] row, final byte[] family, final byte[] qualifier, final CompareOperator op, final byte[] value, final Filter filter, - final TimeRange timeRange, final RowMutations rm) throws IOException { + final TimeRange timeRange, final RowMutations rm, boolean queryMetricsEnabled) + throws IOException { long nonceGroup = getNonceGroup(); long nonce = getNonce(); CancellableRegionServerCallable callable = @@ -835,9 +834,9 @@ private CheckAndMutateResult doCheckAndMutate(final byte[] row, final byte[] fam rm.getMaxPriority(), requestAttributes) { @Override protected MultiResponse rpcCall() throws Exception { - MultiRequest request = - RequestConverter.buildMultiRequest(getLocation().getRegionInfo().getRegionName(), row, - family, qualifier, op, value, filter, timeRange, rm, nonceGroup, nonce); + MultiRequest request = RequestConverter.buildMultiRequest( + getLocation().getRegionInfo().getRegionName(), row, family, qualifier, op, value, + filter, timeRange, rm, nonceGroup, nonce, queryMetricsEnabled); ClientProtos.MultiResponse response = doMulti(request); ClientProtos.RegionActionResult res = response.getRegionActionResultList().get(0); if (res.hasException()) { @@ -880,7 +879,7 @@ public boolean checkAndMutate(final byte[] row, final byte[] family, final byte[ .setTableName(tableName).setOperation(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE) .setContainerOperations(rm); return TraceUtil.trace(() -> doCheckAndMutate(row, family, qualifier, - toCompareOperator(compareOp), value, null, null, rm).isSuccess(), supplier); + toCompareOperator(compareOp), value, null, null, rm, false).isSuccess(), supplier); } @Override @@ -891,7 +890,7 @@ public boolean checkAndMutate(final byte[] row, final byte[] family, final byte[ .setTableName(tableName).setOperation(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE) .setContainerOperations(rm); return TraceUtil.trace( - () -> doCheckAndMutate(row, family, qualifier, op, value, null, null, rm).isSuccess(), + () -> doCheckAndMutate(row, family, qualifier, op, value, null, null, rm, false).isSuccess(), supplier); } @@ -910,18 +909,21 @@ public CheckAndMutateResult checkAndMutate(CheckAndMutate checkAndMutate) throws } return doCheckAndMutate(checkAndMutate.getRow(), checkAndMutate.getFamily(), checkAndMutate.getQualifier(), checkAndMutate.getCompareOp(), checkAndMutate.getValue(), - checkAndMutate.getFilter(), checkAndMutate.getTimeRange(), (Mutation) action); + checkAndMutate.getFilter(), checkAndMutate.getTimeRange(), (Mutation) action, + checkAndMutate.isQueryMetricsEnabled()); } else { return doCheckAndMutate(checkAndMutate.getRow(), checkAndMutate.getFamily(), checkAndMutate.getQualifier(), checkAndMutate.getCompareOp(), checkAndMutate.getValue(), - checkAndMutate.getFilter(), checkAndMutate.getTimeRange(), (RowMutations) action); + checkAndMutate.getFilter(), checkAndMutate.getTimeRange(), (RowMutations) action, + checkAndMutate.isQueryMetricsEnabled()); } }, supplier); } private CheckAndMutateResult doCheckAndMutate(final byte[] row, final byte[] family, final byte[] qualifier, final CompareOperator op, final byte[] value, final Filter filter, - final TimeRange timeRange, final Mutation mutation) throws IOException { + final TimeRange timeRange, final Mutation mutation, boolean queryMetricsEnabled) + throws IOException { long nonceGroup = getNonceGroup(); long nonce = getNonce(); ClientServiceCallable callable = @@ -929,9 +931,9 @@ private CheckAndMutateResult doCheckAndMutate(final byte[] row, final byte[] fam this.rpcControllerFactory.newController(), mutation.getPriority(), requestAttributes) { @Override protected CheckAndMutateResult rpcCall() throws Exception { - MutateRequest request = - RequestConverter.buildMutateRequest(getLocation().getRegionInfo().getRegionName(), row, - family, qualifier, op, value, filter, timeRange, mutation, nonceGroup, nonce); + MutateRequest request = RequestConverter.buildMutateRequest( + getLocation().getRegionInfo().getRegionName(), row, family, qualifier, op, value, + filter, timeRange, mutation, nonceGroup, nonce, queryMetricsEnabled); MutateResponse response = doMutate(request); if (response.hasResult()) { return new CheckAndMutateResult(response.getProcessed(), @@ -1419,7 +1421,7 @@ public boolean thenPut(Put put) throws IOException { return TraceUtil.trace(() -> { validatePut(put); preCheck(); - return doCheckAndMutate(row, family, qualifier, op, value, null, timeRange, put) + return doCheckAndMutate(row, family, qualifier, op, value, null, timeRange, put, false) .isSuccess(); }, supplier); } @@ -1430,7 +1432,7 @@ public boolean thenDelete(Delete delete) throws IOException { .setTableName(tableName).setOperation(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE); return TraceUtil.trace(() -> { preCheck(); - return doCheckAndMutate(row, family, qualifier, op, value, null, timeRange, delete) + return doCheckAndMutate(row, family, qualifier, op, value, null, timeRange, delete, false) .isSuccess(); }, supplier); } @@ -1441,7 +1443,7 @@ public boolean thenMutate(RowMutations mutation) throws IOException { .setTableName(tableName).setOperation(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE); return TraceUtil.trace(() -> { preCheck(); - return doCheckAndMutate(row, family, qualifier, op, value, null, timeRange, mutation) + return doCheckAndMutate(row, family, qualifier, op, value, null, timeRange, mutation, false) .isSuccess(); }, supplier); } @@ -1475,7 +1477,8 @@ public boolean thenPut(Put put) throws IOException { .setTableName(tableName).setOperation(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE); return TraceUtil.trace(() -> { validatePut(put); - return doCheckAndMutate(row, null, null, null, null, filter, timeRange, put).isSuccess(); + return doCheckAndMutate(row, null, null, null, null, filter, timeRange, put, false) + .isSuccess(); }, supplier); } @@ -1483,18 +1486,19 @@ public boolean thenPut(Put put) throws IOException { public boolean thenDelete(Delete delete) throws IOException { final Supplier supplier = new TableOperationSpanBuilder(connection) .setTableName(tableName).setOperation(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE); - return TraceUtil.trace( - () -> doCheckAndMutate(row, null, null, null, null, filter, timeRange, delete).isSuccess(), - supplier); + return TraceUtil + .trace(() -> doCheckAndMutate(row, null, null, null, null, filter, timeRange, delete, false) + .isSuccess(), supplier); } @Override public boolean thenMutate(RowMutations mutation) throws IOException { final Supplier supplier = new TableOperationSpanBuilder(connection) .setTableName(tableName).setOperation(HBaseSemanticAttributes.Operation.CHECK_AND_MUTATE); - return TraceUtil - .trace(() -> doCheckAndMutate(row, null, null, null, null, filter, timeRange, mutation) - .isSuccess(), supplier); + return TraceUtil.trace( + () -> doCheckAndMutate(row, null, null, null, null, filter, timeRange, mutation, false) + .isSuccess(), + supplier); } } } diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Query.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Query.java index 5f129ef9cffe..f31316f01abf 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Query.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Query.java @@ -17,6 +17,7 @@ */ package org.apache.hadoop.hbase.client; +import java.util.List; import java.util.Map; import org.apache.hadoop.hbase.exceptions.DeserializationException; import org.apache.hadoop.hbase.filter.Filter; @@ -46,6 +47,7 @@ public abstract class Query extends OperationWithAttributes { protected Consistency consistency = Consistency.STRONG; protected Map colFamTimeRangeMap = Maps.newTreeMap(Bytes.BYTES_COMPARATOR); protected Boolean loadColumnFamiliesOnDemand = null; + protected boolean queryMetricsEnabled = false; public Filter getFilter() { return filter; @@ -157,6 +159,28 @@ public Query setIsolationLevel(IsolationLevel level) { return this; } + /** + * Enables the return of {@link QueryMetrics} alongside the corresponding result(s) for this query + *

+ * This is intended for advanced users who need result-granular, server-side metrics + *

+ * Does not work with calls to {@link Table#exists(Get)} and {@link Table#exists(List)} + * @param enabled {@code true} to enable collection of per-result query metrics {@code false} to + * disable metrics collection (resulting in {@code null} metrics) + */ + public Query setQueryMetricsEnabled(boolean enabled) { + this.queryMetricsEnabled = enabled; + return this; + } + + /** + * Returns whether query metrics are enabled + * @return {@code true} if query metrics are enabled, {@code false} otherwise + */ + public boolean isQueryMetricsEnabled() { + return queryMetricsEnabled; + } + /** * Returns The isolation level of this query. If no isolation level was set for this query object, * then it returns READ_COMMITTED. diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/QueryMetrics.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/QueryMetrics.java new file mode 100644 index 000000000000..9243b43bb7eb --- /dev/null +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/QueryMetrics.java @@ -0,0 +1,36 @@ +/* + * 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.hadoop.hbase.client; + +import org.apache.yetus.audience.InterfaceAudience; +import org.apache.yetus.audience.InterfaceStability; + +@InterfaceAudience.Public +@InterfaceStability.Evolving +public class QueryMetrics { + private final long blockBytesScanned; + + public QueryMetrics(long blockBytesScanned) { + this.blockBytesScanned = blockBytesScanned; + } + + @InterfaceStability.Evolving + public long getBlockBytesScanned() { + return blockBytesScanned; + } +} diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/RawAsyncTableImpl.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/RawAsyncTableImpl.java index 6f22b35664ba..db8317b6b4d0 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/RawAsyncTableImpl.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/RawAsyncTableImpl.java @@ -396,7 +396,7 @@ public CompletableFuture thenPut(Put put) { () -> RawAsyncTableImpl.this. newCaller(row, put.getPriority(), rpcTimeoutNs) .action((controller, loc, stub) -> RawAsyncTableImpl.mutate(controller, loc, stub, put, (rn, p) -> RequestConverter.buildMutateRequest(rn, row, family, qualifier, op, value, - null, timeRange, p, HConstants.NO_NONCE, HConstants.NO_NONCE), + null, timeRange, p, HConstants.NO_NONCE, HConstants.NO_NONCE, false), (c, r) -> r.getProcessed())) .call(), supplier); @@ -412,7 +412,7 @@ public CompletableFuture thenDelete(Delete delete) { () -> RawAsyncTableImpl.this. newCaller(row, delete.getPriority(), rpcTimeoutNs) .action((controller, loc, stub) -> RawAsyncTableImpl.mutate(controller, loc, stub, delete, (rn, d) -> RequestConverter.buildMutateRequest(rn, row, family, qualifier, op, value, - null, timeRange, d, HConstants.NO_NONCE, HConstants.NO_NONCE), + null, timeRange, d, HConstants.NO_NONCE, HConstants.NO_NONCE, false), (c, r) -> r.getProcessed())) .call(), supplier); @@ -430,7 +430,7 @@ public CompletableFuture thenMutate(RowMutations mutations) { .action((controller, loc, stub) -> RawAsyncTableImpl.this.mutateRow(controller, loc, stub, mutations, (rn, rm) -> RequestConverter.buildMultiRequest(rn, row, family, qualifier, op, value, - null, timeRange, rm, HConstants.NO_NONCE, HConstants.NO_NONCE), + null, timeRange, rm, HConstants.NO_NONCE, HConstants.NO_NONCE, false), CheckAndMutateResult::isSuccess)) .call(), supplier); } @@ -471,7 +471,7 @@ public CompletableFuture thenPut(Put put) { () -> RawAsyncTableImpl.this. newCaller(row, put.getPriority(), rpcTimeoutNs) .action((controller, loc, stub) -> RawAsyncTableImpl.mutate(controller, loc, stub, put, (rn, p) -> RequestConverter.buildMutateRequest(rn, row, null, null, null, null, filter, - timeRange, p, HConstants.NO_NONCE, HConstants.NO_NONCE), + timeRange, p, HConstants.NO_NONCE, HConstants.NO_NONCE, false), (c, r) -> r.getProcessed())) .call(), supplier); @@ -486,7 +486,7 @@ public CompletableFuture thenDelete(Delete delete) { () -> RawAsyncTableImpl.this. newCaller(row, delete.getPriority(), rpcTimeoutNs) .action((controller, loc, stub) -> RawAsyncTableImpl.mutate(controller, loc, stub, delete, (rn, d) -> RequestConverter.buildMutateRequest(rn, row, null, null, null, null, filter, - timeRange, d, HConstants.NO_NONCE, HConstants.NO_NONCE), + timeRange, d, HConstants.NO_NONCE, HConstants.NO_NONCE, false), (c, r) -> r.getProcessed())) .call(), supplier); @@ -503,7 +503,7 @@ public CompletableFuture thenMutate(RowMutations mutations) { .action((controller, loc, stub) -> RawAsyncTableImpl.this.mutateRow(controller, loc, stub, mutations, (rn, rm) -> RequestConverter.buildMultiRequest(rn, row, null, null, null, null, filter, - timeRange, rm, HConstants.NO_NONCE, HConstants.NO_NONCE), + timeRange, rm, HConstants.NO_NONCE, HConstants.NO_NONCE, false), CheckAndMutateResult::isSuccess)) .call(), supplier); } @@ -538,7 +538,8 @@ public CompletableFuture checkAndMutate(CheckAndMutate che (rn, m) -> RequestConverter.buildMutateRequest(rn, checkAndMutate.getRow(), checkAndMutate.getFamily(), checkAndMutate.getQualifier(), checkAndMutate.getCompareOp(), checkAndMutate.getValue(), - checkAndMutate.getFilter(), checkAndMutate.getTimeRange(), m, nonceGroup, nonce), + checkAndMutate.getFilter(), checkAndMutate.getTimeRange(), m, nonceGroup, nonce, + checkAndMutate.isQueryMetricsEnabled()), (c, r) -> ResponseConverter.getCheckAndMutateResult(r, c.cellScanner()))) .call(); } else if (checkAndMutate.getAction() instanceof RowMutations) { @@ -554,7 +555,8 @@ CheckAndMutateResult> mutateRow(controller, loc, stub, rowMutations, (rn, rm) -> RequestConverter.buildMultiRequest(rn, checkAndMutate.getRow(), checkAndMutate.getFamily(), checkAndMutate.getQualifier(), checkAndMutate.getCompareOp(), checkAndMutate.getValue(), - checkAndMutate.getFilter(), checkAndMutate.getTimeRange(), rm, nonceGroup, nonce), + checkAndMutate.getFilter(), checkAndMutate.getTimeRange(), rm, nonceGroup, nonce, + checkAndMutate.isQueryMetricsEnabled()), resp -> resp)) .call(); } else { diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Result.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Result.java index 6915adec0181..faaff681faa9 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Result.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Result.java @@ -97,6 +97,7 @@ public class Result implements CellScannable, CellScanner { */ private int cellScannerIndex = INITIAL_CELLSCANNER_INDEX; private RegionLoadStats stats; + private QueryMetrics metrics = null; private final boolean readonly; @@ -903,6 +904,11 @@ public void setStatistics(RegionLoadStats loadStats) { this.stats = loadStats; } + @InterfaceAudience.Private + public void setMetrics(QueryMetrics metrics) { + this.metrics = metrics; + } + /** * Returns the associated statistics about the region from which this was returned. Can be * null if stats are disabled. @@ -911,6 +917,11 @@ public RegionLoadStats getStats() { return stats; } + /** Returns the query metrics, or {@code null} if we do not enable metrics. */ + public QueryMetrics getMetrics() { + return metrics; + } + /** * All methods modifying state of Result object must call this method to ensure that special * purpose immutable Results can't be accidentally modified. diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Scan.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Scan.java index a96a8a743a09..81186ce656e4 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Scan.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/client/Scan.java @@ -284,6 +284,7 @@ public Scan(Scan scan) throws IOException { setPriority(scan.getPriority()); readType = scan.getReadType(); super.setReplicaId(scan.getReplicaId()); + super.setQueryMetricsEnabled(scan.isQueryMetricsEnabled()); } /** @@ -316,6 +317,7 @@ public Scan(Get get) { this.mvccReadPoint = -1L; setPriority(get.getPriority()); super.setReplicaId(get.getReplicaId()); + super.setQueryMetricsEnabled(get.isQueryMetricsEnabled()); } public boolean isGetScan() { @@ -983,6 +985,7 @@ public Map toMap(int maxCols) { map.put("colFamTimeRangeMap", colFamTimeRangeMapStr); } map.put("priority", getPriority()); + map.put("queryMetricsEnabled", queryMetricsEnabled); return map; } diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/shaded/protobuf/ProtobufUtil.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/shaded/protobuf/ProtobufUtil.java index 46eb86aeb336..a1e28aef8310 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/shaded/protobuf/ProtobufUtil.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/shaded/protobuf/ProtobufUtil.java @@ -88,6 +88,7 @@ import org.apache.hadoop.hbase.client.OnlineLogRecord; import org.apache.hadoop.hbase.client.PackagePrivateFieldAccessor; import org.apache.hadoop.hbase.client.Put; +import org.apache.hadoop.hbase.client.QueryMetrics; import org.apache.hadoop.hbase.client.RegionInfoBuilder; import org.apache.hadoop.hbase.client.RegionLoadStats; import org.apache.hadoop.hbase.client.RegionReplicaUtil; @@ -611,6 +612,7 @@ public static Get toGet(final ClientProtos.Get proto) throws IOException { if (proto.hasLoadColumnFamiliesOnDemand()) { get.setLoadColumnFamiliesOnDemand(proto.getLoadColumnFamiliesOnDemand()); } + get.setQueryMetricsEnabled(proto.getQueryMetricsEnabled()); return get; } @@ -1051,6 +1053,7 @@ public static ClientProtos.Scan toScan(final Scan scan) throws IOException { if (scan.isNeedCursorResult()) { scanBuilder.setNeedCursorResult(true); } + scanBuilder.setQueryMetricsEnabled(scan.isQueryMetricsEnabled()); return scanBuilder.build(); } @@ -1160,6 +1163,7 @@ public static Scan toScan(final ClientProtos.Scan proto) throws IOException { if (proto.getNeedCursorResult()) { scan.setNeedCursorResult(true); } + scan.setQueryMetricsEnabled(proto.getQueryMetricsEnabled()); return scan; } @@ -1239,6 +1243,7 @@ public static ClientProtos.Get toGet(final Get get) throws IOException { if (loadColumnFamiliesOnDemand != null) { builder.setLoadColumnFamiliesOnDemand(loadColumnFamiliesOnDemand); } + builder.setQueryMetricsEnabled(get.isQueryMetricsEnabled()); return builder.build(); } @@ -1393,6 +1398,10 @@ public static ClientProtos.Result toResult(final Result result, boolean encodeTa builder.setStale(result.isStale()); builder.setPartial(result.mayHaveMoreCellsInRow()); + if (result.getMetrics() != null) { + builder.setMetrics(toQueryMetrics(result.getMetrics())); + } + return builder.build(); } @@ -1422,6 +1431,9 @@ public static ClientProtos.Result toResultNoData(final Result result) { ClientProtos.Result.Builder builder = ClientProtos.Result.newBuilder(); builder.setAssociatedCellCount(size); builder.setStale(result.isStale()); + if (result.getMetrics() != null) { + builder.setMetrics(toQueryMetrics(result.getMetrics())); + } return builder.build(); } @@ -1462,7 +1474,11 @@ public static Result toResult(final ClientProtos.Result proto, boolean decodeTag for (CellProtos.Cell c : values) { cells.add(toCell(builder, c, decodeTags)); } - return Result.create(cells, null, proto.getStale(), proto.getPartial()); + Result r = Result.create(cells, null, proto.getStale(), proto.getPartial()); + if (proto.hasMetrics()) { + r.setMetrics(toQueryMetrics(proto.getMetrics())); + } + return r; } /** @@ -1507,9 +1523,15 @@ public static Result toResult(final ClientProtos.Result proto, final CellScanner } } - return (cells == null || cells.isEmpty()) + Result r = (cells == null || cells.isEmpty()) ? (proto.getStale() ? EMPTY_RESULT_STALE : EMPTY_RESULT) : Result.create(cells, null, proto.getStale()); + + if (proto.hasMetrics()) { + r.setMetrics(toQueryMetrics(proto.getMetrics())); + } + + return r; } /** @@ -3509,6 +3531,7 @@ public static CheckAndMutate toCheckAndMutate(ClientProtos.Condition condition, ? ProtobufUtil.toTimeRange(condition.getTimeRange()) : TimeRange.allTime(); builder.timeRange(timeRange); + builder.queryMetricsEnabled(condition.getQueryMetricsEnabled()); try { MutationType type = mutation.getMutateType(); @@ -3546,6 +3569,7 @@ public static CheckAndMutate toCheckAndMutate(ClientProtos.Condition condition, ? ProtobufUtil.toTimeRange(condition.getTimeRange()) : TimeRange.allTime(); builder.timeRange(timeRange); + builder.queryMetricsEnabled(condition.getQueryMetricsEnabled()); try { if (mutations.size() == 1) { @@ -3572,7 +3596,7 @@ public static CheckAndMutate toCheckAndMutate(ClientProtos.Condition condition, public static ClientProtos.Condition toCondition(final byte[] row, final byte[] family, final byte[] qualifier, final CompareOperator op, final byte[] value, final Filter filter, - final TimeRange timeRange) throws IOException { + final TimeRange timeRange, boolean queryMetricsEnabled) throws IOException { ClientProtos.Condition.Builder builder = ClientProtos.Condition.newBuilder().setRow(UnsafeByteOperations.unsafeWrap(row)); @@ -3587,18 +3611,19 @@ public static ClientProtos.Condition toCondition(final byte[] row, final byte[] .setCompareType(HBaseProtos.CompareType.valueOf(op.name())); } + builder.setQueryMetricsEnabled(queryMetricsEnabled); return builder.setTimeRange(ProtobufUtil.toTimeRange(timeRange)).build(); } public static ClientProtos.Condition toCondition(final byte[] row, final Filter filter, final TimeRange timeRange) throws IOException { - return toCondition(row, null, null, null, null, filter, timeRange); + return toCondition(row, null, null, null, null, filter, timeRange, false); } public static ClientProtos.Condition toCondition(final byte[] row, final byte[] family, final byte[] qualifier, final CompareOperator op, final byte[] value, final TimeRange timeRange) throws IOException { - return toCondition(row, family, qualifier, op, value, null, timeRange); + return toCondition(row, family, qualifier, op, value, null, timeRange, false); } public static List toBalancerDecisionResponse(HBaseProtos.LogEntry logEntry) { @@ -3717,6 +3742,15 @@ public static ClusterStatusProtos.ServerTask toServerTask(ServerTask task) { .setStartTime(task.getStartTime()).setCompletionTime(task.getCompletionTime()).build(); } + public static ClientProtos.QueryMetrics toQueryMetrics(QueryMetrics metrics) { + return ClientProtos.QueryMetrics.newBuilder() + .setBlockBytesScanned(metrics.getBlockBytesScanned()).build(); + } + + public static QueryMetrics toQueryMetrics(ClientProtos.QueryMetrics metrics) { + return new QueryMetrics(metrics.getBlockBytesScanned()); + } + /** * Check whether this IPBE indicates EOF or not. *

diff --git a/hbase-client/src/main/java/org/apache/hadoop/hbase/shaded/protobuf/RequestConverter.java b/hbase-client/src/main/java/org/apache/hadoop/hbase/shaded/protobuf/RequestConverter.java index 6ab4817b172f..63bb5ba9f389 100644 --- a/hbase-client/src/main/java/org/apache/hadoop/hbase/shaded/protobuf/RequestConverter.java +++ b/hbase-client/src/main/java/org/apache/hadoop/hbase/shaded/protobuf/RequestConverter.java @@ -234,7 +234,7 @@ public static MutateRequest buildIncrementRequest(final byte[] regionName, final public static MutateRequest buildMutateRequest(final byte[] regionName, final byte[] row, final byte[] family, final byte[] qualifier, final CompareOperator op, final byte[] value, final Filter filter, final TimeRange timeRange, final Mutation mutation, long nonceGroup, - long nonce) throws IOException { + long nonce, boolean queryMetricsEnabled) throws IOException { MutateRequest.Builder builder = MutateRequest.newBuilder(); if (mutation instanceof Increment || mutation instanceof Append) { builder.setMutation(ProtobufUtil.toMutation(getMutationType(mutation), mutation, nonce)) @@ -243,7 +243,8 @@ public static MutateRequest buildMutateRequest(final byte[] regionName, final by builder.setMutation(ProtobufUtil.toMutation(getMutationType(mutation), mutation)); } return builder.setRegion(buildRegionSpecifier(RegionSpecifierType.REGION_NAME, regionName)) - .setCondition(ProtobufUtil.toCondition(row, family, qualifier, op, value, filter, timeRange)) + .setCondition(ProtobufUtil.toCondition(row, family, qualifier, op, value, filter, timeRange, + queryMetricsEnabled)) .build(); } @@ -254,10 +255,10 @@ public static MutateRequest buildMutateRequest(final byte[] regionName, final by public static ClientProtos.MultiRequest buildMultiRequest(final byte[] regionName, final byte[] row, final byte[] family, final byte[] qualifier, final CompareOperator op, final byte[] value, final Filter filter, final TimeRange timeRange, - final RowMutations rowMutations, long nonceGroup, long nonce) throws IOException { - return buildMultiRequest(regionName, rowMutations, - ProtobufUtil.toCondition(row, family, qualifier, op, value, filter, timeRange), nonceGroup, - nonce); + final RowMutations rowMutations, long nonceGroup, long nonce, boolean queryMetricsEnabled) + throws IOException { + return buildMultiRequest(regionName, rowMutations, ProtobufUtil.toCondition(row, family, + qualifier, op, value, filter, timeRange, queryMetricsEnabled), nonceGroup, nonce); } /** @@ -594,9 +595,9 @@ public static void buildRegionActions(final byte[] regionName, final List 0) { // Get the result of the Increment/Append operations from the first element of the // ResultOrException list @@ -222,8 +224,12 @@ private static CheckAndMutateResult getCheckAndMutateResult(RegionActionResult a result = r; } } + + if (roe.hasMetrics()) { + metrics = ProtobufUtil.toQueryMetrics(roe.getMetrics()); + } } - return new CheckAndMutateResult(actionResult.getProcessed(), result); + return new CheckAndMutateResult(actionResult.getProcessed(), result).setMetrics(metrics); } private static Result getMutateRowResult(RegionActionResult actionResult, CellScanner cells) @@ -258,10 +264,16 @@ public static CheckAndMutateResult getCheckAndMutateResult( ClientProtos.MutateResponse mutateResponse, CellScanner cells) throws IOException { boolean success = mutateResponse.getProcessed(); Result result = null; + QueryMetrics metrics = null; if (mutateResponse.hasResult()) { result = ProtobufUtil.toResult(mutateResponse.getResult(), cells); } - return new CheckAndMutateResult(success, result); + + if (mutateResponse.hasMetrics()) { + metrics = ProtobufUtil.toQueryMetrics(mutateResponse.getMetrics()); + } + + return new CheckAndMutateResult(success, result).setMetrics(metrics); } /** @@ -436,6 +448,7 @@ public static Result[] getResults(CellScanner cellScanner, ScanResponse response int noOfResults = cellScanner != null ? response.getCellsPerResultCount() : response.getResultsCount(); Result[] results = new Result[noOfResults]; + List queryMetrics = response.getQueryMetricsList(); for (int i = 0; i < noOfResults; i++) { if (cellScanner != null) { // Cells are out in cellblocks. Group them up again as Results. How many to read at a @@ -472,6 +485,12 @@ public static Result[] getResults(CellScanner cellScanner, ScanResponse response // Result is pure pb. results[i] = ProtobufUtil.toResult(response.getResults(i)); } + + // Populate result metrics if they exist + if (queryMetrics.size() > i) { + QueryMetrics metrics = ProtobufUtil.toQueryMetrics(queryMetrics.get(i)); + results[i].setMetrics(metrics); + } } return results; } diff --git a/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestGet.java b/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestGet.java index 51b22f5366e7..3815951d4e33 100644 --- a/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestGet.java +++ b/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestGet.java @@ -182,6 +182,7 @@ public void TestGetRowFromGetCopyConstructor() throws Exception { get.setMaxResultsPerColumnFamily(10); get.setRowOffsetPerColumnFamily(11); get.setCacheBlocks(true); + get.setQueryMetricsEnabled(true); Get copyGet = new Get(get); assertEquals(0, Bytes.compareTo(get.getRow(), copyGet.getRow())); @@ -196,6 +197,7 @@ public void TestGetRowFromGetCopyConstructor() throws Exception { assertEquals(get.getConsistency(), copyGet.getConsistency()); assertEquals(get.getReplicaId(), copyGet.getReplicaId()); assertEquals(get.getIsolationLevel(), copyGet.getIsolationLevel()); + assertTrue(get.isQueryMetricsEnabled()); // from Get class assertEquals(get.isCheckExistenceOnly(), copyGet.isCheckExistenceOnly()); diff --git a/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestOnlineLogRecord.java b/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestOnlineLogRecord.java index 72013b6f294b..4446ab0f7261 100644 --- a/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestOnlineLogRecord.java +++ b/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestOnlineLogRecord.java @@ -44,21 +44,21 @@ public void itSerializesScan() { Scan scan = new Scan(); scan.withStartRow(Bytes.toBytes(123)); scan.withStopRow(Bytes.toBytes(456)); - String expectedOutput = - "{\n" + " \"startTime\": 1,\n" + " \"processingTime\": 2,\n" + " \"queueTime\": 3,\n" - + " \"responseSize\": 4,\n" + " \"blockBytesScanned\": 5,\n" + " \"fsReadTime\": 6,\n" - + " \"multiGetsCount\": 6,\n" + " \"multiMutationsCount\": 7,\n" + " \"scan\": {\n" - + " \"startRow\": \"\\\\x00\\\\x00\\\\x00{\",\n" + " \"targetReplicaId\": -1,\n" - + " \"batch\": -1,\n" + " \"totalColumns\": 0,\n" + " \"maxResultSize\": -1,\n" - + " \"families\": {},\n" + " \"priority\": -1,\n" + " \"caching\": -1,\n" - + " \"includeStopRow\": false,\n" + " \"consistency\": \"STRONG\",\n" - + " \"maxVersions\": 1,\n" + " \"storeOffset\": 0,\n" + " \"mvccReadPoint\": -1,\n" - + " \"includeStartRow\": true,\n" + " \"needCursorResult\": false,\n" - + " \"stopRow\": \"\\\\x00\\\\x00\\\\x01\\\\xC8\",\n" + " \"storeLimit\": -1,\n" - + " \"limit\": -1,\n" + " \"cacheBlocks\": true,\n" - + " \"readType\": \"DEFAULT\",\n" + " \"allowPartialResults\": false,\n" - + " \"reversed\": false,\n" + " \"timeRange\": [\n" + " 0,\n" - + " 9223372036854775807\n" + " ]\n" + " }\n" + "}"; + String expectedOutput = "{\n" + " \"startTime\": 1,\n" + " \"processingTime\": 2,\n" + + " \"queueTime\": 3,\n" + " \"responseSize\": 4,\n" + " \"blockBytesScanned\": 5,\n" + + " \"fsReadTime\": 6,\n" + " \"multiGetsCount\": 6,\n" + " \"multiMutationsCount\": 7,\n" + + " \"scan\": {\n" + " \"totalColumns\": 0,\n" + " \"maxResultSize\": -1,\n" + + " \"caching\": -1,\n" + " \"includeStopRow\": false,\n" + + " \"consistency\": \"STRONG\",\n" + " \"maxVersions\": 1,\n" + + " \"mvccReadPoint\": -1,\n" + " \"includeStartRow\": true,\n" + + " \"stopRow\": \"\\\\x00\\\\x00\\\\x01\\\\xC8\",\n" + " \"limit\": -1,\n" + + " \"timeRange\": [\n" + " 0,\n" + " 9223372036854775807\n" + " ],\n" + + " \"startRow\": \"\\\\x00\\\\x00\\\\x00{\",\n" + " \"targetReplicaId\": -1,\n" + + " \"batch\": -1,\n" + " \"families\": {},\n" + " \"priority\": -1,\n" + + " \"storeOffset\": 0,\n" + " \"queryMetricsEnabled\": false,\n" + + " \"needCursorResult\": false,\n" + " \"storeLimit\": -1,\n" + + " \"cacheBlocks\": true,\n" + " \"readType\": \"DEFAULT\",\n" + + " \"allowPartialResults\": false,\n" + " \"reversed\": false\n" + " }\n" + "}"; OnlineLogRecord o = new OnlineLogRecord(1, 2, 3, 4, 5, 6, null, null, null, null, null, null, null, 6, 7, 0, scan, Collections.emptyMap(), Collections.emptyMap()); String actualOutput = o.toJsonPrettyPrint(); diff --git a/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestScan.java b/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestScan.java index 66e83959718f..05d0fc865f7c 100644 --- a/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestScan.java +++ b/hbase-client/src/test/java/org/apache/hadoop/hbase/client/TestScan.java @@ -77,7 +77,8 @@ public void testGetToScan() throws Exception { .setAttribute("att_v0", Bytes.toBytes("att_v0")) .setColumnFamilyTimeRange(Bytes.toBytes("cf"), 0, 123).setReplicaId(3) .setACL("test_user", new Permission(Permission.Action.READ)) - .setAuthorizations(new Authorizations("test_label")).setPriority(3); + .setAuthorizations(new Authorizations("test_label")).setQueryMetricsEnabled(true) + .setPriority(3); Scan scan = new Scan(get); assertEquals(get.getCacheBlocks(), scan.getCacheBlocks()); @@ -101,6 +102,7 @@ public void testGetToScan() throws Exception { assertEquals(get.getACL(), scan.getACL()); assertEquals(get.getAuthorizations().getLabels(), scan.getAuthorizations().getLabels()); assertEquals(get.getPriority(), scan.getPriority()); + assertEquals(get.isQueryMetricsEnabled(), scan.isQueryMetricsEnabled()); } @Test @@ -217,7 +219,7 @@ public void testScanCopyConstructor() throws Exception { .setReplicaId(3).setReversed(true).setRowOffsetPerColumnFamily(5) .setStartStopRowForPrefixScan(Bytes.toBytes("row_")).setScanMetricsEnabled(true) .setSmall(true).setReadType(ReadType.STREAM).withStartRow(Bytes.toBytes("row_1")) - .withStopRow(Bytes.toBytes("row_2")).setTimeRange(0, 13); + .withStopRow(Bytes.toBytes("row_2")).setTimeRange(0, 13).setQueryMetricsEnabled(true); // create a copy of existing scan object Scan scanCopy = new Scan(scan); @@ -253,6 +255,7 @@ public void testScanCopyConstructor() throws Exception { assertEquals(scan.getStartRow(), scanCopy.getStartRow()); assertEquals(scan.getStopRow(), scanCopy.getStopRow()); assertEquals(scan.getTimeRange(), scanCopy.getTimeRange()); + assertEquals(scan.isQueryMetricsEnabled(), scanCopy.isQueryMetricsEnabled()); assertTrue("Make sure copy constructor adds all the fields in the copied object", EqualsBuilder.reflectionEquals(scan, scanCopy)); diff --git a/hbase-client/src/test/java/org/apache/hadoop/hbase/shaded/protobuf/TestProtobufUtil.java b/hbase-client/src/test/java/org/apache/hadoop/hbase/shaded/protobuf/TestProtobufUtil.java index 39fa2eb06c0a..a68fc58e6ffb 100644 --- a/hbase-client/src/test/java/org/apache/hadoop/hbase/shaded/protobuf/TestProtobufUtil.java +++ b/hbase-client/src/test/java/org/apache/hadoop/hbase/shaded/protobuf/TestProtobufUtil.java @@ -131,6 +131,7 @@ public void testGet() throws IOException { getBuilder.setMaxVersions(1); getBuilder.setCacheBlocks(true); getBuilder.setTimeRange(ProtobufUtil.toTimeRange(TimeRange.allTime())); + getBuilder.setQueryMetricsEnabled(false); Get get = ProtobufUtil.toGet(proto); assertEquals(getBuilder.build(), ProtobufUtil.toGet(get)); } @@ -259,6 +260,7 @@ public void testScan() throws IOException { scanBuilder.setCaching(1024); scanBuilder.setTimeRange(ProtobufUtil.toTimeRange(TimeRange.allTime())); scanBuilder.setIncludeStopRow(false); + scanBuilder.setQueryMetricsEnabled(false); ClientProtos.Scan expectedProto = scanBuilder.build(); ClientProtos.Scan actualProto = ProtobufUtil.toScan(ProtobufUtil.toScan(expectedProto)); diff --git a/hbase-protocol-shaded/src/main/protobuf/Client.proto b/hbase-protocol-shaded/src/main/protobuf/Client.proto index 13917b6d66cb..78aa138b4f05 100644 --- a/hbase-protocol-shaded/src/main/protobuf/Client.proto +++ b/hbase-protocol-shaded/src/main/protobuf/Client.proto @@ -90,6 +90,7 @@ message Get { optional Consistency consistency = 12 [default = STRONG]; repeated ColumnFamilyTimeRange cf_time_range = 13; optional bool load_column_families_on_demand = 14; /* DO NOT add defaults to load_column_families_on_demand. */ + optional bool query_metrics_enabled = 15 [default = false]; } message Result { @@ -117,6 +118,9 @@ message Result { // to form a complete result. The equivalent flag in o.a.h.h.client.Result is // mayHaveMoreCellsInRow. optional bool partial = 5 [default = false]; + + // Server side metrics about the result + optional QueryMetrics metrics = 6; } /** @@ -145,6 +149,7 @@ message Condition { optional Comparator comparator = 5; optional TimeRange time_range = 6; optional Filter filter = 7; + optional bool queryMetricsEnabled = 8; } @@ -233,6 +238,7 @@ message MutateResponse { // used for mutate to indicate processed only optional bool processed = 2; + optional QueryMetrics metrics = 3; } /** @@ -274,6 +280,7 @@ message Scan { } optional ReadType readType = 23 [default = DEFAULT]; optional bool need_cursor_result = 24 [default = false]; + optional bool query_metrics_enabled = 25 [default = false]; } /** @@ -366,6 +373,9 @@ message ScanResponse { // If the Scan need cursor, return the row key we are scanning in heartbeat message. // If the Scan doesn't need a cursor, don't set this field to reduce network IO. optional Cursor cursor = 12; + + // List of QueryMetrics that maps 1:1 to the results in the response based on index + repeated QueryMetrics query_metrics = 13; } /** @@ -458,6 +468,13 @@ message RegionAction { optional Condition condition = 4; } +/* +* Statistics about the Result's server-side metrics +*/ +message QueryMetrics { + optional uint64 block_bytes_scanned = 1; +} + /* * Statistics about the current load on the region */ @@ -491,6 +508,7 @@ message ResultOrException { optional CoprocessorServiceResult service_result = 4; // current load on the region optional RegionLoadStats loadStats = 5 [deprecated=true]; + optional QueryMetrics metrics = 6; } /** diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/HRegion.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/HRegion.java index 709b38ae926e..b12abf8f835d 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/HRegion.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/HRegion.java @@ -116,6 +116,7 @@ import org.apache.hadoop.hbase.client.IsolationLevel; import org.apache.hadoop.hbase.client.Mutation; import org.apache.hadoop.hbase.client.Put; +import org.apache.hadoop.hbase.client.QueryMetrics; import org.apache.hadoop.hbase.client.RegionInfo; import org.apache.hadoop.hbase.client.RegionInfoBuilder; import org.apache.hadoop.hbase.client.RegionReplicaUtil; @@ -4876,7 +4877,8 @@ private CheckAndMutateResult checkAndMutateInternal(CheckAndMutate checkAndMutat // we'll get the latest on this row. boolean matches = false; long cellTs = 0; - try (RegionScanner scanner = getScanner(new Scan(get))) { + QueryMetrics metrics = null; + try (RegionScannerImpl scanner = getScanner(new Scan(get))) { // NOTE: Please don't use HRegion.get() instead, // because it will copy cells to heap. See HBASE-26036 List result = new ArrayList<>(1); @@ -4901,6 +4903,9 @@ private CheckAndMutateResult checkAndMutateInternal(CheckAndMutate checkAndMutat matches = matches(op, compareResult); } } + if (checkAndMutate.isQueryMetricsEnabled()) { + metrics = new QueryMetrics(scanner.getContext().getBlockSizeProgress()); + } } // If matches, perform the mutation or the rowMutations @@ -4935,10 +4940,10 @@ private CheckAndMutateResult checkAndMutateInternal(CheckAndMutate checkAndMutat r = mutateRow(rowMutations, nonceGroup, nonce); } this.checkAndMutateChecksPassed.increment(); - return new CheckAndMutateResult(true, r); + return new CheckAndMutateResult(true, r).setMetrics(metrics); } this.checkAndMutateChecksFailed.increment(); - return new CheckAndMutateResult(false, null); + return new CheckAndMutateResult(false, null).setMetrics(metrics); } finally { rowLock.release(); } diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/RSRpcServices.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/RSRpcServices.java index b77fcf338a50..9892f48016b8 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/RSRpcServices.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/RSRpcServices.java @@ -78,6 +78,7 @@ import org.apache.hadoop.hbase.client.Mutation; import org.apache.hadoop.hbase.client.OperationWithAttributes; import org.apache.hadoop.hbase.client.Put; +import org.apache.hadoop.hbase.client.QueryMetrics; import org.apache.hadoop.hbase.client.RegionInfo; import org.apache.hadoop.hbase.client.RegionReplicaUtil; import org.apache.hadoop.hbase.client.Result; @@ -2673,6 +2674,7 @@ private Result get(Get get, HRegion region, RegionScannersCloseCallBack closeCal scan.setLoadColumnFamiliesOnDemand(region.isLoadingCfsOnDemandDefault()); } RegionScannerImpl scanner = null; + long blockBytesScannedBefore = context.getBlockBytesScanned(); try { scanner = region.getScanner(scan); scanner.next(results); @@ -2700,7 +2702,13 @@ private Result get(Get get, HRegion region, RegionScannersCloseCallBack closeCal } region.metricsUpdateForGet(results, before); - return Result.create(results, get.isCheckExistenceOnly() ? !results.isEmpty() : null, stale); + Result r = + Result.create(results, get.isCheckExistenceOnly() ? !results.isEmpty() : null, stale); + if (get.isQueryMetricsEnabled()) { + long blockBytesScanned = context.getBlockBytesScanned() - blockBytesScannedBefore; + r.setMetrics(new QueryMetrics(blockBytesScanned)); + } + return r; } private void checkBatchSizeAndLogLargeSize(MultiRequest request) throws ServiceException { @@ -2867,6 +2875,12 @@ public MultiResponse multi(final RpcController rpcc, final MultiRequest request) if (result.getResult() != null) { resultOrExceptionOrBuilder.setResult(ProtobufUtil.toResult(result.getResult())); } + + if (result.getMetrics() != null) { + resultOrExceptionOrBuilder + .setMetrics(ProtobufUtil.toQueryMetrics(result.getMetrics())); + } + regionActionResultBuilder.addResultOrException(resultOrExceptionOrBuilder.build()); } else { CheckAndMutateResult result = checkAndMutate(region, regionAction.getActionList(), @@ -3016,6 +3030,9 @@ public MutateResponse mutate(final RpcController rpcc, final MutateRequest reque if (clientCellBlockSupported) { addSize(context, result.getResult()); } + if (result.getMetrics() != null) { + builder.setMetrics(ProtobufUtil.toQueryMetrics(result.getMetrics())); + } } else { Result r = null; Boolean processed = null; @@ -3428,6 +3445,7 @@ private void scan(HBaseRpcController controller, ScanRequest request, RegionScan contextBuilder.setTrackMetrics(trackMetrics); ScannerContext scannerContext = contextBuilder.build(); boolean limitReached = false; + long blockBytesScannedBefore = 0; while (numOfResults < maxResults) { // Reset the batch progress to 0 before every call to RegionScanner#nextRaw. The // batch limit is a limit on the number of cells per Result. Thus, if progress is @@ -3439,6 +3457,10 @@ private void scan(HBaseRpcController controller, ScanRequest request, RegionScan // Collect values to be returned here moreRows = scanner.nextRaw(values, scannerContext); + + long blockBytesScanned = scannerContext.getBlockSizeProgress() - blockBytesScannedBefore; + blockBytesScannedBefore = scannerContext.getBlockSizeProgress(); + if (rpcCall == null) { // When there is no RpcCallContext,copy EC to heap, then the scanner would close, // This can be an EXPENSIVE call. It may make an extra copy from offheap to onheap @@ -3477,6 +3499,12 @@ private void scan(HBaseRpcController controller, ScanRequest request, RegionScan } boolean mayHaveMoreCellsInRow = scannerContext.mayHaveMoreCellsInRow(); Result r = Result.create(values, null, stale, mayHaveMoreCellsInRow); + + if (request.getScan().getQueryMetricsEnabled()) { + builder.addQueryMetrics(ClientProtos.QueryMetrics.newBuilder() + .setBlockBytesScanned(blockBytesScanned).build()); + } + results.add(r); numOfResults++; if (!mayHaveMoreCellsInRow && limitOfRows > 0) { diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/RegionScannerImpl.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/RegionScannerImpl.java index d829b1961070..4bebf5d927c8 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/RegionScannerImpl.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/RegionScannerImpl.java @@ -146,6 +146,10 @@ private static boolean hasNonce(HRegion region, long nonce) { initializeScanners(scan, additionalScanners); } + public ScannerContext getContext() { + return defaultScannerContext; + } + private void initializeScanners(Scan scan, List additionalScanners) throws IOException { // Here we separate all scanners into two lists - scanner that provide data required diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/client/TestAsyncTableQueryMetrics.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/client/TestAsyncTableQueryMetrics.java new file mode 100644 index 000000000000..7dd9803244f3 --- /dev/null +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/client/TestAsyncTableQueryMetrics.java @@ -0,0 +1,242 @@ +/* + * 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.hadoop.hbase.client; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import org.apache.hadoop.hbase.HBaseClassTestRule; +import org.apache.hadoop.hbase.HBaseTestingUtility; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.regionserver.MetricsRegionServer; +import org.apache.hadoop.hbase.regionserver.MetricsRegionServerSource; +import org.apache.hadoop.hbase.regionserver.MetricsRegionServerSourceImpl; +import org.apache.hadoop.hbase.testclassification.ClientTests; +import org.apache.hadoop.hbase.testclassification.MediumTests; +import org.apache.hadoop.hbase.util.Bytes; +import org.apache.hadoop.hbase.util.JVMClusterUtil; +import org.junit.AfterClass; +import org.junit.Assert; +import org.junit.BeforeClass; +import org.junit.ClassRule; +import org.junit.Test; +import org.junit.experimental.categories.Category; + +import org.apache.hbase.thirdparty.com.google.common.collect.ImmutableList; +import org.apache.hbase.thirdparty.com.google.common.io.Closeables; + +@Category({ MediumTests.class, ClientTests.class }) +public class TestAsyncTableQueryMetrics { + + @ClassRule + public static final HBaseClassTestRule CLASS_RULE = + HBaseClassTestRule.forClass(TestAsyncTableQueryMetrics.class); + + private static final HBaseTestingUtility UTIL = new HBaseTestingUtility(); + + private static final TableName TABLE_NAME = TableName.valueOf("ResultMetrics"); + + private static final byte[] CF = Bytes.toBytes("cf"); + + private static final byte[] CQ = Bytes.toBytes("cq"); + + private static final byte[] VALUE = Bytes.toBytes("value"); + + private static final byte[] ROW_1 = Bytes.toBytes("zzz1"); + private static final byte[] ROW_2 = Bytes.toBytes("zzz2"); + private static final byte[] ROW_3 = Bytes.toBytes("zzz3"); + + private static AsyncConnection CONN; + + @BeforeClass + public static void setUp() throws Exception { + UTIL.startMiniCluster(3); + // Create 3 rows in the table, with rowkeys starting with "zzz*" so that + // scan are forced to hit all the regions. + try (Table table = UTIL.createMultiRegionTable(TABLE_NAME, CF)) { + table.put(Arrays.asList(new Put(ROW_1).addColumn(CF, CQ, VALUE), + new Put(ROW_2).addColumn(CF, CQ, VALUE), new Put(ROW_3).addColumn(CF, CQ, VALUE))); + } + CONN = ConnectionFactory.createAsyncConnection(UTIL.getConfiguration()).get(); + CONN.getAdmin().flush(TABLE_NAME).join(); + } + + @AfterClass + public static void tearDown() throws Exception { + Closeables.close(CONN, true); + UTIL.shutdownMiniCluster(); + } + + @Test + public void itTestsGets() throws Exception { + // Test a single Get + Get g1 = new Get(ROW_1); + g1.setQueryMetricsEnabled(true); + + long bbs = getClusterBlockBytesScanned(); + Result result = CONN.getTable(TABLE_NAME).get(g1).get(); + bbs += result.getMetrics().getBlockBytesScanned(); + Assert.assertNotNull(result.getMetrics()); + Assert.assertEquals(getClusterBlockBytesScanned(), bbs); + + // Test multigets + Get g2 = new Get(ROW_2); + g2.setQueryMetricsEnabled(true); + + Get g3 = new Get(ROW_3); + g3.setQueryMetricsEnabled(true); + + List> futures = + CONN.getTable(TABLE_NAME).get(ImmutableList.of(g1, g2, g3)); + + for (CompletableFuture future : futures) { + result = future.join(); + Assert.assertNotNull(result.getMetrics()); + bbs += result.getMetrics().getBlockBytesScanned(); + } + + Assert.assertEquals(getClusterBlockBytesScanned(), bbs); + } + + @Test + public void itTestsDefaultGetNoMetrics() throws Exception { + // Test a single Get + Get g1 = new Get(ROW_1); + + Result result = CONN.getTable(TABLE_NAME).get(g1).get(); + Assert.assertNull(result.getMetrics()); + + // Test multigets + Get g2 = new Get(ROW_2); + Get g3 = new Get(ROW_3); + List> futures = + CONN.getTable(TABLE_NAME).get(ImmutableList.of(g1, g2, g3)); + futures.forEach(f -> Assert.assertNull(f.join().getMetrics())); + + } + + @Test + public void itTestsScans() { + Scan scan = new Scan(); + scan.setQueryMetricsEnabled(true); + + long bbs = getClusterBlockBytesScanned(); + try (ResultScanner scanner = CONN.getTable(TABLE_NAME).getScanner(scan)) { + for (Result result : scanner) { + Assert.assertNotNull(result.getMetrics()); + bbs += result.getMetrics().getBlockBytesScanned(); + Assert.assertEquals(getClusterBlockBytesScanned(), bbs); + } + } + } + + @Test + public void itTestsDefaultScanNoMetrics() { + Scan scan = new Scan(); + + try (ResultScanner scanner = CONN.getTable(TABLE_NAME).getScanner(scan)) { + for (Result result : scanner) { + Assert.assertNull(result.getMetrics()); + } + } + } + + @Test + public void itTestsAtomicOperations() { + CheckAndMutate cam = CheckAndMutate.newBuilder(ROW_1).ifEquals(CF, CQ, VALUE) + .queryMetricsEnabled(true).build(new Put(ROW_1).addColumn(CF, CQ, VALUE)); + + long bbs = getClusterBlockBytesScanned(); + CheckAndMutateResult result = CONN.getTable(TABLE_NAME).checkAndMutate(cam).join(); + QueryMetrics metrics = result.getMetrics(); + + Assert.assertNotNull(metrics); + Assert.assertEquals(getClusterBlockBytesScanned(), bbs + metrics.getBlockBytesScanned()); + + bbs = getClusterBlockBytesScanned(); + List batch = new ArrayList<>(); + batch.add(cam); + batch.add(CheckAndMutate.newBuilder(ROW_2).queryMetricsEnabled(true).ifEquals(CF, CQ, VALUE) + .build(new Put(ROW_2).addColumn(CF, CQ, VALUE))); + batch.add(CheckAndMutate.newBuilder(ROW_3).queryMetricsEnabled(true).ifEquals(CF, CQ, VALUE) + .build(new Put(ROW_3).addColumn(CF, CQ, VALUE))); + + List res = CONN.getTable(TABLE_NAME).batchAll(batch).join(); + long totalBbs = res.stream() + .mapToLong(r -> ((CheckAndMutateResult) r).getMetrics().getBlockBytesScanned()).sum(); + Assert.assertEquals(getClusterBlockBytesScanned(), bbs + totalBbs); + + bbs = getClusterBlockBytesScanned(); + + // flush to force fetch from disk + CONN.getAdmin().flush(TABLE_NAME).join(); + List> futures = CONN.getTable(TABLE_NAME).batch(batch); + + totalBbs = futures.stream().map(CompletableFuture::join) + .mapToLong(r -> ((CheckAndMutateResult) r).getMetrics().getBlockBytesScanned()).sum(); + Assert.assertEquals(getClusterBlockBytesScanned(), bbs + totalBbs); + } + + @Test + public void itTestsDefaultAtomicOperations() { + CheckAndMutate cam = CheckAndMutate.newBuilder(ROW_1).ifEquals(CF, CQ, VALUE) + .build(new Put(ROW_1).addColumn(CF, CQ, VALUE)); + + CheckAndMutateResult result = CONN.getTable(TABLE_NAME).checkAndMutate(cam).join(); + QueryMetrics metrics = result.getMetrics(); + + Assert.assertNull(metrics); + + List batch = new ArrayList<>(); + batch.add(cam); + batch.add(CheckAndMutate.newBuilder(ROW_2).ifEquals(CF, CQ, VALUE) + .build(new Put(ROW_2).addColumn(CF, CQ, VALUE))); + batch.add(CheckAndMutate.newBuilder(ROW_3).ifEquals(CF, CQ, VALUE) + .build(new Put(ROW_3).addColumn(CF, CQ, VALUE))); + + List res = CONN.getTable(TABLE_NAME).batchAll(batch).join(); + for (Object r : res) { + Assert.assertNull(((CheckAndMutateResult) r).getMetrics()); + } + + // flush to force fetch from disk + CONN.getAdmin().flush(TABLE_NAME).join(); + List> futures = CONN.getTable(TABLE_NAME).batch(batch); + + for (CompletableFuture future : futures) { + Object r = future.join(); + Assert.assertNull(((CheckAndMutateResult) r).getMetrics()); + } + } + + private static long getClusterBlockBytesScanned() { + long bbs = 0L; + + for (JVMClusterUtil.RegionServerThread rs : UTIL.getHBaseCluster().getRegionServerThreads()) { + MetricsRegionServer metrics = rs.getRegionServer().getMetrics(); + MetricsRegionServerSourceImpl source = + (MetricsRegionServerSourceImpl) metrics.getMetricsSource(); + + bbs += source.getMetricsRegistry() + .getCounter(MetricsRegionServerSource.BLOCK_BYTES_SCANNED_KEY, 0L).value(); + } + + return bbs; + } +} diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/client/TestMalformedCellFromClient.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/client/TestMalformedCellFromClient.java index e02b7ef07375..f787a659e138 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/client/TestMalformedCellFromClient.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/client/TestMalformedCellFromClient.java @@ -237,7 +237,7 @@ private static ClientProtos.MultiRequest createRequest(RowMutations rm, byte[] r ClientProtos.Action.Builder actionBuilder = ClientProtos.Action.newBuilder(); ClientProtos.MutationProto.Builder mutationBuilder = ClientProtos.MutationProto.newBuilder(); ClientProtos.Condition condition = ProtobufUtil.toCondition(rm.getRow(), FAMILY, null, - CompareOperator.EQUAL, new byte[10], null, null); + CompareOperator.EQUAL, new byte[10], null, null, false); for (Mutation mutation : rm.getMutations()) { ClientProtos.MutationProto.MutationType mutateType = null; if (mutation instanceof Put) {