001/* 002 * Licensed to the Apache Software Foundation (ASF) under one 003 * or more contributor license agreements. See the NOTICE file 004 * distributed with this work for additional information 005 * regarding copyright ownership. The ASF licenses this file 006 * to you under the Apache License, Version 2.0 (the 007 * "License"); you may not use this file except in compliance 008 * with the License. You may obtain a copy of the License at 009 * 010 * http://www.apache.org/licenses/LICENSE-2.0 011 * 012 * Unless required by applicable law or agreed to in writing, software 013 * distributed under the License is distributed on an "AS IS" BASIS, 014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 015 * See the License for the specific language governing permissions and 016 * limitations under the License. 017 */ 018package org.apache.hadoop.hbase.regionserver; 019 020import com.google.errorprone.annotations.RestrictedApi; 021import java.io.FileNotFoundException; 022import java.io.IOException; 023import java.io.UncheckedIOException; 024import java.lang.reflect.InvocationTargetException; 025import java.lang.reflect.Method; 026import java.net.BindException; 027import java.net.InetAddress; 028import java.net.InetSocketAddress; 029import java.util.ArrayList; 030import java.util.Arrays; 031import java.util.Collections; 032import java.util.HashMap; 033import java.util.Iterator; 034import java.util.List; 035import java.util.Map; 036import java.util.Map.Entry; 037import java.util.NavigableMap; 038import java.util.Set; 039import java.util.TreeSet; 040import java.util.concurrent.ConcurrentHashMap; 041import java.util.concurrent.ConcurrentMap; 042import java.util.concurrent.TimeUnit; 043import java.util.concurrent.atomic.AtomicBoolean; 044import java.util.concurrent.atomic.AtomicLong; 045import java.util.concurrent.atomic.LongAdder; 046import org.apache.hadoop.conf.Configuration; 047import org.apache.hadoop.fs.FileSystem; 048import org.apache.hadoop.fs.Path; 049import org.apache.hadoop.hbase.CacheEvictionStats; 050import org.apache.hadoop.hbase.CacheEvictionStatsBuilder; 051import org.apache.hadoop.hbase.Cell; 052import org.apache.hadoop.hbase.CellScanner; 053import org.apache.hadoop.hbase.CellUtil; 054import org.apache.hadoop.hbase.DoNotRetryIOException; 055import org.apache.hadoop.hbase.DroppedSnapshotException; 056import org.apache.hadoop.hbase.ExtendedCellScannable; 057import org.apache.hadoop.hbase.ExtendedCellScanner; 058import org.apache.hadoop.hbase.HBaseIOException; 059import org.apache.hadoop.hbase.HBaseRpcServicesBase; 060import org.apache.hadoop.hbase.HConstants; 061import org.apache.hadoop.hbase.MultiActionResultTooLarge; 062import org.apache.hadoop.hbase.NotServingRegionException; 063import org.apache.hadoop.hbase.PrivateCellUtil; 064import org.apache.hadoop.hbase.RegionTooBusyException; 065import org.apache.hadoop.hbase.Server; 066import org.apache.hadoop.hbase.ServerName; 067import org.apache.hadoop.hbase.TableName; 068import org.apache.hadoop.hbase.UnknownScannerException; 069import org.apache.hadoop.hbase.client.Append; 070import org.apache.hadoop.hbase.client.CheckAndMutate; 071import org.apache.hadoop.hbase.client.CheckAndMutateResult; 072import org.apache.hadoop.hbase.client.ClientInternalHelper; 073import org.apache.hadoop.hbase.client.Delete; 074import org.apache.hadoop.hbase.client.Durability; 075import org.apache.hadoop.hbase.client.Get; 076import org.apache.hadoop.hbase.client.Increment; 077import org.apache.hadoop.hbase.client.Mutation; 078import org.apache.hadoop.hbase.client.OperationWithAttributes; 079import org.apache.hadoop.hbase.client.Put; 080import org.apache.hadoop.hbase.client.QueryMetrics; 081import org.apache.hadoop.hbase.client.RegionInfo; 082import org.apache.hadoop.hbase.client.RegionReplicaUtil; 083import org.apache.hadoop.hbase.client.Result; 084import org.apache.hadoop.hbase.client.Row; 085import org.apache.hadoop.hbase.client.Scan; 086import org.apache.hadoop.hbase.client.TableDescriptor; 087import org.apache.hadoop.hbase.client.VersionInfoUtil; 088import org.apache.hadoop.hbase.client.metrics.ServerSideScanMetrics; 089import org.apache.hadoop.hbase.exceptions.FailedSanityCheckException; 090import org.apache.hadoop.hbase.exceptions.OutOfOrderScannerNextException; 091import org.apache.hadoop.hbase.exceptions.ScannerResetException; 092import org.apache.hadoop.hbase.exceptions.TimeoutIOException; 093import org.apache.hadoop.hbase.exceptions.UnknownProtocolException; 094import org.apache.hadoop.hbase.io.ByteBuffAllocator; 095import org.apache.hadoop.hbase.io.hfile.BlockCache; 096import org.apache.hadoop.hbase.ipc.HBaseRpcController; 097import org.apache.hadoop.hbase.ipc.PriorityFunction; 098import org.apache.hadoop.hbase.ipc.QosPriority; 099import org.apache.hadoop.hbase.ipc.RpcCall; 100import org.apache.hadoop.hbase.ipc.RpcCallContext; 101import org.apache.hadoop.hbase.ipc.RpcCallback; 102import org.apache.hadoop.hbase.ipc.RpcServer; 103import org.apache.hadoop.hbase.ipc.RpcServer.BlockingServiceAndInterface; 104import org.apache.hadoop.hbase.ipc.RpcServerFactory; 105import org.apache.hadoop.hbase.ipc.RpcServerInterface; 106import org.apache.hadoop.hbase.ipc.ServerNotRunningYetException; 107import org.apache.hadoop.hbase.ipc.ServerRpcController; 108import org.apache.hadoop.hbase.monitoring.ThreadLocalServerSideScanMetrics; 109import org.apache.hadoop.hbase.namequeues.NamedQueueRecorder; 110import org.apache.hadoop.hbase.namequeues.RpcLogDetails; 111import org.apache.hadoop.hbase.namequeues.request.NamedQueueGetRequest; 112import org.apache.hadoop.hbase.namequeues.response.NamedQueueGetResponse; 113import org.apache.hadoop.hbase.net.Address; 114import org.apache.hadoop.hbase.procedure2.RSProcedureCallable; 115import org.apache.hadoop.hbase.quotas.ActivePolicyEnforcement; 116import org.apache.hadoop.hbase.quotas.OperationQuota; 117import org.apache.hadoop.hbase.quotas.QuotaUtil; 118import org.apache.hadoop.hbase.quotas.RegionServerRpcQuotaManager; 119import org.apache.hadoop.hbase.quotas.RegionServerSpaceQuotaManager; 120import org.apache.hadoop.hbase.quotas.SpaceQuotaSnapshot; 121import org.apache.hadoop.hbase.quotas.SpaceViolationPolicyEnforcement; 122import org.apache.hadoop.hbase.regionserver.LeaseManager.Lease; 123import org.apache.hadoop.hbase.regionserver.LeaseManager.LeaseStillHeldException; 124import org.apache.hadoop.hbase.regionserver.Region.Operation; 125import org.apache.hadoop.hbase.regionserver.ScannerContext.LimitScope; 126import org.apache.hadoop.hbase.regionserver.compactions.CompactionLifeCycleTracker; 127import org.apache.hadoop.hbase.regionserver.handler.AssignRegionHandler; 128import org.apache.hadoop.hbase.regionserver.handler.OpenMetaHandler; 129import org.apache.hadoop.hbase.regionserver.handler.OpenPriorityRegionHandler; 130import org.apache.hadoop.hbase.regionserver.handler.OpenRegionHandler; 131import org.apache.hadoop.hbase.regionserver.handler.UnassignRegionHandler; 132import org.apache.hadoop.hbase.replication.ReplicationUtils; 133import org.apache.hadoop.hbase.replication.regionserver.RejectReplicationRequestStateChecker; 134import org.apache.hadoop.hbase.replication.regionserver.RejectRequestsFromClientStateChecker; 135import org.apache.hadoop.hbase.security.Superusers; 136import org.apache.hadoop.hbase.security.access.Permission; 137import org.apache.hadoop.hbase.util.Bytes; 138import org.apache.hadoop.hbase.util.DNS; 139import org.apache.hadoop.hbase.util.DNS.ServerType; 140import org.apache.hadoop.hbase.util.EnvironmentEdgeManager; 141import org.apache.hadoop.hbase.util.Pair; 142import org.apache.hadoop.hbase.util.ServerRegionReplicaUtil; 143import org.apache.hadoop.hbase.wal.WAL; 144import org.apache.hadoop.hbase.wal.WALEdit; 145import org.apache.hadoop.hbase.wal.WALKey; 146import org.apache.hadoop.hbase.wal.WALSplitUtil; 147import org.apache.hadoop.hbase.wal.WALSplitUtil.MutationReplay; 148import org.apache.hadoop.hbase.zookeeper.ZKWatcher; 149import org.apache.yetus.audience.InterfaceAudience; 150import org.slf4j.Logger; 151import org.slf4j.LoggerFactory; 152 153import org.apache.hbase.thirdparty.com.google.common.cache.Cache; 154import org.apache.hbase.thirdparty.com.google.common.cache.CacheBuilder; 155import org.apache.hbase.thirdparty.com.google.common.collect.ImmutableList; 156import org.apache.hbase.thirdparty.com.google.common.collect.Lists; 157import org.apache.hbase.thirdparty.com.google.protobuf.ByteString; 158import org.apache.hbase.thirdparty.com.google.protobuf.Message; 159import org.apache.hbase.thirdparty.com.google.protobuf.RpcController; 160import org.apache.hbase.thirdparty.com.google.protobuf.ServiceException; 161import org.apache.hbase.thirdparty.com.google.protobuf.TextFormat; 162import org.apache.hbase.thirdparty.com.google.protobuf.UnsafeByteOperations; 163import org.apache.hbase.thirdparty.org.apache.commons.collections4.CollectionUtils; 164 165import org.apache.hadoop.hbase.shaded.protobuf.ProtobufUtil; 166import org.apache.hadoop.hbase.shaded.protobuf.RequestConverter; 167import org.apache.hadoop.hbase.shaded.protobuf.ResponseConverter; 168import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.AdminService; 169import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.ClearCompactionQueuesRequest; 170import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.ClearCompactionQueuesResponse; 171import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.ClearRegionBlockCacheRequest; 172import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.ClearRegionBlockCacheResponse; 173import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.CloseRegionRequest; 174import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.CloseRegionResponse; 175import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.CompactRegionRequest; 176import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.CompactRegionResponse; 177import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.CompactionSwitchRequest; 178import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.CompactionSwitchResponse; 179import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.ExecuteProceduresRequest; 180import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.ExecuteProceduresResponse; 181import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.FlushRegionRequest; 182import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.FlushRegionResponse; 183import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetCachedFilesListRequest; 184import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetCachedFilesListResponse; 185import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetOnlineRegionRequest; 186import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetOnlineRegionResponse; 187import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetRegionInfoRequest; 188import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetRegionInfoResponse; 189import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetRegionLoadRequest; 190import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetRegionLoadResponse; 191import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetServerInfoRequest; 192import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetServerInfoResponse; 193import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetStoreFileRequest; 194import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.GetStoreFileResponse; 195import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.OpenRegionRequest; 196import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.OpenRegionRequest.RegionOpenInfo; 197import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.OpenRegionResponse; 198import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.OpenRegionResponse.RegionOpeningState; 199import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.RemoteProcedureRequest; 200import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.ReplicateWALEntryRequest; 201import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.ReplicateWALEntryResponse; 202import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.RollWALWriterRequest; 203import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.RollWALWriterResponse; 204import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.SlowLogResponseRequest; 205import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.SlowLogResponses; 206import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.StopServerRequest; 207import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.StopServerResponse; 208import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.UpdateFavoredNodesRequest; 209import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.UpdateFavoredNodesResponse; 210import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.WALEntry; 211import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.WarmupRegionRequest; 212import org.apache.hadoop.hbase.shaded.protobuf.generated.AdminProtos.WarmupRegionResponse; 213import org.apache.hadoop.hbase.shaded.protobuf.generated.BootstrapNodeProtos.BootstrapNodeService; 214import org.apache.hadoop.hbase.shaded.protobuf.generated.BootstrapNodeProtos.GetAllBootstrapNodesRequest; 215import org.apache.hadoop.hbase.shaded.protobuf.generated.BootstrapNodeProtos.GetAllBootstrapNodesResponse; 216import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos; 217import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.Action; 218import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.BulkLoadHFileRequest; 219import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.BulkLoadHFileRequest.FamilyPath; 220import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.BulkLoadHFileResponse; 221import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.CleanupBulkLoadRequest; 222import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.CleanupBulkLoadResponse; 223import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.ClientService; 224import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.Condition; 225import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.CoprocessorServiceRequest; 226import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.CoprocessorServiceResponse; 227import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.GetRequest; 228import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.GetResponse; 229import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.MultiRegionLoadStats; 230import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.MultiRequest; 231import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.MultiResponse; 232import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.MutateRequest; 233import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.MutateResponse; 234import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.MutationProto; 235import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.MutationProto.MutationType; 236import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.PrepareBulkLoadRequest; 237import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.PrepareBulkLoadResponse; 238import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.RegionAction; 239import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.RegionActionResult; 240import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.ResultOrException; 241import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.ScanRequest; 242import org.apache.hadoop.hbase.shaded.protobuf.generated.ClientProtos.ScanResponse; 243import org.apache.hadoop.hbase.shaded.protobuf.generated.ClusterStatusProtos; 244import org.apache.hadoop.hbase.shaded.protobuf.generated.ClusterStatusProtos.RegionLoad; 245import org.apache.hadoop.hbase.shaded.protobuf.generated.HBaseProtos; 246import org.apache.hadoop.hbase.shaded.protobuf.generated.HBaseProtos.BooleanMsg; 247import org.apache.hadoop.hbase.shaded.protobuf.generated.HBaseProtos.EmptyMsg; 248import org.apache.hadoop.hbase.shaded.protobuf.generated.HBaseProtos.ManagedKeyEntryRequest; 249import org.apache.hadoop.hbase.shaded.protobuf.generated.HBaseProtos.NameBytesPair; 250import org.apache.hadoop.hbase.shaded.protobuf.generated.HBaseProtos.NameInt64Pair; 251import org.apache.hadoop.hbase.shaded.protobuf.generated.HBaseProtos.RegionSpecifier; 252import org.apache.hadoop.hbase.shaded.protobuf.generated.HBaseProtos.RegionSpecifier.RegionSpecifierType; 253import org.apache.hadoop.hbase.shaded.protobuf.generated.MapReduceProtos.ScanMetrics; 254import org.apache.hadoop.hbase.shaded.protobuf.generated.QuotaProtos.GetSpaceQuotaSnapshotsRequest; 255import org.apache.hadoop.hbase.shaded.protobuf.generated.QuotaProtos.GetSpaceQuotaSnapshotsResponse; 256import org.apache.hadoop.hbase.shaded.protobuf.generated.QuotaProtos.GetSpaceQuotaSnapshotsResponse.TableQuotaSnapshot; 257import org.apache.hadoop.hbase.shaded.protobuf.generated.RegistryProtos.ClientMetaService; 258import org.apache.hadoop.hbase.shaded.protobuf.generated.TooSlowLog.SlowLogPayload; 259import org.apache.hadoop.hbase.shaded.protobuf.generated.WALProtos.BulkLoadDescriptor; 260import org.apache.hadoop.hbase.shaded.protobuf.generated.WALProtos.CompactionDescriptor; 261import org.apache.hadoop.hbase.shaded.protobuf.generated.WALProtos.FlushDescriptor; 262import org.apache.hadoop.hbase.shaded.protobuf.generated.WALProtos.RegionEventDescriptor; 263 264/** 265 * Implements the regionserver RPC services. 266 */ 267@InterfaceAudience.Private 268public class RSRpcServices extends HBaseRpcServicesBase<HRegionServer> 269 implements ClientService.BlockingInterface, BootstrapNodeService.BlockingInterface { 270 271 private static final Logger LOG = LoggerFactory.getLogger(RSRpcServices.class); 272 273 /** RPC scheduler to use for the region server. */ 274 public static final String REGION_SERVER_RPC_SCHEDULER_FACTORY_CLASS = 275 "hbase.region.server.rpc.scheduler.factory.class"; 276 277 /** 278 * Minimum allowable time limit delta (in milliseconds) that can be enforced during scans. This 279 * configuration exists to prevent the scenario where a time limit is specified to be so 280 * restrictive that the time limit is reached immediately (before any cells are scanned). 281 */ 282 private static final String REGION_SERVER_RPC_MINIMUM_SCAN_TIME_LIMIT_DELTA = 283 "hbase.region.server.rpc.minimum.scan.time.limit.delta"; 284 /** 285 * Default value of {@link RSRpcServices#REGION_SERVER_RPC_MINIMUM_SCAN_TIME_LIMIT_DELTA} 286 */ 287 static final long DEFAULT_REGION_SERVER_RPC_MINIMUM_SCAN_TIME_LIMIT_DELTA = 10; 288 289 /** 290 * Whether to reject rows with size > threshold defined by 291 * {@link HConstants#BATCH_ROWS_THRESHOLD_NAME} 292 */ 293 private static final String REJECT_BATCH_ROWS_OVER_THRESHOLD = 294 "hbase.rpc.rows.size.threshold.reject"; 295 296 /** 297 * Default value of config {@link RSRpcServices#REJECT_BATCH_ROWS_OVER_THRESHOLD} 298 */ 299 private static final boolean DEFAULT_REJECT_BATCH_ROWS_OVER_THRESHOLD = false; 300 301 // Request counter. (Includes requests that are not serviced by regions.) 302 // Count only once for requests with multiple actions like multi/caching-scan/replayBatch 303 final LongAdder requestCount = new LongAdder(); 304 305 // Request counter for rpc get 306 final LongAdder rpcGetRequestCount = new LongAdder(); 307 308 // Request counter for rpc scan 309 final LongAdder rpcScanRequestCount = new LongAdder(); 310 311 // Request counter for scans that might end up in full scans 312 final LongAdder rpcFullScanRequestCount = new LongAdder(); 313 314 // Request counter for rpc multi 315 final LongAdder rpcMultiRequestCount = new LongAdder(); 316 317 // Request counter for rpc mutate 318 final LongAdder rpcMutateRequestCount = new LongAdder(); 319 320 private volatile long maxScannerResultSize; 321 322 private ScannerIdGenerator scannerIdGenerator; 323 private final ConcurrentMap<String, RegionScannerHolder> scanners = new ConcurrentHashMap<>(); 324 // Hold the name and last sequence number of a closed scanner for a while. This is used 325 // to keep compatible for old clients which may send next or close request to a region 326 // scanner which has already been exhausted. The entries will be removed automatically 327 // after scannerLeaseTimeoutPeriod. 328 private final Cache<String, Long> closedScanners; 329 /** 330 * The lease timeout period for client scanners (milliseconds). 331 */ 332 private final int scannerLeaseTimeoutPeriod; 333 334 /** 335 * The RPC timeout period (milliseconds) 336 */ 337 private final int rpcTimeout; 338 339 /** 340 * The minimum allowable delta to use for the scan limit 341 */ 342 private final long minimumScanTimeLimitDelta; 343 344 /** 345 * Row size threshold for multi requests above which a warning is logged 346 */ 347 private volatile int rowSizeWarnThreshold; 348 /* 349 * Whether we should reject requests with very high no of rows i.e. beyond threshold defined by 350 * rowSizeWarnThreshold 351 */ 352 private volatile boolean rejectRowsWithSizeOverThreshold; 353 354 final AtomicBoolean clearCompactionQueues = new AtomicBoolean(false); 355 356 /** 357 * Services launched in RSRpcServices. By default they are on but you can use the below booleans 358 * to selectively enable/disable these services (Rare is the case where you would ever turn off 359 * one or the other). 360 */ 361 public static final String REGIONSERVER_ADMIN_SERVICE_CONFIG = 362 "hbase.regionserver.admin.executorService"; 363 public static final String REGIONSERVER_CLIENT_SERVICE_CONFIG = 364 "hbase.regionserver.client.executorService"; 365 public static final String REGIONSERVER_CLIENT_META_SERVICE_CONFIG = 366 "hbase.regionserver.client.meta.executorService"; 367 public static final String REGIONSERVER_BOOTSTRAP_NODES_SERVICE_CONFIG = 368 "hbase.regionserver.bootstrap.nodes.executorService"; 369 370 /** 371 * An Rpc callback for closing a RegionScanner. 372 */ 373 private static final class RegionScannerCloseCallBack implements RpcCallback { 374 375 private final RegionScanner scanner; 376 377 public RegionScannerCloseCallBack(RegionScanner scanner) { 378 this.scanner = scanner; 379 } 380 381 @Override 382 public void run() throws IOException { 383 this.scanner.close(); 384 } 385 } 386 387 /** 388 * An Rpc callback for doing shipped() call on a RegionScanner. 389 */ 390 private class RegionScannerShippedCallBack implements RpcCallback { 391 private final String scannerName; 392 private final Shipper shipper; 393 private final Lease lease; 394 395 public RegionScannerShippedCallBack(String scannerName, Shipper shipper, Lease lease) { 396 this.scannerName = scannerName; 397 this.shipper = shipper; 398 this.lease = lease; 399 } 400 401 @Override 402 public void run() throws IOException { 403 this.shipper.shipped(); 404 // We're done. On way out re-add the above removed lease. The lease was temp removed for this 405 // Rpc call and we are at end of the call now. Time to add it back. 406 if (scanners.containsKey(scannerName)) { 407 if (lease != null) { 408 server.getLeaseManager().addLease(lease); 409 } 410 } 411 } 412 } 413 414 /** 415 * An RpcCallBack that creates a list of scanners that needs to perform callBack operation on 416 * completion of multiGets. 417 */ 418 static class RegionScannersCloseCallBack implements RpcCallback { 419 private final List<RegionScanner> scanners = new ArrayList<>(); 420 421 public void addScanner(RegionScanner scanner) { 422 this.scanners.add(scanner); 423 } 424 425 @Override 426 public void run() { 427 for (RegionScanner scanner : scanners) { 428 try { 429 scanner.close(); 430 } catch (IOException e) { 431 LOG.error("Exception while closing the scanner " + scanner, e); 432 } 433 } 434 } 435 } 436 437 static class RegionScannerContext { 438 final String scannerName; 439 final RegionScannerHolder holder; 440 final OperationQuota quota; 441 442 RegionScannerContext(String scannerName, RegionScannerHolder holder, OperationQuota quota) { 443 this.scannerName = scannerName; 444 this.holder = holder; 445 this.quota = quota; 446 } 447 } 448 449 /** 450 * Holder class which holds the RegionScanner, nextCallSeq and RpcCallbacks together. 451 */ 452 static final class RegionScannerHolder { 453 private final AtomicLong nextCallSeq = new AtomicLong(0); 454 private final RegionScanner s; 455 private final HRegion r; 456 private final RpcCallback closeCallBack; 457 private final RpcCallback shippedCallback; 458 private byte[] rowOfLastPartialResult; 459 private boolean needCursor; 460 private boolean fullRegionScan; 461 private final String clientIPAndPort; 462 private final String userName; 463 private volatile long maxBlockBytesScanned = 0; 464 private volatile long prevBlockBytesScanned = 0; 465 private volatile long prevBlockBytesScannedDifference = 0; 466 467 RegionScannerHolder(RegionScanner s, HRegion r, RpcCallback closeCallBack, 468 RpcCallback shippedCallback, boolean needCursor, boolean fullRegionScan, 469 String clientIPAndPort, String userName) { 470 this.s = s; 471 this.r = r; 472 this.closeCallBack = closeCallBack; 473 this.shippedCallback = shippedCallback; 474 this.needCursor = needCursor; 475 this.fullRegionScan = fullRegionScan; 476 this.clientIPAndPort = clientIPAndPort; 477 this.userName = userName; 478 } 479 480 long getNextCallSeq() { 481 return nextCallSeq.get(); 482 } 483 484 boolean incNextCallSeq(long currentSeq) { 485 // Use CAS to prevent multiple scan request running on the same scanner. 486 return nextCallSeq.compareAndSet(currentSeq, currentSeq + 1); 487 } 488 489 long getMaxBlockBytesScanned() { 490 return maxBlockBytesScanned; 491 } 492 493 long getPrevBlockBytesScannedDifference() { 494 return prevBlockBytesScannedDifference; 495 } 496 497 void updateBlockBytesScanned(long blockBytesScanned) { 498 prevBlockBytesScannedDifference = blockBytesScanned - prevBlockBytesScanned; 499 prevBlockBytesScanned = blockBytesScanned; 500 if (blockBytesScanned > maxBlockBytesScanned) { 501 maxBlockBytesScanned = blockBytesScanned; 502 } 503 } 504 505 // Should be called only when we need to print lease expired messages otherwise 506 // cache the String once made. 507 @Override 508 public String toString() { 509 return "clientIPAndPort=" + this.clientIPAndPort + ", userName=" + this.userName 510 + ", regionInfo=" + this.r.getRegionInfo().getRegionNameAsString(); 511 } 512 } 513 514 /** 515 * Instantiated as a scanner lease. If the lease times out, the scanner is closed 516 */ 517 private class ScannerListener implements LeaseListener { 518 private final String scannerName; 519 520 ScannerListener(final String n) { 521 this.scannerName = n; 522 } 523 524 @Override 525 public void leaseExpired() { 526 RegionScannerHolder rsh = scanners.remove(this.scannerName); 527 if (rsh == null) { 528 LOG.warn("Scanner lease {} expired but no outstanding scanner", this.scannerName); 529 return; 530 } 531 LOG.info("Scanner lease {} expired {}", this.scannerName, rsh); 532 server.getMetrics().incrScannerLeaseExpired(); 533 RegionScanner s = rsh.s; 534 HRegion region = null; 535 try { 536 region = server.getRegion(s.getRegionInfo().getRegionName()); 537 if (region != null && region.getCoprocessorHost() != null) { 538 region.getCoprocessorHost().preScannerClose(s); 539 } 540 } catch (IOException e) { 541 LOG.error("Closing scanner {} {}", this.scannerName, rsh, e); 542 } finally { 543 try { 544 s.close(); 545 if (region != null && region.getCoprocessorHost() != null) { 546 region.getCoprocessorHost().postScannerClose(s); 547 } 548 } catch (IOException e) { 549 LOG.error("Closing scanner {} {}", this.scannerName, rsh, e); 550 } 551 } 552 } 553 } 554 555 private static ResultOrException getResultOrException(final ClientProtos.Result r, 556 final int index) { 557 return getResultOrException(ResponseConverter.buildActionResult(r), index); 558 } 559 560 private static ResultOrException getResultOrException(final Exception e, final int index) { 561 return getResultOrException(ResponseConverter.buildActionResult(e), index); 562 } 563 564 private static ResultOrException getResultOrException(final ResultOrException.Builder builder, 565 final int index) { 566 return builder.setIndex(index).build(); 567 } 568 569 /** 570 * Checks for the following pre-checks in order: 571 * <ol> 572 * <li>RegionServer is running</li> 573 * <li>If authorization is enabled, then RPC caller has ADMIN permissions</li> 574 * </ol> 575 * @param requestName name of rpc request. Used in reporting failures to provide context. 576 * @throws ServiceException If any of the above listed pre-check fails. 577 */ 578 private void rpcPreCheck(String requestName) throws ServiceException { 579 try { 580 checkOpen(); 581 requirePermission(requestName, Permission.Action.ADMIN); 582 } catch (IOException ioe) { 583 throw new ServiceException(ioe); 584 } 585 } 586 587 private boolean isClientCellBlockSupport(RpcCallContext context) { 588 return context != null && context.isClientCellBlockSupported(); 589 } 590 591 private void addResult(final MutateResponse.Builder builder, final Result result, 592 final HBaseRpcController rpcc, boolean clientCellBlockSupported) { 593 if (result == null) return; 594 if (clientCellBlockSupported) { 595 builder.setResult(ProtobufUtil.toResultNoData(result)); 596 rpcc.setCellScanner(result.cellScanner()); 597 } else { 598 ClientProtos.Result pbr = ProtobufUtil.toResult(result); 599 builder.setResult(pbr); 600 } 601 } 602 603 private void addResults(ScanResponse.Builder builder, List<Result> results, 604 HBaseRpcController controller, boolean isDefaultRegion, boolean clientCellBlockSupported) { 605 builder.setStale(!isDefaultRegion); 606 if (results.isEmpty()) { 607 return; 608 } 609 if (clientCellBlockSupported) { 610 for (Result res : results) { 611 builder.addCellsPerResult(res.size()); 612 builder.addPartialFlagPerResult(res.mayHaveMoreCellsInRow()); 613 } 614 controller.setCellScanner(PrivateCellUtil.createExtendedCellScanner(results)); 615 } else { 616 for (Result res : results) { 617 ClientProtos.Result pbr = ProtobufUtil.toResult(res); 618 builder.addResults(pbr); 619 } 620 } 621 } 622 623 private CheckAndMutateResult checkAndMutate(HRegion region, List<ClientProtos.Action> actions, 624 CellScanner cellScanner, Condition condition, long nonceGroup, 625 ActivePolicyEnforcement spaceQuotaEnforcement) throws IOException { 626 int countOfCompleteMutation = 0; 627 try { 628 if (!region.getRegionInfo().isMetaRegion()) { 629 server.getMemStoreFlusher().reclaimMemStoreMemory(); 630 } 631 List<Mutation> mutations = new ArrayList<>(); 632 long nonce = HConstants.NO_NONCE; 633 for (ClientProtos.Action action : actions) { 634 if (action.hasGet()) { 635 throw new DoNotRetryIOException( 636 "Atomic put and/or delete only, not a Get=" + action.getGet()); 637 } 638 MutationProto mutation = action.getMutation(); 639 MutationType type = mutation.getMutateType(); 640 switch (type) { 641 case PUT: 642 Put put = ProtobufUtil.toPut(mutation, cellScanner); 643 ++countOfCompleteMutation; 644 checkCellSizeLimit(region, put); 645 spaceQuotaEnforcement.getPolicyEnforcement(region).check(put); 646 mutations.add(put); 647 break; 648 case DELETE: 649 Delete del = ProtobufUtil.toDelete(mutation, cellScanner); 650 ++countOfCompleteMutation; 651 spaceQuotaEnforcement.getPolicyEnforcement(region).check(del); 652 mutations.add(del); 653 break; 654 case INCREMENT: 655 Increment increment = ProtobufUtil.toIncrement(mutation, cellScanner); 656 nonce = mutation.hasNonce() ? mutation.getNonce() : HConstants.NO_NONCE; 657 ++countOfCompleteMutation; 658 checkCellSizeLimit(region, increment); 659 spaceQuotaEnforcement.getPolicyEnforcement(region).check(increment); 660 mutations.add(increment); 661 break; 662 case APPEND: 663 Append append = ProtobufUtil.toAppend(mutation, cellScanner); 664 nonce = mutation.hasNonce() ? mutation.getNonce() : HConstants.NO_NONCE; 665 ++countOfCompleteMutation; 666 checkCellSizeLimit(region, append); 667 spaceQuotaEnforcement.getPolicyEnforcement(region).check(append); 668 mutations.add(append); 669 break; 670 default: 671 throw new DoNotRetryIOException("invalid mutation type : " + type); 672 } 673 } 674 675 if (mutations.size() == 0) { 676 return new CheckAndMutateResult(true, null); 677 } else { 678 CheckAndMutate checkAndMutate = ProtobufUtil.toCheckAndMutate(condition, mutations); 679 CheckAndMutateResult result = null; 680 if (region.getCoprocessorHost() != null) { 681 result = region.getCoprocessorHost().preCheckAndMutate(checkAndMutate); 682 } 683 if (result == null) { 684 result = region.checkAndMutate(checkAndMutate, nonceGroup, nonce); 685 if (region.getCoprocessorHost() != null) { 686 result = region.getCoprocessorHost().postCheckAndMutate(checkAndMutate, result); 687 } 688 } 689 return result; 690 } 691 } finally { 692 // Currently, the checkAndMutate isn't supported by batch so it won't mess up the cell scanner 693 // even if the malformed cells are not skipped. 694 for (int i = countOfCompleteMutation; i < actions.size(); ++i) { 695 skipCellsForMutation(actions.get(i), cellScanner); 696 } 697 } 698 } 699 700 /** 701 * Execute an append mutation. 702 * @return result to return to client if default operation should be bypassed as indicated by 703 * RegionObserver, null otherwise 704 */ 705 private Result append(final HRegion region, final OperationQuota quota, 706 final MutationProto mutation, final CellScanner cellScanner, long nonceGroup, 707 ActivePolicyEnforcement spaceQuota, RpcCallContext context) throws IOException { 708 long before = EnvironmentEdgeManager.currentTime(); 709 Append append = ProtobufUtil.toAppend(mutation, cellScanner); 710 checkCellSizeLimit(region, append); 711 spaceQuota.getPolicyEnforcement(region).check(append); 712 quota.addMutation(append); 713 long blockBytesScannedBefore = context != null ? context.getBlockBytesScanned() : 0; 714 long nonce = mutation.hasNonce() ? mutation.getNonce() : HConstants.NO_NONCE; 715 Result r = region.append(append, nonceGroup, nonce); 716 if (server.getMetrics() != null) { 717 long blockBytesScanned = 718 context != null ? context.getBlockBytesScanned() - blockBytesScannedBefore : 0; 719 server.getMetrics().updateAppend(region, EnvironmentEdgeManager.currentTime() - before, 720 blockBytesScanned); 721 } 722 return r == null ? Result.EMPTY_RESULT : r; 723 } 724 725 /** 726 * Execute an increment mutation. 727 */ 728 private Result increment(final HRegion region, final OperationQuota quota, 729 final MutationProto mutation, final CellScanner cells, long nonceGroup, 730 ActivePolicyEnforcement spaceQuota, RpcCallContext context) throws IOException { 731 long before = EnvironmentEdgeManager.currentTime(); 732 Increment increment = ProtobufUtil.toIncrement(mutation, cells); 733 checkCellSizeLimit(region, increment); 734 spaceQuota.getPolicyEnforcement(region).check(increment); 735 quota.addMutation(increment); 736 long blockBytesScannedBefore = context != null ? context.getBlockBytesScanned() : 0; 737 long nonce = mutation.hasNonce() ? mutation.getNonce() : HConstants.NO_NONCE; 738 Result r = region.increment(increment, nonceGroup, nonce); 739 final MetricsRegionServer metricsRegionServer = server.getMetrics(); 740 if (metricsRegionServer != null) { 741 long blockBytesScanned = 742 context != null ? context.getBlockBytesScanned() - blockBytesScannedBefore : 0; 743 metricsRegionServer.updateIncrement(region, EnvironmentEdgeManager.currentTime() - before, 744 blockBytesScanned); 745 } 746 return r == null ? Result.EMPTY_RESULT : r; 747 } 748 749 /** 750 * Run through the regionMutation <code>rm</code> and per Mutation, do the work, and then when 751 * done, add an instance of a {@link ResultOrException} that corresponds to each Mutation. 752 * @param cellsToReturn Could be null. May be allocated in this method. This is what this method 753 * returns as a 'result'. 754 * @param closeCallBack the callback to be used with multigets 755 * @param context the current RpcCallContext 756 * @return Return the <code>cellScanner</code> passed 757 */ 758 private List<ExtendedCellScannable> doNonAtomicRegionMutation(final HRegion region, 759 final OperationQuota quota, final RegionAction actions, final CellScanner cellScanner, 760 final RegionActionResult.Builder builder, List<ExtendedCellScannable> cellsToReturn, 761 long nonceGroup, final RegionScannersCloseCallBack closeCallBack, RpcCallContext context, 762 ActivePolicyEnforcement spaceQuotaEnforcement) { 763 // Gather up CONTIGUOUS Puts and Deletes in this mutations List. Idea is that rather than do 764 // one at a time, we instead pass them in batch. Be aware that the corresponding 765 // ResultOrException instance that matches each Put or Delete is then added down in the 766 // doNonAtomicBatchOp call. We should be staying aligned though the Put and Delete are 767 // deferred/batched 768 List<ClientProtos.Action> mutations = null; 769 long maxQuotaResultSize = Math.min(maxScannerResultSize, quota.getMaxResultSize()); 770 IOException sizeIOE = null; 771 ClientProtos.ResultOrException.Builder resultOrExceptionBuilder = 772 ResultOrException.newBuilder(); 773 boolean hasResultOrException = false; 774 for (ClientProtos.Action action : actions.getActionList()) { 775 hasResultOrException = false; 776 resultOrExceptionBuilder.clear(); 777 try { 778 Result r = null; 779 long blockBytesScannedBefore = context != null ? context.getBlockBytesScanned() : 0; 780 if ( 781 context != null && context.isRetryImmediatelySupported() 782 && (context.getResponseCellSize() > maxQuotaResultSize 783 || blockBytesScannedBefore + context.getResponseExceptionSize() > maxQuotaResultSize) 784 ) { 785 786 // We're storing the exception since the exception and reason string won't 787 // change after the response size limit is reached. 788 if (sizeIOE == null) { 789 // We don't need the stack un-winding do don't throw the exception. 790 // Throwing will kill the JVM's JIT. 791 // 792 // Instead just create the exception and then store it. 793 sizeIOE = new MultiActionResultTooLarge("Max size exceeded" + " CellSize: " 794 + context.getResponseCellSize() + " BlockSize: " + blockBytesScannedBefore); 795 796 // Only report the exception once since there's only one request that 797 // caused the exception. Otherwise this number will dominate the exceptions count. 798 rpcServer.getMetrics().exception(sizeIOE); 799 } 800 801 // Now that there's an exception is known to be created 802 // use it for the response. 803 // 804 // This will create a copy in the builder. 805 NameBytesPair pair = ResponseConverter.buildException(sizeIOE); 806 resultOrExceptionBuilder.setException(pair); 807 context.incrementResponseExceptionSize(pair.getSerializedSize()); 808 resultOrExceptionBuilder.setIndex(action.getIndex()); 809 builder.addResultOrException(resultOrExceptionBuilder.build()); 810 skipCellsForMutation(action, cellScanner); 811 continue; 812 } 813 if (action.hasGet()) { 814 long before = EnvironmentEdgeManager.currentTime(); 815 ClientProtos.Get pbGet = action.getGet(); 816 // An asynchbase client, https://github.com/OpenTSDB/asynchbase, starts by trying to do 817 // a get closest before. Throwing the UnknownProtocolException signals it that it needs 818 // to switch and do hbase2 protocol (HBase servers do not tell clients what versions 819 // they are; its a problem for non-native clients like asynchbase. HBASE-20225. 820 if (pbGet.hasClosestRowBefore() && pbGet.getClosestRowBefore()) { 821 throw new UnknownProtocolException("Is this a pre-hbase-1.0.0 or asynchbase client? " 822 + "Client is invoking getClosestRowBefore removed in hbase-2.0.0 replaced by " 823 + "reverse Scan."); 824 } 825 try { 826 Get get = ProtobufUtil.toGet(pbGet); 827 if (context != null) { 828 r = get(get, (region), closeCallBack, context); 829 } else { 830 r = region.get(get); 831 } 832 } finally { 833 final MetricsRegionServer metricsRegionServer = server.getMetrics(); 834 if (metricsRegionServer != null) { 835 long blockBytesScanned = 836 context != null ? context.getBlockBytesScanned() - blockBytesScannedBefore : 0; 837 metricsRegionServer.updateGet(region, EnvironmentEdgeManager.currentTime() - before, 838 blockBytesScanned); 839 } 840 } 841 } else if (action.hasServiceCall()) { 842 hasResultOrException = true; 843 Message result = execServiceOnRegion(region, action.getServiceCall()); 844 ClientProtos.CoprocessorServiceResult.Builder serviceResultBuilder = 845 ClientProtos.CoprocessorServiceResult.newBuilder(); 846 resultOrExceptionBuilder.setServiceResult(serviceResultBuilder 847 .setValue(serviceResultBuilder.getValueBuilder().setName(result.getClass().getName()) 848 // TODO: Copy!!! 849 .setValue(UnsafeByteOperations.unsafeWrap(result.toByteArray())))); 850 } else if (action.hasMutation()) { 851 MutationType type = action.getMutation().getMutateType(); 852 if ( 853 type != MutationType.PUT && type != MutationType.DELETE && mutations != null 854 && !mutations.isEmpty() 855 ) { 856 // Flush out any Puts or Deletes already collected. 857 doNonAtomicBatchOp(builder, region, quota, mutations, cellScanner, 858 spaceQuotaEnforcement); 859 mutations.clear(); 860 } 861 switch (type) { 862 case APPEND: 863 r = append(region, quota, action.getMutation(), cellScanner, nonceGroup, 864 spaceQuotaEnforcement, context); 865 break; 866 case INCREMENT: 867 r = increment(region, quota, action.getMutation(), cellScanner, nonceGroup, 868 spaceQuotaEnforcement, context); 869 break; 870 case PUT: 871 case DELETE: 872 // Collect the individual mutations and apply in a batch 873 if (mutations == null) { 874 mutations = new ArrayList<>(actions.getActionCount()); 875 } 876 mutations.add(action); 877 break; 878 default: 879 throw new DoNotRetryIOException("Unsupported mutate type: " + type.name()); 880 } 881 } else { 882 throw new HBaseIOException("Unexpected Action type"); 883 } 884 if (r != null) { 885 ClientProtos.Result pbResult = null; 886 if (isClientCellBlockSupport(context)) { 887 pbResult = ProtobufUtil.toResultNoData(r); 888 // Hard to guess the size here. Just make a rough guess. 889 if (cellsToReturn == null) { 890 cellsToReturn = new ArrayList<>(); 891 } 892 cellsToReturn.add(r); 893 } else { 894 pbResult = ProtobufUtil.toResult(r); 895 } 896 addSize(context, r); 897 hasResultOrException = true; 898 resultOrExceptionBuilder.setResult(pbResult); 899 } 900 // Could get to here and there was no result and no exception. Presumes we added 901 // a Put or Delete to the collecting Mutations List for adding later. In this 902 // case the corresponding ResultOrException instance for the Put or Delete will be added 903 // down in the doNonAtomicBatchOp method call rather than up here. 904 } catch (IOException ie) { 905 rpcServer.getMetrics().exception(ie); 906 hasResultOrException = true; 907 NameBytesPair pair = ResponseConverter.buildException(ie); 908 resultOrExceptionBuilder.setException(pair); 909 context.incrementResponseExceptionSize(pair.getSerializedSize()); 910 } 911 if (hasResultOrException) { 912 // Propagate index. 913 resultOrExceptionBuilder.setIndex(action.getIndex()); 914 builder.addResultOrException(resultOrExceptionBuilder.build()); 915 } 916 } 917 // Finish up any outstanding mutations 918 if (!CollectionUtils.isEmpty(mutations)) { 919 doNonAtomicBatchOp(builder, region, quota, mutations, cellScanner, spaceQuotaEnforcement); 920 } 921 return cellsToReturn; 922 } 923 924 private void checkCellSizeLimit(final HRegion r, final Mutation m) throws IOException { 925 if (r.maxCellSize > 0) { 926 CellScanner cells = m.cellScanner(); 927 while (cells.advance()) { 928 int size = PrivateCellUtil.estimatedSerializedSizeOf(cells.current()); 929 if (size > r.maxCellSize) { 930 String msg = "Cell[" + cells.current() + "] with size " + size + " exceeds limit of " 931 + r.maxCellSize + " bytes"; 932 LOG.debug(msg); 933 throw new DoNotRetryIOException(msg); 934 } 935 } 936 } 937 } 938 939 private void doAtomicBatchOp(final RegionActionResult.Builder builder, final HRegion region, 940 final OperationQuota quota, final List<ClientProtos.Action> mutations, final CellScanner cells, 941 long nonceGroup, ActivePolicyEnforcement spaceQuotaEnforcement) throws IOException { 942 // Just throw the exception. The exception will be caught and then added to region-level 943 // exception for RegionAction. Leaving the null to action result is ok since the null 944 // result is viewed as failure by hbase client. And the region-lever exception will be used 945 // to replaced the null result. see AsyncRequestFutureImpl#receiveMultiAction and 946 // AsyncBatchRpcRetryingCaller#onComplete for more details. 947 doBatchOp(builder, region, quota, mutations, cells, nonceGroup, spaceQuotaEnforcement, true); 948 } 949 950 private void doNonAtomicBatchOp(final RegionActionResult.Builder builder, final HRegion region, 951 final OperationQuota quota, final List<ClientProtos.Action> mutations, final CellScanner cells, 952 ActivePolicyEnforcement spaceQuotaEnforcement) { 953 try { 954 doBatchOp(builder, region, quota, mutations, cells, HConstants.NO_NONCE, 955 spaceQuotaEnforcement, false); 956 } catch (IOException e) { 957 // Set the exception for each action. The mutations in same RegionAction are group to 958 // different batch and then be processed individually. Hence, we don't set the region-level 959 // exception here for whole RegionAction. 960 for (Action mutation : mutations) { 961 builder.addResultOrException(getResultOrException(e, mutation.getIndex())); 962 } 963 } 964 } 965 966 /** 967 * Execute a list of mutations. 968 */ 969 private void doBatchOp(final RegionActionResult.Builder builder, final HRegion region, 970 final OperationQuota quota, final List<ClientProtos.Action> mutations, final CellScanner cells, 971 long nonceGroup, ActivePolicyEnforcement spaceQuotaEnforcement, boolean atomic) 972 throws IOException { 973 Mutation[] mArray = new Mutation[mutations.size()]; 974 long before = EnvironmentEdgeManager.currentTime(); 975 boolean batchContainsPuts = false, batchContainsDelete = false; 976 try { 977 /** 978 * HBASE-17924 mutationActionMap is a map to map the relation between mutations and actions 979 * since mutation array may have been reoredered.In order to return the right result or 980 * exception to the corresponding actions, We need to know which action is the mutation belong 981 * to. We can't sort ClientProtos.Action array, since they are bonded to cellscanners. 982 */ 983 Map<Mutation, ClientProtos.Action> mutationActionMap = new HashMap<>(); 984 int i = 0; 985 long nonce = HConstants.NO_NONCE; 986 for (ClientProtos.Action action : mutations) { 987 if (action.hasGet()) { 988 throw new DoNotRetryIOException( 989 "Atomic put and/or delete only, not a Get=" + action.getGet()); 990 } 991 MutationProto m = action.getMutation(); 992 Mutation mutation; 993 switch (m.getMutateType()) { 994 case PUT: 995 mutation = ProtobufUtil.toPut(m, cells); 996 batchContainsPuts = true; 997 break; 998 999 case DELETE: 1000 mutation = ProtobufUtil.toDelete(m, cells); 1001 batchContainsDelete = true; 1002 break; 1003 1004 case INCREMENT: 1005 mutation = ProtobufUtil.toIncrement(m, cells); 1006 nonce = m.hasNonce() ? m.getNonce() : HConstants.NO_NONCE; 1007 break; 1008 1009 case APPEND: 1010 mutation = ProtobufUtil.toAppend(m, cells); 1011 nonce = m.hasNonce() ? m.getNonce() : HConstants.NO_NONCE; 1012 break; 1013 1014 default: 1015 throw new DoNotRetryIOException("Invalid mutation type : " + m.getMutateType()); 1016 } 1017 mutationActionMap.put(mutation, action); 1018 mArray[i++] = mutation; 1019 checkCellSizeLimit(region, mutation); 1020 // Check if a space quota disallows this mutation 1021 spaceQuotaEnforcement.getPolicyEnforcement(region).check(mutation); 1022 quota.addMutation(mutation); 1023 } 1024 1025 if (!region.getRegionInfo().isMetaRegion()) { 1026 server.getMemStoreFlusher().reclaimMemStoreMemory(); 1027 } 1028 1029 // HBASE-17924 1030 // Sort to improve lock efficiency for non-atomic batch of operations. If atomic 1031 // order is preserved as its expected from the client 1032 if (!atomic) { 1033 Arrays.sort(mArray, (v1, v2) -> Row.COMPARATOR.compare(v1, v2)); 1034 } 1035 1036 OperationStatus[] codes = region.batchMutate(mArray, atomic, nonceGroup, nonce); 1037 1038 // When atomic is true, it indicates that the mutateRow API or the batch API with 1039 // RowMutations is called. In this case, we need to merge the results of the 1040 // Increment/Append operations if the mutations include those operations, and set the merged 1041 // result to the first element of the ResultOrException list 1042 if (atomic) { 1043 List<ResultOrException> resultOrExceptions = new ArrayList<>(); 1044 List<Result> results = new ArrayList<>(); 1045 for (i = 0; i < codes.length; i++) { 1046 if (codes[i].getResult() != null) { 1047 results.add(codes[i].getResult()); 1048 } 1049 if (i != 0) { 1050 resultOrExceptions 1051 .add(getResultOrException(ClientProtos.Result.getDefaultInstance(), i)); 1052 } 1053 } 1054 1055 if (results.isEmpty()) { 1056 builder.addResultOrException( 1057 getResultOrException(ClientProtos.Result.getDefaultInstance(), 0)); 1058 } else { 1059 // Merge the results of the Increment/Append operations 1060 List<Cell> cellList = new ArrayList<>(); 1061 for (Result result : results) { 1062 if (result.rawCells() != null) { 1063 cellList.addAll(Arrays.asList(result.rawCells())); 1064 } 1065 } 1066 Result result = Result.create(cellList); 1067 1068 // Set the merged result of the Increment/Append operations to the first element of the 1069 // ResultOrException list 1070 builder.addResultOrException(getResultOrException(ProtobufUtil.toResult(result), 0)); 1071 } 1072 1073 builder.addAllResultOrException(resultOrExceptions); 1074 return; 1075 } 1076 1077 for (i = 0; i < codes.length; i++) { 1078 Mutation currentMutation = mArray[i]; 1079 ClientProtos.Action currentAction = mutationActionMap.get(currentMutation); 1080 int index = currentAction.hasIndex() ? currentAction.getIndex() : i; 1081 Exception e; 1082 switch (codes[i].getOperationStatusCode()) { 1083 case BAD_FAMILY: 1084 e = new NoSuchColumnFamilyException(codes[i].getExceptionMsg()); 1085 builder.addResultOrException(getResultOrException(e, index)); 1086 break; 1087 1088 case SANITY_CHECK_FAILURE: 1089 e = new FailedSanityCheckException(codes[i].getExceptionMsg()); 1090 builder.addResultOrException(getResultOrException(e, index)); 1091 break; 1092 1093 default: 1094 e = new DoNotRetryIOException(codes[i].getExceptionMsg()); 1095 builder.addResultOrException(getResultOrException(e, index)); 1096 break; 1097 1098 case SUCCESS: 1099 ClientProtos.Result result = codes[i].getResult() == null 1100 ? ClientProtos.Result.getDefaultInstance() 1101 : ProtobufUtil.toResult(codes[i].getResult()); 1102 builder.addResultOrException(getResultOrException(result, index)); 1103 break; 1104 1105 case STORE_TOO_BUSY: 1106 e = new RegionTooBusyException(codes[i].getExceptionMsg()); 1107 builder.addResultOrException(getResultOrException(e, index)); 1108 break; 1109 } 1110 } 1111 } finally { 1112 int processedMutationIndex = 0; 1113 for (Action mutation : mutations) { 1114 // The non-null mArray[i] means the cell scanner has been read. 1115 if (mArray[processedMutationIndex++] == null) { 1116 skipCellsForMutation(mutation, cells); 1117 } 1118 } 1119 updateMutationMetrics(region, before, batchContainsPuts, batchContainsDelete); 1120 } 1121 } 1122 1123 private void updateMutationMetrics(HRegion region, long starttime, boolean batchContainsPuts, 1124 boolean batchContainsDelete) { 1125 final MetricsRegionServer metricsRegionServer = server.getMetrics(); 1126 if (metricsRegionServer != null) { 1127 long after = EnvironmentEdgeManager.currentTime(); 1128 if (batchContainsPuts) { 1129 metricsRegionServer.updatePutBatch(region, after - starttime); 1130 } 1131 if (batchContainsDelete) { 1132 metricsRegionServer.updateDeleteBatch(region, after - starttime); 1133 } 1134 } 1135 } 1136 1137 /** 1138 * Execute a list of Put/Delete mutations. The function returns OperationStatus instead of 1139 * constructing MultiResponse to save a possible loop if caller doesn't need MultiResponse. 1140 * @return an array of OperationStatus which internally contains the OperationStatusCode and the 1141 * exceptionMessage if any 1142 * @deprecated Since 3.0.0, will be removed in 4.0.0. We do not use this method for replaying 1143 * edits for secondary replicas any more, see 1144 * {@link #replicateToReplica(RpcController, ReplicateWALEntryRequest)}. 1145 */ 1146 @Deprecated 1147 private OperationStatus[] doReplayBatchOp(final HRegion region, 1148 final List<MutationReplay> mutations, long replaySeqId) throws IOException { 1149 long before = EnvironmentEdgeManager.currentTime(); 1150 boolean batchContainsPuts = false, batchContainsDelete = false; 1151 try { 1152 for (Iterator<MutationReplay> it = mutations.iterator(); it.hasNext();) { 1153 MutationReplay m = it.next(); 1154 1155 if (m.getType() == MutationType.PUT) { 1156 batchContainsPuts = true; 1157 } else { 1158 batchContainsDelete = true; 1159 } 1160 1161 NavigableMap<byte[], List<Cell>> map = m.mutation.getFamilyCellMap(); 1162 List<Cell> metaCells = map.get(WALEdit.METAFAMILY); 1163 if (metaCells != null && !metaCells.isEmpty()) { 1164 for (Cell metaCell : metaCells) { 1165 CompactionDescriptor compactionDesc = WALEdit.getCompaction(metaCell); 1166 boolean isDefaultReplica = RegionReplicaUtil.isDefaultReplica(region.getRegionInfo()); 1167 HRegion hRegion = region; 1168 if (compactionDesc != null) { 1169 // replay the compaction. Remove the files from stores only if we are the primary 1170 // region replica (thus own the files) 1171 hRegion.replayWALCompactionMarker(compactionDesc, !isDefaultReplica, isDefaultReplica, 1172 replaySeqId); 1173 continue; 1174 } 1175 FlushDescriptor flushDesc = WALEdit.getFlushDescriptor(metaCell); 1176 if (flushDesc != null && !isDefaultReplica) { 1177 hRegion.replayWALFlushMarker(flushDesc, replaySeqId); 1178 continue; 1179 } 1180 RegionEventDescriptor regionEvent = WALEdit.getRegionEventDescriptor(metaCell); 1181 if (regionEvent != null && !isDefaultReplica) { 1182 hRegion.replayWALRegionEventMarker(regionEvent); 1183 continue; 1184 } 1185 BulkLoadDescriptor bulkLoadEvent = WALEdit.getBulkLoadDescriptor(metaCell); 1186 if (bulkLoadEvent != null) { 1187 hRegion.replayWALBulkLoadEventMarker(bulkLoadEvent); 1188 continue; 1189 } 1190 } 1191 it.remove(); 1192 } 1193 } 1194 requestCount.increment(); 1195 if (!region.getRegionInfo().isMetaRegion()) { 1196 server.getMemStoreFlusher().reclaimMemStoreMemory(); 1197 } 1198 return region.batchReplay(mutations.toArray(new MutationReplay[mutations.size()]), 1199 replaySeqId); 1200 } finally { 1201 updateMutationMetrics(region, before, batchContainsPuts, batchContainsDelete); 1202 } 1203 } 1204 1205 private void closeAllScanners() { 1206 // Close any outstanding scanners. Means they'll get an UnknownScanner 1207 // exception next time they come in. 1208 for (Map.Entry<String, RegionScannerHolder> e : scanners.entrySet()) { 1209 try { 1210 e.getValue().s.close(); 1211 } catch (IOException ioe) { 1212 LOG.warn("Closing scanner " + e.getKey(), ioe); 1213 } 1214 } 1215 } 1216 1217 // Directly invoked only for testing 1218 public RSRpcServices(final HRegionServer rs) throws IOException { 1219 super(rs, rs.getProcessName()); 1220 final Configuration conf = rs.getConfiguration(); 1221 setReloadableGuardrails(conf); 1222 scannerLeaseTimeoutPeriod = conf.getInt(HConstants.HBASE_CLIENT_SCANNER_TIMEOUT_PERIOD, 1223 HConstants.DEFAULT_HBASE_CLIENT_SCANNER_TIMEOUT_PERIOD); 1224 rpcTimeout = 1225 conf.getInt(HConstants.HBASE_RPC_TIMEOUT_KEY, HConstants.DEFAULT_HBASE_RPC_TIMEOUT); 1226 minimumScanTimeLimitDelta = conf.getLong(REGION_SERVER_RPC_MINIMUM_SCAN_TIME_LIMIT_DELTA, 1227 DEFAULT_REGION_SERVER_RPC_MINIMUM_SCAN_TIME_LIMIT_DELTA); 1228 rpcServer.setNamedQueueRecorder(rs.getNamedQueueRecorder()); 1229 closedScanners = CacheBuilder.newBuilder() 1230 .expireAfterAccess(scannerLeaseTimeoutPeriod, TimeUnit.MILLISECONDS).build(); 1231 } 1232 1233 @Override 1234 protected boolean defaultReservoirEnabled() { 1235 return true; 1236 } 1237 1238 @Override 1239 protected ServerType getDNSServerType() { 1240 return DNS.ServerType.REGIONSERVER; 1241 } 1242 1243 @Override 1244 protected String getHostname(Configuration conf, String defaultHostname) { 1245 return conf.get("hbase.regionserver.ipc.address", defaultHostname); 1246 } 1247 1248 @Override 1249 protected String getPortConfigName() { 1250 return HConstants.REGIONSERVER_PORT; 1251 } 1252 1253 @Override 1254 protected int getDefaultPort() { 1255 return HConstants.DEFAULT_REGIONSERVER_PORT; 1256 } 1257 1258 @Override 1259 protected Class<?> getRpcSchedulerFactoryClass(Configuration conf) { 1260 return conf.getClass(REGION_SERVER_RPC_SCHEDULER_FACTORY_CLASS, 1261 SimpleRpcSchedulerFactory.class); 1262 } 1263 1264 protected RpcServerInterface createRpcServer(final Server server, 1265 final RpcSchedulerFactory rpcSchedulerFactory, final InetSocketAddress bindAddress, 1266 final String name) throws IOException { 1267 final Configuration conf = server.getConfiguration(); 1268 boolean reservoirEnabled = conf.getBoolean(ByteBuffAllocator.ALLOCATOR_POOL_ENABLED_KEY, true); 1269 try { 1270 return RpcServerFactory.createRpcServer(server, name, getServices(), bindAddress, // use final 1271 // bindAddress 1272 // for this 1273 // server. 1274 conf, rpcSchedulerFactory.create(conf, this, server), reservoirEnabled); 1275 } catch (BindException be) { 1276 throw new IOException(be.getMessage() + ". To switch ports use the '" 1277 + HConstants.REGIONSERVER_PORT + "' configuration property.", 1278 be.getCause() != null ? be.getCause() : be); 1279 } 1280 } 1281 1282 protected Class<?> getRpcSchedulerFactoryClass() { 1283 final Configuration conf = server.getConfiguration(); 1284 return conf.getClass(REGION_SERVER_RPC_SCHEDULER_FACTORY_CLASS, 1285 SimpleRpcSchedulerFactory.class); 1286 } 1287 1288 protected PriorityFunction createPriority() { 1289 return new RSAnnotationReadingPriorityFunction(this); 1290 } 1291 1292 public int getScannersCount() { 1293 return scanners.size(); 1294 } 1295 1296 /** Returns The outstanding RegionScanner for <code>scannerId</code> or null if none found. */ 1297 RegionScanner getScanner(long scannerId) { 1298 RegionScannerHolder rsh = checkQuotaAndGetRegionScannerContext(scannerId); 1299 return rsh == null ? null : rsh.s; 1300 } 1301 1302 /** Returns The associated RegionScannerHolder for <code>scannerId</code> or null. */ 1303 private RegionScannerHolder checkQuotaAndGetRegionScannerContext(long scannerId) { 1304 return scanners.get(toScannerName(scannerId)); 1305 } 1306 1307 public String getScanDetailsWithId(long scannerId) { 1308 RegionScanner scanner = getScanner(scannerId); 1309 if (scanner == null) { 1310 return null; 1311 } 1312 StringBuilder builder = new StringBuilder(); 1313 builder.append("table: ").append(scanner.getRegionInfo().getTable().getNameAsString()); 1314 builder.append(" region: ").append(scanner.getRegionInfo().getRegionNameAsString()); 1315 builder.append(" operation_id: ").append(scanner.getOperationId()); 1316 return builder.toString(); 1317 } 1318 1319 public String getScanDetailsWithRequest(ScanRequest request) { 1320 try { 1321 if (!request.hasRegion()) { 1322 return null; 1323 } 1324 Region region = getRegion(request.getRegion()); 1325 StringBuilder builder = new StringBuilder(); 1326 builder.append("table: ").append(region.getRegionInfo().getTable().getNameAsString()); 1327 builder.append(" region: ").append(region.getRegionInfo().getRegionNameAsString()); 1328 for (NameBytesPair pair : request.getScan().getAttributeList()) { 1329 if (OperationWithAttributes.ID_ATRIBUTE.equals(pair.getName())) { 1330 builder.append(" operation_id: ").append(Bytes.toString(pair.getValue().toByteArray())); 1331 break; 1332 } 1333 } 1334 return builder.toString(); 1335 } catch (IOException ignored) { 1336 return null; 1337 } 1338 } 1339 1340 /** 1341 * Get the vtime associated with the scanner. Currently the vtime is the number of "next" calls. 1342 */ 1343 long getScannerVirtualTime(long scannerId) { 1344 RegionScannerHolder rsh = checkQuotaAndGetRegionScannerContext(scannerId); 1345 return rsh == null ? 0L : rsh.getNextCallSeq(); 1346 } 1347 1348 /** 1349 * Method to account for the size of retained cells. 1350 * @param context rpc call context 1351 * @param r result to add size. 1352 * @return an object that represents the last referenced block from this response. 1353 */ 1354 void addSize(RpcCallContext context, Result r) { 1355 if (context != null && r != null && !r.isEmpty()) { 1356 for (Cell c : r.rawCells()) { 1357 context.incrementResponseCellSize(PrivateCellUtil.estimatedSerializedSizeOf(c)); 1358 } 1359 } 1360 } 1361 1362 /** Returns Remote client's ip and port else null if can't be determined. */ 1363 @RestrictedApi(explanation = "Should only be called in TestRSRpcServices and RSRpcServices", 1364 link = "", allowedOnPath = ".*(TestRSRpcServices|RSRpcServices).java") 1365 static String getRemoteClientIpAndPort() { 1366 RpcCall rpcCall = RpcServer.getCurrentCall().orElse(null); 1367 if (rpcCall == null) { 1368 return HConstants.EMPTY_STRING; 1369 } 1370 InetAddress address = rpcCall.getRemoteAddress(); 1371 if (address == null) { 1372 return HConstants.EMPTY_STRING; 1373 } 1374 // Be careful here with InetAddress. Do InetAddress#getHostAddress. It will not do a name 1375 // resolution. Just use the IP. It is generally a smaller amount of info to keep around while 1376 // scanning than a hostname anyways. 1377 return Address.fromParts(address.getHostAddress(), rpcCall.getRemotePort()).toString(); 1378 } 1379 1380 /** Returns Remote client's username. */ 1381 @RestrictedApi(explanation = "Should only be called in TestRSRpcServices and RSRpcServices", 1382 link = "", allowedOnPath = ".*(TestRSRpcServices|RSRpcServices).java") 1383 static String getUserName() { 1384 RpcCall rpcCall = RpcServer.getCurrentCall().orElse(null); 1385 if (rpcCall == null) { 1386 return HConstants.EMPTY_STRING; 1387 } 1388 return rpcCall.getRequestUserName().orElse(HConstants.EMPTY_STRING); 1389 } 1390 1391 private RegionScannerHolder addScanner(String scannerName, RegionScanner s, Shipper shipper, 1392 HRegion r, boolean needCursor, boolean fullRegionScan) throws LeaseStillHeldException { 1393 Lease lease = server.getLeaseManager().createLease(scannerName, this.scannerLeaseTimeoutPeriod, 1394 new ScannerListener(scannerName)); 1395 RpcCallback shippedCallback = new RegionScannerShippedCallBack(scannerName, shipper, lease); 1396 RpcCallback closeCallback = 1397 s instanceof RpcCallback ? (RpcCallback) s : new RegionScannerCloseCallBack(s); 1398 RegionScannerHolder rsh = new RegionScannerHolder(s, r, closeCallback, shippedCallback, 1399 needCursor, fullRegionScan, getRemoteClientIpAndPort(), getUserName()); 1400 RegionScannerHolder existing = scanners.putIfAbsent(scannerName, rsh); 1401 assert existing == null : "scannerId must be unique within regionserver's whole lifecycle! " 1402 + scannerName + ", " + existing; 1403 return rsh; 1404 } 1405 1406 private boolean isFullRegionScan(Scan scan, HRegion region) { 1407 // If the scan start row equals or less than the start key of the region 1408 // and stop row greater than equals end key (if stop row present) 1409 // or if the stop row is empty 1410 // account this as a full region scan 1411 if ( 1412 Bytes.compareTo(scan.getStartRow(), region.getRegionInfo().getStartKey()) <= 0 1413 && (Bytes.compareTo(scan.getStopRow(), region.getRegionInfo().getEndKey()) >= 0 1414 && !Bytes.equals(region.getRegionInfo().getEndKey(), HConstants.EMPTY_END_ROW) 1415 || Bytes.equals(scan.getStopRow(), HConstants.EMPTY_END_ROW)) 1416 ) { 1417 return true; 1418 } 1419 return false; 1420 } 1421 1422 /** 1423 * Find the HRegion based on a region specifier 1424 * @param regionSpecifier the region specifier 1425 * @return the corresponding region 1426 * @throws IOException if the specifier is not null, but failed to find the region 1427 */ 1428 public HRegion getRegion(final RegionSpecifier regionSpecifier) throws IOException { 1429 return server.getRegion(regionSpecifier.getValue().toByteArray()); 1430 } 1431 1432 /** 1433 * Find the List of HRegions based on a list of region specifiers 1434 * @param regionSpecifiers the list of region specifiers 1435 * @return the corresponding list of regions 1436 * @throws IOException if any of the specifiers is not null, but failed to find the region 1437 */ 1438 private List<HRegion> getRegions(final List<RegionSpecifier> regionSpecifiers, 1439 final CacheEvictionStatsBuilder stats) { 1440 List<HRegion> regions = Lists.newArrayListWithCapacity(regionSpecifiers.size()); 1441 for (RegionSpecifier regionSpecifier : regionSpecifiers) { 1442 try { 1443 regions.add(server.getRegion(regionSpecifier.getValue().toByteArray())); 1444 } catch (NotServingRegionException e) { 1445 stats.addException(regionSpecifier.getValue().toByteArray(), e); 1446 } 1447 } 1448 return regions; 1449 } 1450 1451 PriorityFunction getPriority() { 1452 return priority; 1453 } 1454 1455 private RegionServerRpcQuotaManager getRpcQuotaManager() { 1456 return server.getRegionServerRpcQuotaManager(); 1457 } 1458 1459 private RegionServerSpaceQuotaManager getSpaceQuotaManager() { 1460 return server.getRegionServerSpaceQuotaManager(); 1461 } 1462 1463 void start(ZKWatcher zkWatcher) { 1464 this.scannerIdGenerator = new ScannerIdGenerator(this.server.getServerName()); 1465 internalStart(zkWatcher); 1466 } 1467 1468 void stop() { 1469 closeAllScanners(); 1470 internalStop(); 1471 } 1472 1473 /** 1474 * Called to verify that this server is up and running. 1475 */ 1476 // TODO : Rename this and HMaster#checkInitialized to isRunning() (or a better name). 1477 protected void checkOpen() throws IOException { 1478 if (server.isAborted()) { 1479 throw new RegionServerAbortedException("Server " + server.getServerName() + " aborting"); 1480 } 1481 if (server.isStopped()) { 1482 throw new RegionServerStoppedException("Server " + server.getServerName() + " stopping"); 1483 } 1484 if (!server.isDataFileSystemOk()) { 1485 throw new RegionServerStoppedException("File system not available"); 1486 } 1487 if (!server.isOnline()) { 1488 throw new ServerNotRunningYetException( 1489 "Server " + server.getServerName() + " is not running yet"); 1490 } 1491 } 1492 1493 /** 1494 * By default, put up an Admin and a Client Service. Set booleans 1495 * <code>hbase.regionserver.admin.executorService</code> and 1496 * <code>hbase.regionserver.client.executorService</code> if you want to enable/disable services. 1497 * Default is that both are enabled. 1498 * @return immutable list of blocking services and the security info classes that this server 1499 * supports 1500 */ 1501 protected List<BlockingServiceAndInterface> getServices() { 1502 boolean admin = getConfiguration().getBoolean(REGIONSERVER_ADMIN_SERVICE_CONFIG, true); 1503 boolean client = getConfiguration().getBoolean(REGIONSERVER_CLIENT_SERVICE_CONFIG, true); 1504 boolean clientMeta = 1505 getConfiguration().getBoolean(REGIONSERVER_CLIENT_META_SERVICE_CONFIG, true); 1506 boolean bootstrapNodes = 1507 getConfiguration().getBoolean(REGIONSERVER_BOOTSTRAP_NODES_SERVICE_CONFIG, true); 1508 List<BlockingServiceAndInterface> bssi = new ArrayList<>(); 1509 if (client) { 1510 bssi.add(new BlockingServiceAndInterface(ClientService.newReflectiveBlockingService(this), 1511 ClientService.BlockingInterface.class)); 1512 } 1513 if (admin) { 1514 bssi.add(new BlockingServiceAndInterface(AdminService.newReflectiveBlockingService(this), 1515 AdminService.BlockingInterface.class)); 1516 } 1517 if (clientMeta) { 1518 bssi.add(new BlockingServiceAndInterface(ClientMetaService.newReflectiveBlockingService(this), 1519 ClientMetaService.BlockingInterface.class)); 1520 } 1521 if (bootstrapNodes) { 1522 bssi.add( 1523 new BlockingServiceAndInterface(BootstrapNodeService.newReflectiveBlockingService(this), 1524 BootstrapNodeService.BlockingInterface.class)); 1525 } 1526 return new ImmutableList.Builder<BlockingServiceAndInterface>().addAll(bssi).build(); 1527 } 1528 1529 /** 1530 * Close a region on the region server. 1531 * @param controller the RPC controller 1532 * @param request the request 1533 */ 1534 @Override 1535 @QosPriority(priority = HConstants.ADMIN_QOS) 1536 public CloseRegionResponse closeRegion(final RpcController controller, 1537 final CloseRegionRequest request) throws ServiceException { 1538 final ServerName sn = (request.hasDestinationServer() 1539 ? ProtobufUtil.toServerName(request.getDestinationServer()) 1540 : null); 1541 1542 try { 1543 checkOpen(); 1544 throwOnWrongStartCode(request); 1545 final String encodedRegionName = ProtobufUtil.getRegionEncodedName(request.getRegion()); 1546 1547 requestCount.increment(); 1548 if (sn == null) { 1549 LOG.info("Close " + encodedRegionName + " without moving"); 1550 } else { 1551 LOG.info("Close " + encodedRegionName + ", moving to " + sn); 1552 } 1553 boolean closed = server.closeRegion(encodedRegionName, false, sn); 1554 CloseRegionResponse.Builder builder = CloseRegionResponse.newBuilder().setClosed(closed); 1555 return builder.build(); 1556 } catch (IOException ie) { 1557 throw new ServiceException(ie); 1558 } 1559 } 1560 1561 /** 1562 * Compact a region on the region server. 1563 * @param controller the RPC controller 1564 * @param request the request 1565 */ 1566 @Override 1567 @QosPriority(priority = HConstants.ADMIN_QOS) 1568 public CompactRegionResponse compactRegion(final RpcController controller, 1569 final CompactRegionRequest request) throws ServiceException { 1570 try { 1571 checkOpen(); 1572 requestCount.increment(); 1573 HRegion region = getRegion(request.getRegion()); 1574 // Quota support is enabled, the requesting user is not system/super user 1575 // and a quota policy is enforced that disables compactions. 1576 if ( 1577 QuotaUtil.isQuotaEnabled(getConfiguration()) 1578 && !Superusers.isSuperUser(RpcServer.getRequestUser().orElse(null)) 1579 && this.server.getRegionServerSpaceQuotaManager() 1580 .areCompactionsDisabled(region.getTableDescriptor().getTableName()) 1581 ) { 1582 throw new DoNotRetryIOException( 1583 "Compactions on this region are " + "disabled due to a space quota violation."); 1584 } 1585 region.startRegionOperation(Operation.COMPACT_REGION); 1586 LOG.info("Compacting " + region.getRegionInfo().getRegionNameAsString()); 1587 boolean major = request.hasMajor() && request.getMajor(); 1588 if (request.hasFamily()) { 1589 byte[] family = request.getFamily().toByteArray(); 1590 String log = "User-triggered " + (major ? "major " : "") + "compaction for region " 1591 + region.getRegionInfo().getRegionNameAsString() + " and family " 1592 + Bytes.toString(family); 1593 LOG.trace(log); 1594 region.requestCompaction(family, log, Store.PRIORITY_USER, major, 1595 CompactionLifeCycleTracker.DUMMY); 1596 } else { 1597 String log = "User-triggered " + (major ? "major " : "") + "compaction for region " 1598 + region.getRegionInfo().getRegionNameAsString(); 1599 LOG.trace(log); 1600 region.requestCompaction(log, Store.PRIORITY_USER, major, CompactionLifeCycleTracker.DUMMY); 1601 } 1602 return CompactRegionResponse.newBuilder().build(); 1603 } catch (IOException ie) { 1604 throw new ServiceException(ie); 1605 } 1606 } 1607 1608 @Override 1609 @QosPriority(priority = HConstants.ADMIN_QOS) 1610 public CompactionSwitchResponse compactionSwitch(RpcController controller, 1611 CompactionSwitchRequest request) throws ServiceException { 1612 rpcPreCheck("compactionSwitch"); 1613 final CompactSplit compactSplitThread = server.getCompactSplitThread(); 1614 requestCount.increment(); 1615 boolean prevState = compactSplitThread.isCompactionsEnabled(); 1616 CompactionSwitchResponse response = 1617 CompactionSwitchResponse.newBuilder().setPrevState(prevState).build(); 1618 if (prevState == request.getEnabled()) { 1619 // passed in requested state is same as current state. No action required 1620 return response; 1621 } 1622 compactSplitThread.switchCompaction(request.getEnabled()); 1623 return response; 1624 } 1625 1626 /** 1627 * Flush a region on the region server. 1628 * @param controller the RPC controller 1629 * @param request the request 1630 */ 1631 @Override 1632 @QosPriority(priority = HConstants.ADMIN_QOS) 1633 public FlushRegionResponse flushRegion(final RpcController controller, 1634 final FlushRegionRequest request) throws ServiceException { 1635 try { 1636 checkOpen(); 1637 requestCount.increment(); 1638 HRegion region = getRegion(request.getRegion()); 1639 LOG.info("Flushing " + region.getRegionInfo().getRegionNameAsString()); 1640 boolean shouldFlush = true; 1641 if (request.hasIfOlderThanTs()) { 1642 shouldFlush = region.getEarliestFlushTimeForAllStores() < request.getIfOlderThanTs(); 1643 } 1644 FlushRegionResponse.Builder builder = FlushRegionResponse.newBuilder(); 1645 if (shouldFlush) { 1646 boolean writeFlushWalMarker = 1647 request.hasWriteFlushWalMarker() ? request.getWriteFlushWalMarker() : false; 1648 // Go behind the curtain so we can manage writing of the flush WAL marker 1649 HRegion.FlushResultImpl flushResult = null; 1650 if (request.hasFamily()) { 1651 List<byte[]> families = new ArrayList(); 1652 families.add(request.getFamily().toByteArray()); 1653 TableDescriptor tableDescriptor = region.getTableDescriptor(); 1654 List<String> noSuchFamilies = families.stream() 1655 .filter(f -> !tableDescriptor.hasColumnFamily(f)).map(Bytes::toString).toList(); 1656 if (!noSuchFamilies.isEmpty()) { 1657 throw new NoSuchColumnFamilyException("Column families " + noSuchFamilies 1658 + " don't exist in table " + tableDescriptor.getTableName().getNameAsString()); 1659 } 1660 flushResult = 1661 region.flushcache(families, writeFlushWalMarker, FlushLifeCycleTracker.DUMMY); 1662 } else { 1663 flushResult = region.flushcache(true, writeFlushWalMarker, FlushLifeCycleTracker.DUMMY); 1664 } 1665 boolean compactionNeeded = flushResult.isCompactionNeeded(); 1666 if (compactionNeeded) { 1667 server.getCompactSplitThread().requestSystemCompaction(region, 1668 "Compaction through user triggered flush"); 1669 } 1670 builder.setFlushed(flushResult.isFlushSucceeded()); 1671 builder.setWroteFlushWalMarker(flushResult.wroteFlushWalMarker); 1672 } 1673 builder.setLastFlushTime(region.getEarliestFlushTimeForAllStores()); 1674 return builder.build(); 1675 } catch (DroppedSnapshotException ex) { 1676 // Cache flush can fail in a few places. If it fails in a critical 1677 // section, we get a DroppedSnapshotException and a replay of wal 1678 // is required. Currently the only way to do this is a restart of 1679 // the server. 1680 server.abort("Replay of WAL required. Forcing server shutdown", ex); 1681 throw new ServiceException(ex); 1682 } catch (IOException ie) { 1683 throw new ServiceException(ie); 1684 } 1685 } 1686 1687 @Override 1688 @QosPriority(priority = HConstants.ADMIN_QOS) 1689 public GetOnlineRegionResponse getOnlineRegion(final RpcController controller, 1690 final GetOnlineRegionRequest request) throws ServiceException { 1691 try { 1692 checkOpen(); 1693 requestCount.increment(); 1694 Map<String, HRegion> onlineRegions = server.getOnlineRegions(); 1695 List<RegionInfo> list = new ArrayList<>(onlineRegions.size()); 1696 for (HRegion region : onlineRegions.values()) { 1697 list.add(region.getRegionInfo()); 1698 } 1699 list.sort(RegionInfo.COMPARATOR); 1700 return ResponseConverter.buildGetOnlineRegionResponse(list); 1701 } catch (IOException ie) { 1702 throw new ServiceException(ie); 1703 } 1704 } 1705 1706 // Master implementation of this Admin Service differs given it is not 1707 // able to supply detail only known to RegionServer. See note on 1708 // MasterRpcServers#getRegionInfo. 1709 @Override 1710 @QosPriority(priority = HConstants.ADMIN_QOS) 1711 public GetRegionInfoResponse getRegionInfo(final RpcController controller, 1712 final GetRegionInfoRequest request) throws ServiceException { 1713 try { 1714 checkOpen(); 1715 requestCount.increment(); 1716 HRegion region = getRegion(request.getRegion()); 1717 RegionInfo info = region.getRegionInfo(); 1718 byte[] bestSplitRow; 1719 if (request.hasBestSplitRow() && request.getBestSplitRow()) { 1720 bestSplitRow = region.checkSplit(true).orElse(null); 1721 // when all table data are in memstore, bestSplitRow = null 1722 // try to flush region first 1723 if (bestSplitRow == null) { 1724 region.flush(true); 1725 bestSplitRow = region.checkSplit(true).orElse(null); 1726 } 1727 } else { 1728 bestSplitRow = null; 1729 } 1730 GetRegionInfoResponse.Builder builder = GetRegionInfoResponse.newBuilder(); 1731 builder.setRegionInfo(ProtobufUtil.toRegionInfo(info)); 1732 if (request.hasCompactionState() && request.getCompactionState()) { 1733 builder.setCompactionState(ProtobufUtil.createCompactionState(region.getCompactionState())); 1734 } 1735 builder.setSplittable(region.isSplittable()); 1736 builder.setMergeable(region.isMergeable()); 1737 if (request.hasBestSplitRow() && request.getBestSplitRow() && bestSplitRow != null) { 1738 builder.setBestSplitRow(UnsafeByteOperations.unsafeWrap(bestSplitRow)); 1739 } 1740 return builder.build(); 1741 } catch (IOException ie) { 1742 throw new ServiceException(ie); 1743 } 1744 } 1745 1746 @Override 1747 @QosPriority(priority = HConstants.ADMIN_QOS) 1748 public GetRegionLoadResponse getRegionLoad(RpcController controller, GetRegionLoadRequest request) 1749 throws ServiceException { 1750 1751 List<HRegion> regions; 1752 if (request.hasTableName()) { 1753 TableName tableName = ProtobufUtil.toTableName(request.getTableName()); 1754 regions = server.getRegions(tableName); 1755 } else { 1756 regions = server.getRegions(); 1757 } 1758 List<RegionLoad> rLoads = new ArrayList<>(regions.size()); 1759 RegionLoad.Builder regionLoadBuilder = ClusterStatusProtos.RegionLoad.newBuilder(); 1760 RegionSpecifier.Builder regionSpecifier = RegionSpecifier.newBuilder(); 1761 1762 try { 1763 for (HRegion region : regions) { 1764 rLoads.add(server.createRegionLoad(region, regionLoadBuilder, regionSpecifier)); 1765 } 1766 } catch (IOException e) { 1767 throw new ServiceException(e); 1768 } 1769 GetRegionLoadResponse.Builder builder = GetRegionLoadResponse.newBuilder(); 1770 builder.addAllRegionLoads(rLoads); 1771 return builder.build(); 1772 } 1773 1774 @Override 1775 @QosPriority(priority = HConstants.ADMIN_QOS) 1776 public ClearCompactionQueuesResponse clearCompactionQueues(RpcController controller, 1777 ClearCompactionQueuesRequest request) throws ServiceException { 1778 LOG.debug("Client=" + RpcServer.getRequestUserName().orElse(null) + "/" 1779 + RpcServer.getRemoteAddress().orElse(null) + " clear compactions queue"); 1780 ClearCompactionQueuesResponse.Builder respBuilder = ClearCompactionQueuesResponse.newBuilder(); 1781 requestCount.increment(); 1782 if (clearCompactionQueues.compareAndSet(false, true)) { 1783 final CompactSplit compactSplitThread = server.getCompactSplitThread(); 1784 try { 1785 checkOpen(); 1786 server.getRegionServerCoprocessorHost().preClearCompactionQueues(); 1787 for (String queueName : request.getQueueNameList()) { 1788 LOG.debug("clear " + queueName + " compaction queue"); 1789 switch (queueName) { 1790 case "long": 1791 compactSplitThread.clearLongCompactionsQueue(); 1792 break; 1793 case "short": 1794 compactSplitThread.clearShortCompactionsQueue(); 1795 break; 1796 default: 1797 LOG.warn("Unknown queue name " + queueName); 1798 throw new IOException("Unknown queue name " + queueName); 1799 } 1800 } 1801 server.getRegionServerCoprocessorHost().postClearCompactionQueues(); 1802 } catch (IOException ie) { 1803 throw new ServiceException(ie); 1804 } finally { 1805 clearCompactionQueues.set(false); 1806 } 1807 } else { 1808 LOG.warn("Clear compactions queue is executing by other admin."); 1809 } 1810 return respBuilder.build(); 1811 } 1812 1813 /** 1814 * Get some information of the region server. 1815 * @param controller the RPC controller 1816 * @param request the request 1817 */ 1818 @Override 1819 @QosPriority(priority = HConstants.ADMIN_QOS) 1820 public GetServerInfoResponse getServerInfo(final RpcController controller, 1821 final GetServerInfoRequest request) throws ServiceException { 1822 try { 1823 checkOpen(); 1824 } catch (IOException ie) { 1825 throw new ServiceException(ie); 1826 } 1827 requestCount.increment(); 1828 int infoPort = server.getInfoServer() != null ? server.getInfoServer().getPort() : -1; 1829 return ResponseConverter.buildGetServerInfoResponse(server.getServerName(), infoPort); 1830 } 1831 1832 @Override 1833 @QosPriority(priority = HConstants.ADMIN_QOS) 1834 public GetStoreFileResponse getStoreFile(final RpcController controller, 1835 final GetStoreFileRequest request) throws ServiceException { 1836 try { 1837 checkOpen(); 1838 HRegion region = getRegion(request.getRegion()); 1839 requestCount.increment(); 1840 Set<byte[]> columnFamilies; 1841 if (request.getFamilyCount() == 0) { 1842 columnFamilies = region.getTableDescriptor().getColumnFamilyNames(); 1843 } else { 1844 columnFamilies = new TreeSet<>(Bytes.BYTES_RAWCOMPARATOR); 1845 for (ByteString cf : request.getFamilyList()) { 1846 columnFamilies.add(cf.toByteArray()); 1847 } 1848 } 1849 int nCF = columnFamilies.size(); 1850 List<String> fileList = region.getStoreFileList(columnFamilies.toArray(new byte[nCF][])); 1851 GetStoreFileResponse.Builder builder = GetStoreFileResponse.newBuilder(); 1852 builder.addAllStoreFile(fileList); 1853 return builder.build(); 1854 } catch (IOException ie) { 1855 throw new ServiceException(ie); 1856 } 1857 } 1858 1859 private void throwOnWrongStartCode(OpenRegionRequest request) throws ServiceException { 1860 if (!request.hasServerStartCode()) { 1861 LOG.warn("OpenRegionRequest for {} does not have a start code", request.getOpenInfoList()); 1862 return; 1863 } 1864 throwOnWrongStartCode(request.getServerStartCode()); 1865 } 1866 1867 private void throwOnWrongStartCode(CloseRegionRequest request) throws ServiceException { 1868 if (!request.hasServerStartCode()) { 1869 LOG.warn("CloseRegionRequest for {} does not have a start code", request.getRegion()); 1870 return; 1871 } 1872 throwOnWrongStartCode(request.getServerStartCode()); 1873 } 1874 1875 private void throwOnWrongStartCode(long serverStartCode) throws ServiceException { 1876 // check that we are the same server that this RPC is intended for. 1877 if (server.getServerName().getStartcode() != serverStartCode) { 1878 throw new ServiceException(new DoNotRetryIOException( 1879 "This RPC was intended for a " + "different server with startCode: " + serverStartCode 1880 + ", this server is: " + server.getServerName())); 1881 } 1882 } 1883 1884 private void throwOnWrongStartCode(ExecuteProceduresRequest req) throws ServiceException { 1885 if (req.getOpenRegionCount() > 0) { 1886 for (OpenRegionRequest openReq : req.getOpenRegionList()) { 1887 throwOnWrongStartCode(openReq); 1888 } 1889 } 1890 if (req.getCloseRegionCount() > 0) { 1891 for (CloseRegionRequest closeReq : req.getCloseRegionList()) { 1892 throwOnWrongStartCode(closeReq); 1893 } 1894 } 1895 } 1896 1897 /** 1898 * Open asynchronously a region or a set of regions on the region server. The opening is 1899 * coordinated by ZooKeeper, and this method requires the znode to be created before being called. 1900 * As a consequence, this method should be called only from the master. 1901 * <p> 1902 * Different manages states for the region are: 1903 * </p> 1904 * <ul> 1905 * <li>region not opened: the region opening will start asynchronously.</li> 1906 * <li>a close is already in progress: this is considered as an error.</li> 1907 * <li>an open is already in progress: this new open request will be ignored. This is important 1908 * because the Master can do multiple requests if it crashes.</li> 1909 * <li>the region is already opened: this new open request will be ignored.</li> 1910 * </ul> 1911 * <p> 1912 * Bulk assign: If there are more than 1 region to open, it will be considered as a bulk assign. 1913 * For a single region opening, errors are sent through a ServiceException. For bulk assign, 1914 * errors are put in the response as FAILED_OPENING. 1915 * </p> 1916 * @param controller the RPC controller 1917 * @param request the request 1918 */ 1919 @Override 1920 @QosPriority(priority = HConstants.ADMIN_QOS) 1921 public OpenRegionResponse openRegion(final RpcController controller, 1922 final OpenRegionRequest request) throws ServiceException { 1923 requestCount.increment(); 1924 throwOnWrongStartCode(request); 1925 1926 OpenRegionResponse.Builder builder = OpenRegionResponse.newBuilder(); 1927 final int regionCount = request.getOpenInfoCount(); 1928 final Map<TableName, TableDescriptor> htds = new HashMap<>(regionCount); 1929 final boolean isBulkAssign = regionCount > 1; 1930 try { 1931 checkOpen(); 1932 } catch (IOException ie) { 1933 TableName tableName = null; 1934 if (regionCount == 1) { 1935 org.apache.hadoop.hbase.shaded.protobuf.generated.HBaseProtos.RegionInfo ri = 1936 request.getOpenInfo(0).getRegion(); 1937 if (ri != null) { 1938 tableName = ProtobufUtil.toTableName(ri.getTableName()); 1939 } 1940 } 1941 if (!TableName.META_TABLE_NAME.equals(tableName)) { 1942 throw new ServiceException(ie); 1943 } 1944 // We are assigning meta, wait a little for regionserver to finish initialization. 1945 // Default to quarter of RPC timeout 1946 int timeout = server.getConfiguration().getInt(HConstants.HBASE_RPC_TIMEOUT_KEY, 1947 HConstants.DEFAULT_HBASE_RPC_TIMEOUT) >> 2; 1948 long endTime = EnvironmentEdgeManager.currentTime() + timeout; 1949 synchronized (server.online) { 1950 try { 1951 while ( 1952 EnvironmentEdgeManager.currentTime() <= endTime && !server.isStopped() 1953 && !server.isOnline() 1954 ) { 1955 server.online.wait(server.getMsgInterval()); 1956 } 1957 checkOpen(); 1958 } catch (InterruptedException t) { 1959 Thread.currentThread().interrupt(); 1960 throw new ServiceException(t); 1961 } catch (IOException e) { 1962 throw new ServiceException(e); 1963 } 1964 } 1965 } 1966 1967 long masterSystemTime = request.hasMasterSystemTime() ? request.getMasterSystemTime() : -1; 1968 1969 for (RegionOpenInfo regionOpenInfo : request.getOpenInfoList()) { 1970 final RegionInfo region = ProtobufUtil.toRegionInfo(regionOpenInfo.getRegion()); 1971 TableDescriptor htd; 1972 try { 1973 String encodedName = region.getEncodedName(); 1974 byte[] encodedNameBytes = region.getEncodedNameAsBytes(); 1975 final HRegion onlineRegion = server.getRegion(encodedName); 1976 if (onlineRegion != null) { 1977 // The region is already online. This should not happen any more. 1978 String error = "Received OPEN for the region:" + region.getRegionNameAsString() 1979 + ", which is already online"; 1980 LOG.warn(error); 1981 // server.abort(error); 1982 // throw new IOException(error); 1983 builder.addOpeningState(RegionOpeningState.OPENED); 1984 continue; 1985 } 1986 LOG.info("Open " + region.getRegionNameAsString()); 1987 1988 final Boolean previous = 1989 server.getRegionsInTransitionInRS().putIfAbsent(encodedNameBytes, Boolean.TRUE); 1990 1991 if (Boolean.FALSE.equals(previous)) { 1992 if (server.getRegion(encodedName) != null) { 1993 // There is a close in progress. This should not happen any more. 1994 String error = "Received OPEN for the region:" + region.getRegionNameAsString() 1995 + ", which we are already trying to CLOSE"; 1996 server.abort(error); 1997 throw new IOException(error); 1998 } 1999 server.getRegionsInTransitionInRS().put(encodedNameBytes, Boolean.TRUE); 2000 } 2001 2002 if (Boolean.TRUE.equals(previous)) { 2003 // An open is in progress. This is supported, but let's log this. 2004 LOG.info("Receiving OPEN for the region:" + region.getRegionNameAsString() 2005 + ", which we are already trying to OPEN" 2006 + " - ignoring this new request for this region."); 2007 } 2008 2009 // We are opening this region. If it moves back and forth for whatever reason, we don't 2010 // want to keep returning the stale moved record while we are opening/if we close again. 2011 server.removeFromMovedRegions(region.getEncodedName()); 2012 2013 if (previous == null || !previous.booleanValue()) { 2014 htd = htds.get(region.getTable()); 2015 if (htd == null) { 2016 htd = server.getTableDescriptors().get(region.getTable()); 2017 htds.put(region.getTable(), htd); 2018 } 2019 if (htd == null) { 2020 throw new IOException("Missing table descriptor for " + region.getEncodedName()); 2021 } 2022 // If there is no action in progress, we can submit a specific handler. 2023 // Need to pass the expected version in the constructor. 2024 if (server.getExecutorService() == null) { 2025 LOG.info("No executor executorService; skipping open request"); 2026 } else { 2027 if (region.isMetaRegion()) { 2028 server.getExecutorService() 2029 .submit(new OpenMetaHandler(server, server, region, htd, masterSystemTime)); 2030 } else { 2031 if (regionOpenInfo.getFavoredNodesCount() > 0) { 2032 server.updateRegionFavoredNodesMapping(region.getEncodedName(), 2033 regionOpenInfo.getFavoredNodesList()); 2034 } 2035 if (htd.getPriority() >= HConstants.ADMIN_QOS || region.getTable().isSystemTable()) { 2036 server.getExecutorService().submit( 2037 new OpenPriorityRegionHandler(server, server, region, htd, masterSystemTime)); 2038 } else { 2039 server.getExecutorService() 2040 .submit(new OpenRegionHandler(server, server, region, htd, masterSystemTime)); 2041 } 2042 } 2043 } 2044 } 2045 2046 builder.addOpeningState(RegionOpeningState.OPENED); 2047 } catch (IOException ie) { 2048 LOG.warn("Failed opening region " + region.getRegionNameAsString(), ie); 2049 if (isBulkAssign) { 2050 builder.addOpeningState(RegionOpeningState.FAILED_OPENING); 2051 } else { 2052 throw new ServiceException(ie); 2053 } 2054 } 2055 } 2056 return builder.build(); 2057 } 2058 2059 /** 2060 * Warmup a region on this server. This method should only be called by Master. It synchronously 2061 * opens the region and closes the region bringing the most important pages in cache. 2062 */ 2063 @Override 2064 public WarmupRegionResponse warmupRegion(final RpcController controller, 2065 final WarmupRegionRequest request) throws ServiceException { 2066 final RegionInfo region = ProtobufUtil.toRegionInfo(request.getRegionInfo()); 2067 WarmupRegionResponse response = WarmupRegionResponse.getDefaultInstance(); 2068 try { 2069 checkOpen(); 2070 String encodedName = region.getEncodedName(); 2071 byte[] encodedNameBytes = region.getEncodedNameAsBytes(); 2072 final HRegion onlineRegion = server.getRegion(encodedName); 2073 if (onlineRegion != null) { 2074 LOG.info("{} is online; skipping warmup", region); 2075 return response; 2076 } 2077 TableDescriptor htd = server.getTableDescriptors().get(region.getTable()); 2078 if (server.getRegionsInTransitionInRS().containsKey(encodedNameBytes)) { 2079 LOG.info("{} is in transition; skipping warmup", region); 2080 return response; 2081 } 2082 LOG.info("Warmup {}", region.getRegionNameAsString()); 2083 HRegion.warmupHRegion(region, htd, server.getWAL(region), server.getConfiguration(), server, 2084 null); 2085 } catch (IOException ie) { 2086 LOG.error("Failed warmup of {}", region.getRegionNameAsString(), ie); 2087 throw new ServiceException(ie); 2088 } 2089 2090 return response; 2091 } 2092 2093 private ExtendedCellScanner getAndReset(RpcController controller) { 2094 HBaseRpcController hrc = (HBaseRpcController) controller; 2095 ExtendedCellScanner cells = hrc.cellScanner(); 2096 hrc.setCellScanner(null); 2097 return cells; 2098 } 2099 2100 /** 2101 * Replay the given changes when distributedLogReplay WAL edits from a failed RS. The guarantee is 2102 * that the given mutations will be durable on the receiving RS if this method returns without any 2103 * exception. 2104 * @param controller the RPC controller 2105 * @param request the request 2106 * @deprecated Since 3.0.0, will be removed in 4.0.0. Not used any more, put here only for 2107 * compatibility with old region replica implementation. Now we will use 2108 * {@code replicateToReplica} method instead. 2109 */ 2110 @Deprecated 2111 @Override 2112 @QosPriority(priority = HConstants.REPLAY_QOS) 2113 public ReplicateWALEntryResponse replay(final RpcController controller, 2114 final ReplicateWALEntryRequest request) throws ServiceException { 2115 long before = EnvironmentEdgeManager.currentTime(); 2116 ExtendedCellScanner cells = getAndReset(controller); 2117 try { 2118 checkOpen(); 2119 List<WALEntry> entries = request.getEntryList(); 2120 if (entries == null || entries.isEmpty()) { 2121 // empty input 2122 return ReplicateWALEntryResponse.newBuilder().build(); 2123 } 2124 ByteString regionName = entries.get(0).getKey().getEncodedRegionName(); 2125 HRegion region = server.getRegionByEncodedName(regionName.toStringUtf8()); 2126 RegionCoprocessorHost coprocessorHost = 2127 ServerRegionReplicaUtil.isDefaultReplica(region.getRegionInfo()) 2128 ? region.getCoprocessorHost() 2129 : null; // do not invoke coprocessors if this is a secondary region replica 2130 List<Pair<WALKey, WALEdit>> walEntries = new ArrayList<>(); 2131 2132 // Skip adding the edits to WAL if this is a secondary region replica 2133 boolean isPrimary = RegionReplicaUtil.isDefaultReplica(region.getRegionInfo()); 2134 Durability durability = isPrimary ? Durability.USE_DEFAULT : Durability.SKIP_WAL; 2135 2136 for (WALEntry entry : entries) { 2137 if (!regionName.equals(entry.getKey().getEncodedRegionName())) { 2138 throw new NotServingRegionException("Replay request contains entries from multiple " 2139 + "regions. First region:" + regionName.toStringUtf8() + " , other region:" 2140 + entry.getKey().getEncodedRegionName()); 2141 } 2142 if (server.nonceManager != null && isPrimary) { 2143 long nonceGroup = 2144 entry.getKey().hasNonceGroup() ? entry.getKey().getNonceGroup() : HConstants.NO_NONCE; 2145 long nonce = entry.getKey().hasNonce() ? entry.getKey().getNonce() : HConstants.NO_NONCE; 2146 server.nonceManager.reportOperationFromWal(nonceGroup, nonce, 2147 entry.getKey().getWriteTime()); 2148 } 2149 Pair<WALKey, WALEdit> walEntry = (coprocessorHost == null) ? null : new Pair<>(); 2150 List<MutationReplay> edits = 2151 WALSplitUtil.getMutationsFromWALEntry(entry, cells, walEntry, durability); 2152 if (coprocessorHost != null) { 2153 // Start coprocessor replay here. The coprocessor is for each WALEdit instead of a 2154 // KeyValue. 2155 if ( 2156 coprocessorHost.preWALRestore(region.getRegionInfo(), walEntry.getFirst(), 2157 walEntry.getSecond()) 2158 ) { 2159 // if bypass this log entry, ignore it ... 2160 continue; 2161 } 2162 walEntries.add(walEntry); 2163 } 2164 if (edits != null && !edits.isEmpty()) { 2165 // HBASE-17924 2166 // sort to improve lock efficiency 2167 Collections.sort(edits, (v1, v2) -> Row.COMPARATOR.compare(v1.mutation, v2.mutation)); 2168 long replaySeqId = (entry.getKey().hasOrigSequenceNumber()) 2169 ? entry.getKey().getOrigSequenceNumber() 2170 : entry.getKey().getLogSequenceNumber(); 2171 OperationStatus[] result = doReplayBatchOp(region, edits, replaySeqId); 2172 // check if it's a partial success 2173 for (int i = 0; result != null && i < result.length; i++) { 2174 if (result[i] != OperationStatus.SUCCESS) { 2175 throw new IOException(result[i].getExceptionMsg()); 2176 } 2177 } 2178 } 2179 } 2180 2181 // sync wal at the end because ASYNC_WAL is used above 2182 WAL wal = region.getWAL(); 2183 if (wal != null) { 2184 wal.sync(); 2185 } 2186 2187 if (coprocessorHost != null) { 2188 for (Pair<WALKey, WALEdit> entry : walEntries) { 2189 coprocessorHost.postWALRestore(region.getRegionInfo(), entry.getFirst(), 2190 entry.getSecond()); 2191 } 2192 } 2193 return ReplicateWALEntryResponse.newBuilder().build(); 2194 } catch (IOException ie) { 2195 throw new ServiceException(ie); 2196 } finally { 2197 final MetricsRegionServer metricsRegionServer = server.getMetrics(); 2198 if (metricsRegionServer != null) { 2199 metricsRegionServer.updateReplay(EnvironmentEdgeManager.currentTime() - before); 2200 } 2201 } 2202 } 2203 2204 /** 2205 * Replay the given changes on a secondary replica 2206 */ 2207 @Override 2208 public ReplicateWALEntryResponse replicateToReplica(RpcController controller, 2209 ReplicateWALEntryRequest request) throws ServiceException { 2210 CellScanner cells = getAndReset(controller); 2211 try { 2212 checkOpen(); 2213 List<WALEntry> entries = request.getEntryList(); 2214 if (entries == null || entries.isEmpty()) { 2215 // empty input 2216 return ReplicateWALEntryResponse.newBuilder().build(); 2217 } 2218 ByteString regionName = entries.get(0).getKey().getEncodedRegionName(); 2219 HRegion region = server.getRegionByEncodedName(regionName.toStringUtf8()); 2220 if (RegionReplicaUtil.isDefaultReplica(region.getRegionInfo())) { 2221 throw new DoNotRetryIOException( 2222 "Should not replicate to primary replica " + region.getRegionInfo() + ", CODE BUG?"); 2223 } 2224 for (WALEntry entry : entries) { 2225 if (!regionName.equals(entry.getKey().getEncodedRegionName())) { 2226 throw new NotServingRegionException( 2227 "ReplicateToReplica request contains entries from multiple " + "regions. First region:" 2228 + regionName.toStringUtf8() + " , other region:" 2229 + entry.getKey().getEncodedRegionName()); 2230 } 2231 region.replayWALEntry(entry, cells); 2232 } 2233 return ReplicateWALEntryResponse.newBuilder().build(); 2234 } catch (IOException ie) { 2235 throw new ServiceException(ie); 2236 } 2237 } 2238 2239 private void checkShouldRejectReplicationRequest(List<WALEntry> entries) throws IOException { 2240 ReplicationSourceService replicationSource = server.getReplicationSourceService(); 2241 if (replicationSource == null || entries.isEmpty()) { 2242 return; 2243 } 2244 // We can ensure that all entries are for one peer, so only need to check one entry's 2245 // table name. if the table hit sync replication at peer side and the peer cluster 2246 // is (or is transiting to) state ACTIVE or DOWNGRADE_ACTIVE, we should reject to apply 2247 // those entries according to the design doc. 2248 TableName table = TableName.valueOf(entries.get(0).getKey().getTableName().toByteArray()); 2249 if ( 2250 replicationSource.getSyncReplicationPeerInfoProvider().checkState(table, 2251 RejectReplicationRequestStateChecker.get()) 2252 ) { 2253 throw new DoNotRetryIOException( 2254 "Reject to apply to sink cluster because sync replication state of sink cluster " 2255 + "is ACTIVE or DOWNGRADE_ACTIVE, table: " + table); 2256 } 2257 } 2258 2259 /** 2260 * Replicate WAL entries on the region server. 2261 * @param controller the RPC controller 2262 * @param request the request 2263 */ 2264 @Override 2265 @QosPriority(priority = HConstants.REPLICATION_QOS) 2266 public ReplicateWALEntryResponse replicateWALEntry(final RpcController controller, 2267 final ReplicateWALEntryRequest request) throws ServiceException { 2268 try { 2269 checkOpen(); 2270 if (server.getReplicationSinkService() != null) { 2271 requestCount.increment(); 2272 List<WALEntry> entries = request.getEntryList(); 2273 checkShouldRejectReplicationRequest(entries); 2274 ExtendedCellScanner cellScanner = getAndReset(controller); 2275 server.getRegionServerCoprocessorHost().preReplicateLogEntries(); 2276 server.getReplicationSinkService().replicateLogEntries(entries, cellScanner, 2277 request.getReplicationClusterId(), request.getSourceBaseNamespaceDirPath(), 2278 request.getSourceHFileArchiveDirPath()); 2279 server.getRegionServerCoprocessorHost().postReplicateLogEntries(); 2280 return ReplicateWALEntryResponse.newBuilder().build(); 2281 } else { 2282 throw new ServiceException("Replication services are not initialized yet"); 2283 } 2284 } catch (IOException ie) { 2285 throw new ServiceException(ie); 2286 } 2287 } 2288 2289 /** 2290 * Roll the WAL writer of the region server. 2291 * @param controller the RPC controller 2292 * @param request the request 2293 */ 2294 @Override 2295 @QosPriority(priority = HConstants.ADMIN_QOS) 2296 public RollWALWriterResponse rollWALWriter(final RpcController controller, 2297 final RollWALWriterRequest request) throws ServiceException { 2298 try { 2299 checkOpen(); 2300 requestCount.increment(); 2301 server.getRegionServerCoprocessorHost().preRollWALWriterRequest(); 2302 server.getWalRoller().requestRollAll(); 2303 server.getRegionServerCoprocessorHost().postRollWALWriterRequest(); 2304 RollWALWriterResponse.Builder builder = RollWALWriterResponse.newBuilder(); 2305 return builder.build(); 2306 } catch (IOException ie) { 2307 throw new ServiceException(ie); 2308 } 2309 } 2310 2311 /** 2312 * Stop the region server. 2313 * @param controller the RPC controller 2314 * @param request the request 2315 */ 2316 @Override 2317 @QosPriority(priority = HConstants.ADMIN_QOS) 2318 public StopServerResponse stopServer(final RpcController controller, 2319 final StopServerRequest request) throws ServiceException { 2320 rpcPreCheck("stopServer"); 2321 requestCount.increment(); 2322 String reason = request.getReason(); 2323 server.stop(reason); 2324 return StopServerResponse.newBuilder().build(); 2325 } 2326 2327 @Override 2328 public UpdateFavoredNodesResponse updateFavoredNodes(RpcController controller, 2329 UpdateFavoredNodesRequest request) throws ServiceException { 2330 rpcPreCheck("updateFavoredNodes"); 2331 List<UpdateFavoredNodesRequest.RegionUpdateInfo> openInfoList = request.getUpdateInfoList(); 2332 UpdateFavoredNodesResponse.Builder respBuilder = UpdateFavoredNodesResponse.newBuilder(); 2333 for (UpdateFavoredNodesRequest.RegionUpdateInfo regionUpdateInfo : openInfoList) { 2334 RegionInfo hri = ProtobufUtil.toRegionInfo(regionUpdateInfo.getRegion()); 2335 if (regionUpdateInfo.getFavoredNodesCount() > 0) { 2336 server.updateRegionFavoredNodesMapping(hri.getEncodedName(), 2337 regionUpdateInfo.getFavoredNodesList()); 2338 } 2339 } 2340 respBuilder.setResponse(openInfoList.size()); 2341 return respBuilder.build(); 2342 } 2343 2344 /** 2345 * Atomically bulk load several HFiles into an open region 2346 * @return true if successful, false is failed but recoverably (no action) 2347 * @throws ServiceException if failed unrecoverably 2348 */ 2349 @Override 2350 public BulkLoadHFileResponse bulkLoadHFile(final RpcController controller, 2351 final BulkLoadHFileRequest request) throws ServiceException { 2352 long start = EnvironmentEdgeManager.currentTime(); 2353 List<String> clusterIds = new ArrayList<>(request.getClusterIdsList()); 2354 if (clusterIds.contains(this.server.getClusterId())) { 2355 return BulkLoadHFileResponse.newBuilder().setLoaded(true).build(); 2356 } else { 2357 clusterIds.add(this.server.getClusterId()); 2358 } 2359 try { 2360 checkOpen(); 2361 requestCount.increment(); 2362 HRegion region = getRegion(request.getRegion()); 2363 final boolean spaceQuotaEnabled = QuotaUtil.isQuotaEnabled(getConfiguration()); 2364 long sizeToBeLoaded = -1; 2365 2366 // Check to see if this bulk load would exceed the space quota for this table 2367 if (spaceQuotaEnabled) { 2368 ActivePolicyEnforcement activeSpaceQuotas = getSpaceQuotaManager().getActiveEnforcements(); 2369 SpaceViolationPolicyEnforcement enforcement = 2370 activeSpaceQuotas.getPolicyEnforcement(region); 2371 if (enforcement != null) { 2372 // Bulk loads must still be atomic. We must enact all or none. 2373 List<String> filePaths = new ArrayList<>(request.getFamilyPathCount()); 2374 for (FamilyPath familyPath : request.getFamilyPathList()) { 2375 filePaths.add(familyPath.getPath()); 2376 } 2377 // Check if the batch of files exceeds the current quota 2378 sizeToBeLoaded = enforcement.computeBulkLoadSize(getFileSystem(filePaths), filePaths); 2379 } 2380 } 2381 // secure bulk load 2382 Map<byte[], List<Path>> map = 2383 server.getSecureBulkLoadManager().secureBulkLoadHFiles(region, request, clusterIds); 2384 BulkLoadHFileResponse.Builder builder = BulkLoadHFileResponse.newBuilder(); 2385 builder.setLoaded(map != null); 2386 if (map != null) { 2387 // Treat any negative size as a flag to "ignore" updating the region size as that is 2388 // not possible to occur in real life (cannot bulk load a file with negative size) 2389 if (spaceQuotaEnabled && sizeToBeLoaded > 0) { 2390 if (LOG.isTraceEnabled()) { 2391 LOG.trace("Incrementing space use of " + region.getRegionInfo() + " by " 2392 + sizeToBeLoaded + " bytes"); 2393 } 2394 // Inform space quotas of the new files for this region 2395 getSpaceQuotaManager().getRegionSizeStore().incrementRegionSize(region.getRegionInfo(), 2396 sizeToBeLoaded); 2397 } 2398 } 2399 return builder.build(); 2400 } catch (IOException ie) { 2401 throw new ServiceException(ie); 2402 } finally { 2403 final MetricsRegionServer metricsRegionServer = server.getMetrics(); 2404 if (metricsRegionServer != null) { 2405 metricsRegionServer.updateBulkLoad(EnvironmentEdgeManager.currentTime() - start); 2406 } 2407 } 2408 } 2409 2410 @Override 2411 public PrepareBulkLoadResponse prepareBulkLoad(RpcController controller, 2412 PrepareBulkLoadRequest request) throws ServiceException { 2413 try { 2414 checkOpen(); 2415 requestCount.increment(); 2416 2417 HRegion region = getRegion(request.getRegion()); 2418 2419 String bulkToken = server.getSecureBulkLoadManager().prepareBulkLoad(region, request); 2420 PrepareBulkLoadResponse.Builder builder = PrepareBulkLoadResponse.newBuilder(); 2421 builder.setBulkToken(bulkToken); 2422 return builder.build(); 2423 } catch (IOException ie) { 2424 throw new ServiceException(ie); 2425 } 2426 } 2427 2428 @Override 2429 public CleanupBulkLoadResponse cleanupBulkLoad(RpcController controller, 2430 CleanupBulkLoadRequest request) throws ServiceException { 2431 try { 2432 checkOpen(); 2433 requestCount.increment(); 2434 2435 HRegion region = getRegion(request.getRegion()); 2436 2437 server.getSecureBulkLoadManager().cleanupBulkLoad(region, request); 2438 return CleanupBulkLoadResponse.newBuilder().build(); 2439 } catch (IOException ie) { 2440 throw new ServiceException(ie); 2441 } 2442 } 2443 2444 @Override 2445 public CoprocessorServiceResponse execService(final RpcController controller, 2446 final CoprocessorServiceRequest request) throws ServiceException { 2447 try { 2448 checkOpen(); 2449 requestCount.increment(); 2450 HRegion region = getRegion(request.getRegion()); 2451 Message result = execServiceOnRegion(region, request.getCall()); 2452 CoprocessorServiceResponse.Builder builder = CoprocessorServiceResponse.newBuilder(); 2453 builder.setRegion(RequestConverter.buildRegionSpecifier(RegionSpecifierType.REGION_NAME, 2454 region.getRegionInfo().getRegionName())); 2455 // TODO: COPIES!!!!!! 2456 builder.setValue(builder.getValueBuilder().setName(result.getClass().getName()).setValue( 2457 org.apache.hbase.thirdparty.com.google.protobuf.ByteString.copyFrom(result.toByteArray()))); 2458 return builder.build(); 2459 } catch (IOException ie) { 2460 throw new ServiceException(ie); 2461 } 2462 } 2463 2464 private FileSystem getFileSystem(List<String> filePaths) throws IOException { 2465 if (filePaths.isEmpty()) { 2466 // local hdfs 2467 return server.getFileSystem(); 2468 } 2469 // source hdfs 2470 return new Path(filePaths.get(0)).getFileSystem(server.getConfiguration()); 2471 } 2472 2473 private Message execServiceOnRegion(HRegion region, 2474 final ClientProtos.CoprocessorServiceCall serviceCall) throws IOException { 2475 // ignore the passed in controller (from the serialized call) 2476 ServerRpcController execController = new ServerRpcController(); 2477 return region.execService(execController, serviceCall); 2478 } 2479 2480 private boolean shouldRejectRequestsFromClient(HRegion region) { 2481 TableName table = region.getRegionInfo().getTable(); 2482 ReplicationSourceService service = server.getReplicationSourceService(); 2483 return service != null && service.getSyncReplicationPeerInfoProvider().checkState(table, 2484 RejectRequestsFromClientStateChecker.get()); 2485 } 2486 2487 private void rejectIfInStandByState(HRegion region) throws DoNotRetryIOException { 2488 if (shouldRejectRequestsFromClient(region)) { 2489 throw new DoNotRetryIOException( 2490 region.getRegionInfo().getRegionNameAsString() + " is in STANDBY state."); 2491 } 2492 } 2493 2494 /** 2495 * Get data from a table. 2496 * @param controller the RPC controller 2497 * @param request the get request 2498 */ 2499 @Override 2500 public GetResponse get(final RpcController controller, final GetRequest request) 2501 throws ServiceException { 2502 long before = EnvironmentEdgeManager.currentTime(); 2503 OperationQuota quota = null; 2504 HRegion region = null; 2505 RpcCallContext context = RpcServer.getCurrentCall().orElse(null); 2506 try { 2507 checkOpen(); 2508 requestCount.increment(); 2509 rpcGetRequestCount.increment(); 2510 region = getRegion(request.getRegion()); 2511 rejectIfInStandByState(region); 2512 2513 GetResponse.Builder builder = GetResponse.newBuilder(); 2514 ClientProtos.Get get = request.getGet(); 2515 // An asynchbase client, https://github.com/OpenTSDB/asynchbase, starts by trying to do 2516 // a get closest before. Throwing the UnknownProtocolException signals it that it needs 2517 // to switch and do hbase2 protocol (HBase servers do not tell clients what versions 2518 // they are; its a problem for non-native clients like asynchbase. HBASE-20225. 2519 if (get.hasClosestRowBefore() && get.getClosestRowBefore()) { 2520 throw new UnknownProtocolException("Is this a pre-hbase-1.0.0 or asynchbase client? " 2521 + "Client is invoking getClosestRowBefore removed in hbase-2.0.0 replaced by " 2522 + "reverse Scan."); 2523 } 2524 Boolean existence = null; 2525 Result r = null; 2526 quota = getRpcQuotaManager().checkBatchQuota(region, OperationQuota.OperationType.GET); 2527 2528 Get clientGet = ProtobufUtil.toGet(get); 2529 if (get.getExistenceOnly() && region.getCoprocessorHost() != null) { 2530 existence = region.getCoprocessorHost().preExists(clientGet); 2531 } 2532 if (existence == null) { 2533 if (context != null) { 2534 r = get(clientGet, (region), null, context); 2535 } else { 2536 // for test purpose 2537 r = region.get(clientGet); 2538 } 2539 if (get.getExistenceOnly()) { 2540 boolean exists = r.getExists(); 2541 if (region.getCoprocessorHost() != null) { 2542 exists = region.getCoprocessorHost().postExists(clientGet, exists); 2543 } 2544 existence = exists; 2545 } 2546 } 2547 if (existence != null) { 2548 ClientProtos.Result pbr = ProtobufUtil.toResult(existence, 2549 region.getRegionInfo().getReplicaId() != 0, r != null ? r.getMetrics() : null); 2550 builder.setResult(pbr); 2551 } else if (r != null) { 2552 ClientProtos.Result pbr; 2553 if ( 2554 isClientCellBlockSupport(context) && controller instanceof HBaseRpcController 2555 && VersionInfoUtil.hasMinimumVersion(context.getClientVersionInfo(), 1, 3) 2556 ) { 2557 pbr = ProtobufUtil.toResultNoData(r); 2558 ((HBaseRpcController) controller).setCellScanner( 2559 PrivateCellUtil.createExtendedCellScanner(ClientInternalHelper.getExtendedRawCells(r))); 2560 addSize(context, r); 2561 } else { 2562 pbr = ProtobufUtil.toResult(r); 2563 } 2564 builder.setResult(pbr); 2565 } 2566 // r.cells is null when an table.exists(get) call 2567 if (r != null && r.rawCells() != null) { 2568 quota.addGetResult(r); 2569 } 2570 return builder.build(); 2571 } catch (IOException ie) { 2572 throw new ServiceException(ie); 2573 } finally { 2574 final MetricsRegionServer metricsRegionServer = server.getMetrics(); 2575 if (metricsRegionServer != null && region != null) { 2576 long blockBytesScanned = context != null ? context.getBlockBytesScanned() : 0; 2577 metricsRegionServer.updateGet(region, EnvironmentEdgeManager.currentTime() - before, 2578 blockBytesScanned); 2579 } 2580 if (quota != null) { 2581 quota.close(); 2582 } 2583 } 2584 } 2585 2586 private Result get(Get get, HRegion region, RegionScannersCloseCallBack closeCallBack, 2587 RpcCallContext context) throws IOException { 2588 region.prepareGet(get); 2589 boolean stale = region.getRegionInfo().getReplicaId() != 0; 2590 2591 // This method is almost the same as HRegion#get. 2592 List<Cell> results = new ArrayList<>(); 2593 // pre-get CP hook 2594 if (region.getCoprocessorHost() != null) { 2595 if (region.getCoprocessorHost().preGet(get, results)) { 2596 region.metricsUpdateForGet(); 2597 return Result.create(results, get.isCheckExistenceOnly() ? !results.isEmpty() : null, 2598 stale); 2599 } 2600 } 2601 Scan scan = new Scan(get); 2602 if (scan.getLoadColumnFamiliesOnDemandValue() == null) { 2603 scan.setLoadColumnFamiliesOnDemand(region.isLoadingCfsOnDemandDefault()); 2604 } 2605 RegionScannerImpl scanner = null; 2606 long blockBytesScannedBefore = context.getBlockBytesScanned(); 2607 try { 2608 scanner = region.getScanner(scan); 2609 scanner.next(results); 2610 } finally { 2611 if (scanner != null) { 2612 if (closeCallBack == null) { 2613 // If there is a context then the scanner can be added to the current 2614 // RpcCallContext. The rpc callback will take care of closing the 2615 // scanner, for eg in case 2616 // of get() 2617 context.setCallBack(scanner); 2618 } else { 2619 // The call is from multi() where the results from the get() are 2620 // aggregated and then send out to the 2621 // rpc. The rpccall back will close all such scanners created as part 2622 // of multi(). 2623 closeCallBack.addScanner(scanner); 2624 } 2625 } 2626 } 2627 2628 // post-get CP hook 2629 if (region.getCoprocessorHost() != null) { 2630 region.getCoprocessorHost().postGet(get, results); 2631 } 2632 region.metricsUpdateForGet(); 2633 2634 Result r = 2635 Result.create(results, get.isCheckExistenceOnly() ? !results.isEmpty() : null, stale); 2636 if (get.isQueryMetricsEnabled()) { 2637 long blockBytesScanned = context.getBlockBytesScanned() - blockBytesScannedBefore; 2638 r.setMetrics(new QueryMetrics(blockBytesScanned)); 2639 } 2640 return r; 2641 } 2642 2643 private void checkBatchSizeAndLogLargeSize(MultiRequest request) throws ServiceException { 2644 int sum = 0; 2645 String firstRegionName = null; 2646 for (RegionAction regionAction : request.getRegionActionList()) { 2647 if (sum == 0) { 2648 firstRegionName = Bytes.toStringBinary(regionAction.getRegion().getValue().toByteArray()); 2649 } 2650 sum += regionAction.getActionCount(); 2651 } 2652 if (sum > rowSizeWarnThreshold) { 2653 LOG.warn("Large batch operation detected (greater than " + rowSizeWarnThreshold 2654 + ") (HBASE-18023)." + " Requested Number of Rows: " + sum + " Client: " 2655 + RpcServer.getRequestUserName().orElse(null) + "/" 2656 + RpcServer.getRemoteAddress().orElse(null) + " first region in multi=" + firstRegionName); 2657 if (rejectRowsWithSizeOverThreshold) { 2658 throw new ServiceException( 2659 "Rejecting large batch operation for current batch with firstRegionName: " 2660 + firstRegionName + " , Requested Number of Rows: " + sum + " , Size Threshold: " 2661 + rowSizeWarnThreshold); 2662 } 2663 } 2664 } 2665 2666 private void failRegionAction(MultiResponse.Builder responseBuilder, 2667 RegionActionResult.Builder regionActionResultBuilder, RegionAction regionAction, 2668 CellScanner cellScanner, Throwable error) { 2669 rpcServer.getMetrics().exception(error); 2670 regionActionResultBuilder.setException(ResponseConverter.buildException(error)); 2671 responseBuilder.addRegionActionResult(regionActionResultBuilder.build()); 2672 // All Mutations in this RegionAction not executed as we can not see the Region online here 2673 // in this RS. Will be retried from Client. Skipping all the Cells in CellScanner 2674 // corresponding to these Mutations. 2675 if (cellScanner != null) { 2676 skipCellsForMutations(regionAction.getActionList(), cellScanner); 2677 } 2678 } 2679 2680 private boolean isReplicationRequest(Action action) { 2681 // replication request can only be put or delete. 2682 if (!action.hasMutation()) { 2683 return false; 2684 } 2685 MutationProto mutation = action.getMutation(); 2686 MutationType type = mutation.getMutateType(); 2687 if (type != MutationType.PUT && type != MutationType.DELETE) { 2688 return false; 2689 } 2690 // replication will set a special attribute so we can make use of it to decide whether a request 2691 // is for replication. 2692 return mutation.getAttributeList().stream().map(p -> p.getName()) 2693 .filter(n -> n.equals(ReplicationUtils.REPLICATION_ATTR_NAME)).findAny().isPresent(); 2694 } 2695 2696 /** 2697 * Execute multiple actions on a table: get, mutate, and/or execCoprocessor 2698 * @param rpcc the RPC controller 2699 * @param request the multi request 2700 */ 2701 @Override 2702 public MultiResponse multi(final RpcController rpcc, final MultiRequest request) 2703 throws ServiceException { 2704 try { 2705 checkOpen(); 2706 } catch (IOException ie) { 2707 throw new ServiceException(ie); 2708 } 2709 2710 checkBatchSizeAndLogLargeSize(request); 2711 2712 // rpc controller is how we bring in data via the back door; it is unprotobuf'ed data. 2713 // It is also the conduit via which we pass back data. 2714 HBaseRpcController controller = (HBaseRpcController) rpcc; 2715 CellScanner cellScanner = controller != null ? getAndReset(controller) : null; 2716 2717 long nonceGroup = request.hasNonceGroup() ? request.getNonceGroup() : HConstants.NO_NONCE; 2718 2719 MultiResponse.Builder responseBuilder = MultiResponse.newBuilder(); 2720 RegionActionResult.Builder regionActionResultBuilder = RegionActionResult.newBuilder(); 2721 this.rpcMultiRequestCount.increment(); 2722 this.requestCount.increment(); 2723 ActivePolicyEnforcement spaceQuotaEnforcement = getSpaceQuotaManager().getActiveEnforcements(); 2724 2725 // We no longer use MultiRequest#condition. Instead, we use RegionAction#condition. The 2726 // following logic is for backward compatibility as old clients still use 2727 // MultiRequest#condition in case of checkAndMutate with RowMutations. 2728 if (request.hasCondition()) { 2729 if (request.getRegionActionList().isEmpty()) { 2730 // If the region action list is empty, do nothing. 2731 responseBuilder.setProcessed(true); 2732 return responseBuilder.build(); 2733 } 2734 2735 RegionAction regionAction = request.getRegionAction(0); 2736 2737 // When request.hasCondition() is true, regionAction.getAtomic() should be always true. So 2738 // we can assume regionAction.getAtomic() is true here. 2739 assert regionAction.getAtomic(); 2740 2741 OperationQuota quota; 2742 HRegion region; 2743 RegionSpecifier regionSpecifier = regionAction.getRegion(); 2744 2745 try { 2746 region = getRegion(regionSpecifier); 2747 quota = getRpcQuotaManager().checkBatchQuota(region, regionAction.getActionList(), 2748 regionAction.hasCondition()); 2749 } catch (IOException e) { 2750 failRegionAction(responseBuilder, regionActionResultBuilder, regionAction, cellScanner, e); 2751 return responseBuilder.build(); 2752 } 2753 2754 try { 2755 boolean rejectIfFromClient = shouldRejectRequestsFromClient(region); 2756 // We only allow replication in standby state and it will not set the atomic flag. 2757 if (rejectIfFromClient) { 2758 failRegionAction(responseBuilder, regionActionResultBuilder, regionAction, cellScanner, 2759 new DoNotRetryIOException( 2760 region.getRegionInfo().getRegionNameAsString() + " is in STANDBY state")); 2761 return responseBuilder.build(); 2762 } 2763 2764 try { 2765 CheckAndMutateResult result = checkAndMutate(region, regionAction.getActionList(), 2766 cellScanner, request.getCondition(), nonceGroup, spaceQuotaEnforcement); 2767 responseBuilder.setProcessed(result.isSuccess()); 2768 ClientProtos.ResultOrException.Builder resultOrExceptionOrBuilder = 2769 ClientProtos.ResultOrException.newBuilder(); 2770 for (int i = 0; i < regionAction.getActionCount(); i++) { 2771 // To unify the response format with doNonAtomicRegionMutation and read through 2772 // client's AsyncProcess we have to add an empty result instance per operation 2773 resultOrExceptionOrBuilder.clear(); 2774 resultOrExceptionOrBuilder.setIndex(i); 2775 regionActionResultBuilder.addResultOrException(resultOrExceptionOrBuilder.build()); 2776 } 2777 } catch (IOException e) { 2778 rpcServer.getMetrics().exception(e); 2779 // As it's an atomic operation with a condition, we may expect it's a global failure. 2780 regionActionResultBuilder.setException(ResponseConverter.buildException(e)); 2781 } 2782 } finally { 2783 quota.close(); 2784 } 2785 2786 responseBuilder.addRegionActionResult(regionActionResultBuilder.build()); 2787 ClientProtos.RegionLoadStats regionLoadStats = region.getLoadStatistics(); 2788 if (regionLoadStats != null) { 2789 responseBuilder.setRegionStatistics(MultiRegionLoadStats.newBuilder() 2790 .addRegion(regionSpecifier).addStat(regionLoadStats).build()); 2791 } 2792 return responseBuilder.build(); 2793 } 2794 2795 // this will contain all the cells that we need to return. It's created later, if needed. 2796 List<ExtendedCellScannable> cellsToReturn = null; 2797 RegionScannersCloseCallBack closeCallBack = null; 2798 RpcCallContext context = RpcServer.getCurrentCall().orElse(null); 2799 Map<RegionSpecifier, ClientProtos.RegionLoadStats> regionStats = 2800 new HashMap<>(request.getRegionActionCount()); 2801 2802 for (RegionAction regionAction : request.getRegionActionList()) { 2803 OperationQuota quota; 2804 HRegion region; 2805 RegionSpecifier regionSpecifier = regionAction.getRegion(); 2806 regionActionResultBuilder.clear(); 2807 2808 try { 2809 region = getRegion(regionSpecifier); 2810 quota = getRpcQuotaManager().checkBatchQuota(region, regionAction.getActionList(), 2811 regionAction.hasCondition()); 2812 } catch (IOException e) { 2813 failRegionAction(responseBuilder, regionActionResultBuilder, regionAction, cellScanner, e); 2814 continue; // For this region it's a failure. 2815 } 2816 2817 try { 2818 boolean rejectIfFromClient = shouldRejectRequestsFromClient(region); 2819 2820 if (regionAction.hasCondition()) { 2821 // We only allow replication in standby state and it will not set the atomic flag. 2822 if (rejectIfFromClient) { 2823 failRegionAction(responseBuilder, regionActionResultBuilder, regionAction, cellScanner, 2824 new DoNotRetryIOException( 2825 region.getRegionInfo().getRegionNameAsString() + " is in STANDBY state")); 2826 continue; 2827 } 2828 2829 try { 2830 ClientProtos.ResultOrException.Builder resultOrExceptionOrBuilder = 2831 ClientProtos.ResultOrException.newBuilder(); 2832 if (regionAction.getActionCount() == 1) { 2833 CheckAndMutateResult result = 2834 checkAndMutate(region, quota, regionAction.getAction(0).getMutation(), cellScanner, 2835 regionAction.getCondition(), nonceGroup, spaceQuotaEnforcement, context); 2836 regionActionResultBuilder.setProcessed(result.isSuccess()); 2837 resultOrExceptionOrBuilder.setIndex(0); 2838 if (result.getResult() != null) { 2839 resultOrExceptionOrBuilder.setResult(ProtobufUtil.toResult(result.getResult())); 2840 } 2841 2842 if (result.getMetrics() != null) { 2843 resultOrExceptionOrBuilder 2844 .setMetrics(ProtobufUtil.toQueryMetrics(result.getMetrics())); 2845 } 2846 2847 regionActionResultBuilder.addResultOrException(resultOrExceptionOrBuilder.build()); 2848 } else { 2849 CheckAndMutateResult result = checkAndMutate(region, regionAction.getActionList(), 2850 cellScanner, regionAction.getCondition(), nonceGroup, spaceQuotaEnforcement); 2851 regionActionResultBuilder.setProcessed(result.isSuccess()); 2852 for (int i = 0; i < regionAction.getActionCount(); i++) { 2853 if (i == 0 && result.getResult() != null) { 2854 // Set the result of the Increment/Append operations to the first element of the 2855 // ResultOrException list 2856 resultOrExceptionOrBuilder.setIndex(i); 2857 regionActionResultBuilder.addResultOrException(resultOrExceptionOrBuilder 2858 .setResult(ProtobufUtil.toResult(result.getResult())).build()); 2859 continue; 2860 } 2861 // To unify the response format with doNonAtomicRegionMutation and read through 2862 // client's AsyncProcess we have to add an empty result instance per operation 2863 resultOrExceptionOrBuilder.clear(); 2864 resultOrExceptionOrBuilder.setIndex(i); 2865 regionActionResultBuilder.addResultOrException(resultOrExceptionOrBuilder.build()); 2866 } 2867 } 2868 } catch (IOException e) { 2869 rpcServer.getMetrics().exception(e); 2870 // As it's an atomic operation with a condition, we may expect it's a global failure. 2871 regionActionResultBuilder.setException(ResponseConverter.buildException(e)); 2872 } 2873 } else if (regionAction.hasAtomic() && regionAction.getAtomic()) { 2874 // We only allow replication in standby state and it will not set the atomic flag. 2875 if (rejectIfFromClient) { 2876 failRegionAction(responseBuilder, regionActionResultBuilder, regionAction, cellScanner, 2877 new DoNotRetryIOException( 2878 region.getRegionInfo().getRegionNameAsString() + " is in STANDBY state")); 2879 continue; 2880 } 2881 try { 2882 doAtomicBatchOp(regionActionResultBuilder, region, quota, regionAction.getActionList(), 2883 cellScanner, nonceGroup, spaceQuotaEnforcement); 2884 regionActionResultBuilder.setProcessed(true); 2885 // We no longer use MultiResponse#processed. Instead, we use 2886 // RegionActionResult#processed. This is for backward compatibility for old clients. 2887 responseBuilder.setProcessed(true); 2888 } catch (IOException e) { 2889 rpcServer.getMetrics().exception(e); 2890 // As it's atomic, we may expect it's a global failure. 2891 regionActionResultBuilder.setException(ResponseConverter.buildException(e)); 2892 } 2893 } else { 2894 if ( 2895 rejectIfFromClient && regionAction.getActionCount() > 0 2896 && !isReplicationRequest(regionAction.getAction(0)) 2897 ) { 2898 // fail if it is not a replication request 2899 failRegionAction(responseBuilder, regionActionResultBuilder, regionAction, cellScanner, 2900 new DoNotRetryIOException( 2901 region.getRegionInfo().getRegionNameAsString() + " is in STANDBY state")); 2902 continue; 2903 } 2904 // doNonAtomicRegionMutation manages the exception internally 2905 if (context != null && closeCallBack == null) { 2906 // An RpcCallBack that creates a list of scanners that needs to perform callBack 2907 // operation on completion of multiGets. 2908 // Set this only once 2909 closeCallBack = new RegionScannersCloseCallBack(); 2910 context.setCallBack(closeCallBack); 2911 } 2912 cellsToReturn = doNonAtomicRegionMutation(region, quota, regionAction, cellScanner, 2913 regionActionResultBuilder, cellsToReturn, nonceGroup, closeCallBack, context, 2914 spaceQuotaEnforcement); 2915 } 2916 } finally { 2917 quota.close(); 2918 } 2919 2920 responseBuilder.addRegionActionResult(regionActionResultBuilder.build()); 2921 ClientProtos.RegionLoadStats regionLoadStats = region.getLoadStatistics(); 2922 if (regionLoadStats != null) { 2923 regionStats.put(regionSpecifier, regionLoadStats); 2924 } 2925 } 2926 // Load the controller with the Cells to return. 2927 if (cellsToReturn != null && !cellsToReturn.isEmpty() && controller != null) { 2928 controller.setCellScanner(PrivateCellUtil.createExtendedCellScanner(cellsToReturn)); 2929 } 2930 2931 MultiRegionLoadStats.Builder builder = MultiRegionLoadStats.newBuilder(); 2932 for (Entry<RegionSpecifier, ClientProtos.RegionLoadStats> stat : regionStats.entrySet()) { 2933 builder.addRegion(stat.getKey()); 2934 builder.addStat(stat.getValue()); 2935 } 2936 responseBuilder.setRegionStatistics(builder); 2937 return responseBuilder.build(); 2938 } 2939 2940 private void skipCellsForMutations(List<Action> actions, CellScanner cellScanner) { 2941 if (cellScanner == null) { 2942 return; 2943 } 2944 for (Action action : actions) { 2945 skipCellsForMutation(action, cellScanner); 2946 } 2947 } 2948 2949 private void skipCellsForMutation(Action action, CellScanner cellScanner) { 2950 if (cellScanner == null) { 2951 return; 2952 } 2953 try { 2954 if (action.hasMutation()) { 2955 MutationProto m = action.getMutation(); 2956 if (m.hasAssociatedCellCount()) { 2957 for (int i = 0; i < m.getAssociatedCellCount(); i++) { 2958 cellScanner.advance(); 2959 } 2960 } 2961 } 2962 } catch (IOException e) { 2963 // No need to handle these Individual Muatation level issue. Any way this entire RegionAction 2964 // marked as failed as we could not see the Region here. At client side the top level 2965 // RegionAction exception will be considered first. 2966 LOG.error("Error while skipping Cells in CellScanner for invalid Region Mutations", e); 2967 } 2968 } 2969 2970 /** 2971 * Mutate data in a table. 2972 * @param rpcc the RPC controller 2973 * @param request the mutate request 2974 */ 2975 @Override 2976 public MutateResponse mutate(final RpcController rpcc, final MutateRequest request) 2977 throws ServiceException { 2978 // rpc controller is how we bring in data via the back door; it is unprotobuf'ed data. 2979 // It is also the conduit via which we pass back data. 2980 HBaseRpcController controller = (HBaseRpcController) rpcc; 2981 CellScanner cellScanner = controller != null ? controller.cellScanner() : null; 2982 OperationQuota quota = null; 2983 RpcCallContext context = RpcServer.getCurrentCall().orElse(null); 2984 // Clear scanner so we are not holding on to reference across call. 2985 if (controller != null) { 2986 controller.setCellScanner(null); 2987 } 2988 try { 2989 checkOpen(); 2990 requestCount.increment(); 2991 rpcMutateRequestCount.increment(); 2992 HRegion region = getRegion(request.getRegion()); 2993 rejectIfInStandByState(region); 2994 MutateResponse.Builder builder = MutateResponse.newBuilder(); 2995 MutationProto mutation = request.getMutation(); 2996 if (!region.getRegionInfo().isMetaRegion()) { 2997 server.getMemStoreFlusher().reclaimMemStoreMemory(); 2998 } 2999 long nonceGroup = request.hasNonceGroup() ? request.getNonceGroup() : HConstants.NO_NONCE; 3000 OperationQuota.OperationType operationType = QuotaUtil.getQuotaOperationType(request); 3001 quota = getRpcQuotaManager().checkBatchQuota(region, operationType); 3002 ActivePolicyEnforcement spaceQuotaEnforcement = 3003 getSpaceQuotaManager().getActiveEnforcements(); 3004 3005 if (request.hasCondition()) { 3006 CheckAndMutateResult result = checkAndMutate(region, quota, mutation, cellScanner, 3007 request.getCondition(), nonceGroup, spaceQuotaEnforcement, context); 3008 builder.setProcessed(result.isSuccess()); 3009 boolean clientCellBlockSupported = isClientCellBlockSupport(context); 3010 addResult(builder, result.getResult(), controller, clientCellBlockSupported); 3011 if (clientCellBlockSupported) { 3012 addSize(context, result.getResult()); 3013 } 3014 if (result.getMetrics() != null) { 3015 builder.setMetrics(ProtobufUtil.toQueryMetrics(result.getMetrics())); 3016 } 3017 } else { 3018 Result r = null; 3019 Boolean processed = null; 3020 MutationType type = mutation.getMutateType(); 3021 switch (type) { 3022 case APPEND: 3023 // TODO: this doesn't actually check anything. 3024 r = append(region, quota, mutation, cellScanner, nonceGroup, spaceQuotaEnforcement, 3025 context); 3026 break; 3027 case INCREMENT: 3028 // TODO: this doesn't actually check anything. 3029 r = increment(region, quota, mutation, cellScanner, nonceGroup, spaceQuotaEnforcement, 3030 context); 3031 break; 3032 case PUT: 3033 put(region, quota, mutation, cellScanner, spaceQuotaEnforcement); 3034 processed = Boolean.TRUE; 3035 break; 3036 case DELETE: 3037 delete(region, quota, mutation, cellScanner, spaceQuotaEnforcement); 3038 processed = Boolean.TRUE; 3039 break; 3040 default: 3041 throw new DoNotRetryIOException("Unsupported mutate type: " + type.name()); 3042 } 3043 if (processed != null) { 3044 builder.setProcessed(processed); 3045 } 3046 boolean clientCellBlockSupported = isClientCellBlockSupport(context); 3047 addResult(builder, r, controller, clientCellBlockSupported); 3048 if (clientCellBlockSupported) { 3049 addSize(context, r); 3050 } 3051 } 3052 return builder.build(); 3053 } catch (IOException ie) { 3054 server.checkFileSystem(); 3055 throw new ServiceException(ie); 3056 } finally { 3057 if (quota != null) { 3058 quota.close(); 3059 } 3060 } 3061 } 3062 3063 private void put(HRegion region, OperationQuota quota, MutationProto mutation, 3064 CellScanner cellScanner, ActivePolicyEnforcement spaceQuota) throws IOException { 3065 long before = EnvironmentEdgeManager.currentTime(); 3066 Put put = ProtobufUtil.toPut(mutation, cellScanner); 3067 checkCellSizeLimit(region, put); 3068 spaceQuota.getPolicyEnforcement(region).check(put); 3069 quota.addMutation(put); 3070 region.put(put); 3071 3072 MetricsRegionServer metricsRegionServer = server.getMetrics(); 3073 if (metricsRegionServer != null) { 3074 long after = EnvironmentEdgeManager.currentTime(); 3075 metricsRegionServer.updatePut(region, after - before); 3076 } 3077 } 3078 3079 private void delete(HRegion region, OperationQuota quota, MutationProto mutation, 3080 CellScanner cellScanner, ActivePolicyEnforcement spaceQuota) throws IOException { 3081 long before = EnvironmentEdgeManager.currentTime(); 3082 Delete delete = ProtobufUtil.toDelete(mutation, cellScanner); 3083 checkCellSizeLimit(region, delete); 3084 spaceQuota.getPolicyEnforcement(region).check(delete); 3085 quota.addMutation(delete); 3086 region.delete(delete); 3087 3088 MetricsRegionServer metricsRegionServer = server.getMetrics(); 3089 if (metricsRegionServer != null) { 3090 long after = EnvironmentEdgeManager.currentTime(); 3091 metricsRegionServer.updateDelete(region, after - before); 3092 } 3093 } 3094 3095 private CheckAndMutateResult checkAndMutate(HRegion region, OperationQuota quota, 3096 MutationProto mutation, CellScanner cellScanner, Condition condition, long nonceGroup, 3097 ActivePolicyEnforcement spaceQuota, RpcCallContext context) throws IOException { 3098 long before = EnvironmentEdgeManager.currentTime(); 3099 long blockBytesScannedBefore = context != null ? context.getBlockBytesScanned() : 0; 3100 CheckAndMutate checkAndMutate = ProtobufUtil.toCheckAndMutate(condition, mutation, cellScanner); 3101 long nonce = mutation.hasNonce() ? mutation.getNonce() : HConstants.NO_NONCE; 3102 checkCellSizeLimit(region, (Mutation) checkAndMutate.getAction()); 3103 spaceQuota.getPolicyEnforcement(region).check((Mutation) checkAndMutate.getAction()); 3104 quota.addMutation((Mutation) checkAndMutate.getAction()); 3105 3106 CheckAndMutateResult result = null; 3107 if (region.getCoprocessorHost() != null) { 3108 result = region.getCoprocessorHost().preCheckAndMutate(checkAndMutate); 3109 } 3110 if (result == null) { 3111 result = region.checkAndMutate(checkAndMutate, nonceGroup, nonce); 3112 if (region.getCoprocessorHost() != null) { 3113 result = region.getCoprocessorHost().postCheckAndMutate(checkAndMutate, result); 3114 } 3115 } 3116 MetricsRegionServer metricsRegionServer = server.getMetrics(); 3117 if (metricsRegionServer != null) { 3118 long after = EnvironmentEdgeManager.currentTime(); 3119 long blockBytesScanned = 3120 context != null ? context.getBlockBytesScanned() - blockBytesScannedBefore : 0; 3121 metricsRegionServer.updateCheckAndMutate(region, after - before, blockBytesScanned); 3122 3123 MutationType type = mutation.getMutateType(); 3124 switch (type) { 3125 case PUT: 3126 metricsRegionServer.updateCheckAndPut(region, after - before); 3127 break; 3128 case DELETE: 3129 metricsRegionServer.updateCheckAndDelete(region, after - before); 3130 break; 3131 default: 3132 break; 3133 } 3134 } 3135 return result; 3136 } 3137 3138 // This is used to keep compatible with the old client implementation. Consider remove it if we 3139 // decide to drop the support of the client that still sends close request to a region scanner 3140 // which has already been exhausted. 3141 @Deprecated 3142 private static final IOException SCANNER_ALREADY_CLOSED = new IOException() { 3143 3144 private static final long serialVersionUID = -4305297078988180130L; 3145 3146 @Override 3147 public synchronized Throwable fillInStackTrace() { 3148 return this; 3149 } 3150 }; 3151 3152 private RegionScannerHolder getRegionScanner(ScanRequest request) throws IOException { 3153 String scannerName = toScannerName(request.getScannerId()); 3154 RegionScannerHolder rsh = this.scanners.get(scannerName); 3155 if (rsh == null) { 3156 // just ignore the next or close request if scanner does not exists. 3157 Long lastCallSeq = closedScanners.getIfPresent(scannerName); 3158 if (lastCallSeq != null) { 3159 // Check the sequence number to catch if the last call was incorrectly retried. 3160 // The only allowed scenario is when the scanner is exhausted and one more scan 3161 // request arrives - in this case returning 0 rows is correct. 3162 if (request.hasNextCallSeq() && request.getNextCallSeq() != lastCallSeq + 1) { 3163 throw new OutOfOrderScannerNextException("Expected nextCallSeq for closed request: " 3164 + (lastCallSeq + 1) + " But the nextCallSeq got from client: " 3165 + request.getNextCallSeq() + "; request=" + TextFormat.shortDebugString(request)); 3166 } else { 3167 throw SCANNER_ALREADY_CLOSED; 3168 } 3169 } else { 3170 LOG.warn("Client tried to access missing scanner " + scannerName); 3171 throw new UnknownScannerException( 3172 "Unknown scanner '" + scannerName + "'. This can happen due to any of the following " 3173 + "reasons: a) Scanner id given is wrong, b) Scanner lease expired because of " 3174 + "long wait between consecutive client checkins, c) Server may be closing down, " 3175 + "d) RegionServer restart during upgrade.\nIf the issue is due to reason (b), a " 3176 + "possible fix would be increasing the value of" 3177 + "'hbase.client.scanner.timeout.period' configuration."); 3178 } 3179 } 3180 rejectIfInStandByState(rsh.r); 3181 RegionInfo hri = rsh.s.getRegionInfo(); 3182 // Yes, should be the same instance 3183 if (server.getOnlineRegion(hri.getRegionName()) != rsh.r) { 3184 String msg = "Region has changed on the scanner " + scannerName + ": regionName=" 3185 + hri.getRegionNameAsString() + ", scannerRegionName=" + rsh.r; 3186 LOG.warn(msg + ", closing..."); 3187 scanners.remove(scannerName); 3188 try { 3189 rsh.s.close(); 3190 } catch (IOException e) { 3191 LOG.warn("Getting exception closing " + scannerName, e); 3192 } finally { 3193 try { 3194 server.getLeaseManager().cancelLease(scannerName); 3195 } catch (LeaseException e) { 3196 LOG.warn("Getting exception closing " + scannerName, e); 3197 } 3198 } 3199 throw new NotServingRegionException(msg); 3200 } 3201 return rsh; 3202 } 3203 3204 /** 3205 * @return Pair with scannerName key to use with this new Scanner and its RegionScannerHolder 3206 * value. 3207 */ 3208 private Pair<String, RegionScannerHolder> newRegionScanner(ScanRequest request, HRegion region, 3209 ScanResponse.Builder builder) throws IOException { 3210 rejectIfInStandByState(region); 3211 ClientProtos.Scan protoScan = request.getScan(); 3212 boolean isLoadingCfsOnDemandSet = protoScan.hasLoadColumnFamiliesOnDemand(); 3213 Scan scan = ProtobufUtil.toScan(protoScan); 3214 // if the request doesn't set this, get the default region setting. 3215 if (!isLoadingCfsOnDemandSet) { 3216 scan.setLoadColumnFamiliesOnDemand(region.isLoadingCfsOnDemandDefault()); 3217 } 3218 3219 if (!scan.hasFamilies()) { 3220 // Adding all families to scanner 3221 for (byte[] family : region.getTableDescriptor().getColumnFamilyNames()) { 3222 scan.addFamily(family); 3223 } 3224 } 3225 if (region.getCoprocessorHost() != null) { 3226 // preScannerOpen is not allowed to return a RegionScanner. Only post hook can create a 3227 // wrapper for the core created RegionScanner 3228 region.getCoprocessorHost().preScannerOpen(scan); 3229 } 3230 RegionScannerImpl coreScanner = region.getScanner(scan); 3231 Shipper shipper = coreScanner; 3232 RegionScanner scanner = coreScanner; 3233 try { 3234 if (region.getCoprocessorHost() != null) { 3235 scanner = region.getCoprocessorHost().postScannerOpen(scan, scanner); 3236 } 3237 } catch (Exception e) { 3238 // Although region coprocessor is for advanced users and they should take care of the 3239 // implementation to not damage the HBase system, closing the scanner on exception here does 3240 // not have any bad side effect, so let's do it 3241 scanner.close(); 3242 throw e; 3243 } 3244 long scannerId = scannerIdGenerator.generateNewScannerId(); 3245 builder.setScannerId(scannerId); 3246 builder.setMvccReadPoint(scanner.getMvccReadPoint()); 3247 builder.setTtl(scannerLeaseTimeoutPeriod); 3248 String scannerName = toScannerName(scannerId); 3249 3250 boolean fullRegionScan = 3251 !region.getRegionInfo().getTable().isSystemTable() && isFullRegionScan(scan, region); 3252 3253 return new Pair<String, RegionScannerHolder>(scannerName, 3254 addScanner(scannerName, scanner, shipper, region, scan.isNeedCursorResult(), fullRegionScan)); 3255 } 3256 3257 /** 3258 * The returned String is used as key doing look up of outstanding Scanners in this Servers' 3259 * this.scanners, the Map of outstanding scanners and their current state. 3260 * @param scannerId A scanner long id. 3261 * @return The long id as a String. 3262 */ 3263 private static String toScannerName(long scannerId) { 3264 return Long.toString(scannerId); 3265 } 3266 3267 private void checkScanNextCallSeq(ScanRequest request, RegionScannerHolder rsh) 3268 throws OutOfOrderScannerNextException { 3269 // if nextCallSeq does not match throw Exception straight away. This needs to be 3270 // performed even before checking of Lease. 3271 // See HBASE-5974 3272 if (request.hasNextCallSeq()) { 3273 long callSeq = request.getNextCallSeq(); 3274 if (!rsh.incNextCallSeq(callSeq)) { 3275 throw new OutOfOrderScannerNextException( 3276 "Expected nextCallSeq: " + rsh.getNextCallSeq() + " But the nextCallSeq got from client: " 3277 + request.getNextCallSeq() + "; request=" + TextFormat.shortDebugString(request) 3278 + "; region=" + rsh.r.getRegionInfo().getRegionNameAsString()); 3279 } 3280 } 3281 } 3282 3283 private void addScannerLeaseBack(LeaseManager.Lease lease) { 3284 try { 3285 server.getLeaseManager().addLease(lease); 3286 } catch (LeaseStillHeldException e) { 3287 // should not happen as the scanner id is unique. 3288 throw new AssertionError(e); 3289 } 3290 } 3291 3292 // visible for testing only 3293 long getTimeLimit(RpcCall rpcCall, HBaseRpcController controller, 3294 boolean allowHeartbeatMessages) { 3295 // Set the time limit to be half of the more restrictive timeout value (one of the 3296 // timeout values must be positive). In the event that both values are positive, the 3297 // more restrictive of the two is used to calculate the limit. 3298 if (allowHeartbeatMessages) { 3299 long now = EnvironmentEdgeManager.currentTime(); 3300 long remainingTimeout = getRemainingRpcTimeout(rpcCall, controller, now); 3301 if (scannerLeaseTimeoutPeriod > 0 || remainingTimeout > 0) { 3302 long timeLimitDelta; 3303 if (scannerLeaseTimeoutPeriod > 0 && remainingTimeout > 0) { 3304 timeLimitDelta = Math.min(scannerLeaseTimeoutPeriod, remainingTimeout); 3305 } else { 3306 timeLimitDelta = 3307 scannerLeaseTimeoutPeriod > 0 ? scannerLeaseTimeoutPeriod : remainingTimeout; 3308 } 3309 3310 // Use half of whichever timeout value was more restrictive... But don't allow 3311 // the time limit to be less than the allowable minimum (could cause an 3312 // immediate timeout before scanning any data). 3313 timeLimitDelta = Math.max(timeLimitDelta / 2, minimumScanTimeLimitDelta); 3314 return now + timeLimitDelta; 3315 } 3316 } 3317 // Default value of timeLimit is negative to indicate no timeLimit should be 3318 // enforced. 3319 return -1L; 3320 } 3321 3322 private long getRemainingRpcTimeout(RpcCall call, HBaseRpcController controller, long now) { 3323 long timeout; 3324 if (controller != null && controller.getCallTimeout() > 0) { 3325 timeout = controller.getCallTimeout(); 3326 } else if (rpcTimeout > 0) { 3327 timeout = rpcTimeout; 3328 } else { 3329 return -1; 3330 } 3331 if (call != null) { 3332 timeout -= (now - call.getReceiveTime()); 3333 } 3334 // getTimeLimit ignores values <= 0, but timeout may now be negative if queue time was high. 3335 // return minimum value here in that case so we count this in calculating the final delta. 3336 return Math.max(minimumScanTimeLimitDelta, timeout); 3337 } 3338 3339 private void checkLimitOfRows(int numOfCompleteRows, int limitOfRows, boolean moreRows, 3340 ScannerContext scannerContext, ScanResponse.Builder builder) { 3341 if (numOfCompleteRows >= limitOfRows) { 3342 if (LOG.isTraceEnabled()) { 3343 LOG.trace("Done scanning, limit of rows reached, moreRows: " + moreRows 3344 + " scannerContext: " + scannerContext); 3345 } 3346 builder.setMoreResults(false); 3347 } 3348 } 3349 3350 // return whether we have more results in region. 3351 private void scan(HBaseRpcController controller, ScanRequest request, RegionScannerHolder rsh, 3352 long maxQuotaResultSize, int maxResults, int limitOfRows, List<Result> results, 3353 ScanResponse.Builder builder, RpcCall rpcCall, ServerSideScanMetrics scanMetrics) 3354 throws IOException { 3355 HRegion region = rsh.r; 3356 RegionScanner scanner = rsh.s; 3357 long maxResultSize; 3358 if (scanner.getMaxResultSize() > 0) { 3359 maxResultSize = Math.min(scanner.getMaxResultSize(), maxQuotaResultSize); 3360 } else { 3361 maxResultSize = maxQuotaResultSize; 3362 } 3363 // This is cells inside a row. Default size is 10 so if many versions or many cfs, 3364 // then we'll resize. Resizings show in profiler. Set it higher than 10. For now 3365 // arbitrary 32. TODO: keep record of general size of results being returned. 3366 ArrayList<Cell> values = new ArrayList<>(32); 3367 region.startRegionOperation(Operation.SCAN); 3368 long before = EnvironmentEdgeManager.currentTime(); 3369 // Used to check if we've matched the row limit set on the Scan 3370 int numOfCompleteRows = 0; 3371 // Count of times we call nextRaw; can be > numOfCompleteRows. 3372 int numOfNextRawCalls = 0; 3373 try { 3374 int numOfResults = 0; 3375 synchronized (scanner) { 3376 boolean stale = (region.getRegionInfo().getReplicaId() != 0); 3377 boolean clientHandlesPartials = 3378 request.hasClientHandlesPartials() && request.getClientHandlesPartials(); 3379 boolean clientHandlesHeartbeats = 3380 request.hasClientHandlesHeartbeats() && request.getClientHandlesHeartbeats(); 3381 3382 // On the server side we must ensure that the correct ordering of partial results is 3383 // returned to the client to allow them to properly reconstruct the partial results. 3384 // If the coprocessor host is adding to the result list, we cannot guarantee the 3385 // correct ordering of partial results and so we prevent partial results from being 3386 // formed. 3387 boolean serverGuaranteesOrderOfPartials = results.isEmpty(); 3388 boolean allowPartialResults = clientHandlesPartials && serverGuaranteesOrderOfPartials; 3389 boolean moreRows = false; 3390 3391 // Heartbeat messages occur when the processing of the ScanRequest is exceeds a 3392 // certain time threshold on the server. When the time threshold is exceeded, the 3393 // server stops the scan and sends back whatever Results it has accumulated within 3394 // that time period (may be empty). Since heartbeat messages have the potential to 3395 // create partial Results (in the event that the timeout occurs in the middle of a 3396 // row), we must only generate heartbeat messages when the client can handle both 3397 // heartbeats AND partials 3398 boolean allowHeartbeatMessages = clientHandlesHeartbeats && allowPartialResults; 3399 3400 long timeLimit = getTimeLimit(rpcCall, controller, allowHeartbeatMessages); 3401 3402 final LimitScope sizeScope = 3403 allowPartialResults ? LimitScope.BETWEEN_CELLS : LimitScope.BETWEEN_ROWS; 3404 final LimitScope timeScope = 3405 allowHeartbeatMessages ? LimitScope.BETWEEN_CELLS : LimitScope.BETWEEN_ROWS; 3406 3407 // Configure with limits for this RPC. Set keep progress true since size progress 3408 // towards size limit should be kept between calls to nextRaw 3409 ScannerContext.Builder contextBuilder = ScannerContext.newBuilder(true); 3410 // maxResultSize - either we can reach this much size for all cells(being read) data or sum 3411 // of heap size occupied by cells(being read). Cell data means its key and value parts. 3412 // maxQuotaResultSize - max results just from server side configuration and quotas, without 3413 // user's specified max. We use this for evaluating limits based on blocks (not cells). 3414 // We may have accumulated some results in coprocessor preScannerNext call. Subtract any 3415 // cell or block size from maximum here so we adhere to total limits of request. 3416 // Note: we track block size in StoreScanner. If the CP hook got cells from hbase, it will 3417 // have accumulated block bytes. If not, this will be 0 for block size. 3418 long maxCellSize = maxResultSize; 3419 long maxBlockSize = maxQuotaResultSize; 3420 if (rpcCall != null) { 3421 maxBlockSize -= rpcCall.getBlockBytesScanned(); 3422 maxCellSize -= rpcCall.getResponseCellSize(); 3423 } 3424 3425 contextBuilder.setSizeLimit(sizeScope, maxCellSize, maxCellSize, maxBlockSize); 3426 contextBuilder.setBatchLimit(scanner.getBatch()); 3427 contextBuilder.setTimeLimit(timeScope, timeLimit); 3428 contextBuilder.setTrackMetrics(scanMetrics != null); 3429 contextBuilder.setScanMetrics(scanMetrics); 3430 ScannerContext scannerContext = contextBuilder.build(); 3431 boolean limitReached = false; 3432 long blockBytesScannedBefore = 0; 3433 while (numOfResults < maxResults) { 3434 // Reset the batch progress to 0 before every call to RegionScanner#nextRaw. The 3435 // batch limit is a limit on the number of cells per Result. Thus, if progress is 3436 // being tracked (i.e. scannerContext.keepProgress() is true) then we need to 3437 // reset the batch progress between nextRaw invocations since we don't want the 3438 // batch progress from previous calls to affect future calls 3439 scannerContext.setBatchProgress(0); 3440 assert values.isEmpty(); 3441 3442 // Collect values to be returned here 3443 moreRows = scanner.nextRaw(values, scannerContext); 3444 3445 long blockBytesScanned = scannerContext.getBlockSizeProgress() - blockBytesScannedBefore; 3446 blockBytesScannedBefore = scannerContext.getBlockSizeProgress(); 3447 3448 if (rpcCall == null) { 3449 // When there is no RpcCallContext,copy EC to heap, then the scanner would close, 3450 // This can be an EXPENSIVE call. It may make an extra copy from offheap to onheap 3451 // buffers.See more details in HBASE-26036. 3452 CellUtil.cloneIfNecessary(values); 3453 } 3454 numOfNextRawCalls++; 3455 3456 if (!values.isEmpty()) { 3457 if (limitOfRows > 0) { 3458 // First we need to check if the last result is partial and we have a row change. If 3459 // so then we need to increase the numOfCompleteRows. 3460 if (results.isEmpty()) { 3461 if ( 3462 rsh.rowOfLastPartialResult != null 3463 && !CellUtil.matchingRows(values.get(0), rsh.rowOfLastPartialResult) 3464 ) { 3465 numOfCompleteRows++; 3466 checkLimitOfRows(numOfCompleteRows, limitOfRows, moreRows, scannerContext, 3467 builder); 3468 } 3469 } else { 3470 Result lastResult = results.get(results.size() - 1); 3471 if ( 3472 lastResult.mayHaveMoreCellsInRow() 3473 && !CellUtil.matchingRows(values.get(0), lastResult.getRow()) 3474 ) { 3475 numOfCompleteRows++; 3476 checkLimitOfRows(numOfCompleteRows, limitOfRows, moreRows, scannerContext, 3477 builder); 3478 } 3479 } 3480 if (builder.hasMoreResults() && !builder.getMoreResults()) { 3481 break; 3482 } 3483 } 3484 boolean mayHaveMoreCellsInRow = scannerContext.mayHaveMoreCellsInRow(); 3485 Result r = Result.create(values, null, stale, mayHaveMoreCellsInRow); 3486 3487 if (request.getScan().getQueryMetricsEnabled()) { 3488 builder.addQueryMetrics(ClientProtos.QueryMetrics.newBuilder() 3489 .setBlockBytesScanned(blockBytesScanned).build()); 3490 } 3491 3492 results.add(r); 3493 numOfResults++; 3494 if (!mayHaveMoreCellsInRow && limitOfRows > 0) { 3495 numOfCompleteRows++; 3496 checkLimitOfRows(numOfCompleteRows, limitOfRows, moreRows, scannerContext, builder); 3497 if (builder.hasMoreResults() && !builder.getMoreResults()) { 3498 break; 3499 } 3500 } 3501 } else if (!moreRows && !results.isEmpty()) { 3502 // No more cells for the scan here, we need to ensure that the mayHaveMoreCellsInRow of 3503 // last result is false. Otherwise it's possible that: the first nextRaw returned 3504 // because BATCH_LIMIT_REACHED (BTW it happen to exhaust all cells of the scan),so the 3505 // last result's mayHaveMoreCellsInRow will be true. while the following nextRaw will 3506 // return with moreRows=false, which means moreResultsInRegion would be false, it will 3507 // be a contradictory state (HBASE-21206). 3508 int lastIdx = results.size() - 1; 3509 Result r = results.get(lastIdx); 3510 if (r.mayHaveMoreCellsInRow()) { 3511 results.set(lastIdx, ClientInternalHelper.createResult( 3512 ClientInternalHelper.getExtendedRawCells(r), r.getExists(), r.isStale(), false)); 3513 } 3514 } 3515 boolean sizeLimitReached = scannerContext.checkSizeLimit(LimitScope.BETWEEN_ROWS); 3516 boolean timeLimitReached = scannerContext.checkTimeLimit(LimitScope.BETWEEN_ROWS); 3517 boolean resultsLimitReached = numOfResults >= maxResults; 3518 limitReached = sizeLimitReached || timeLimitReached || resultsLimitReached; 3519 3520 if (limitReached || !moreRows) { 3521 // With block size limit, we may exceed size limit without collecting any results. 3522 // In this case we want to send heartbeat and/or cursor. We don't want to send heartbeat 3523 // or cursor if results were collected, for example for cell size or heap size limits. 3524 boolean sizeLimitReachedWithoutResults = sizeLimitReached && results.isEmpty(); 3525 // We only want to mark a ScanResponse as a heartbeat message in the event that 3526 // there are more values to be read server side. If there aren't more values, 3527 // marking it as a heartbeat is wasteful because the client will need to issue 3528 // another ScanRequest only to realize that they already have all the values 3529 if (moreRows && (timeLimitReached || sizeLimitReachedWithoutResults)) { 3530 // Heartbeat messages occur when the time limit has been reached, or size limit has 3531 // been reached before collecting any results. This can happen for heavily filtered 3532 // scans which scan over too many blocks. 3533 builder.setHeartbeatMessage(true); 3534 if (rsh.needCursor) { 3535 Cell cursorCell = scannerContext.getLastPeekedCell(); 3536 if (cursorCell != null) { 3537 builder.setCursor(ProtobufUtil.toCursor(cursorCell)); 3538 } 3539 } 3540 } 3541 break; 3542 } 3543 values.clear(); 3544 } 3545 if (rpcCall != null) { 3546 rpcCall.incrementResponseCellSize(scannerContext.getHeapSizeProgress()); 3547 } 3548 builder.setMoreResultsInRegion(moreRows); 3549 // Check to see if the client requested that we track metrics server side. If the 3550 // client requested metrics, retrieve the metrics from the scanner context. 3551 if (scanMetrics != null) { 3552 // rather than increment yet another counter in StoreScanner, just set the value here 3553 // from block size progress before writing into the response 3554 scanMetrics.setCounter(ServerSideScanMetrics.BLOCK_BYTES_SCANNED_KEY_METRIC_NAME, 3555 scannerContext.getBlockSizeProgress()); 3556 } 3557 } 3558 } finally { 3559 region.closeRegionOperation(); 3560 // Update serverside metrics, even on error. 3561 long end = EnvironmentEdgeManager.currentTime(); 3562 3563 long responseCellSize = 0; 3564 long blockBytesScanned = 0; 3565 if (rpcCall != null) { 3566 responseCellSize = rpcCall.getResponseCellSize(); 3567 blockBytesScanned = rpcCall.getBlockBytesScanned(); 3568 rsh.updateBlockBytesScanned(blockBytesScanned); 3569 } 3570 region.getMetrics().updateScan(); 3571 final MetricsRegionServer metricsRegionServer = server.getMetrics(); 3572 if (metricsRegionServer != null) { 3573 metricsRegionServer.updateScan(region, end - before, responseCellSize, blockBytesScanned); 3574 metricsRegionServer.updateReadQueryMeter(region, numOfNextRawCalls); 3575 } 3576 } 3577 // coprocessor postNext hook 3578 if (region.getCoprocessorHost() != null) { 3579 region.getCoprocessorHost().postScannerNext(scanner, results, maxResults, true); 3580 } 3581 } 3582 3583 /** 3584 * Scan data in a table. 3585 * @param controller the RPC controller 3586 * @param request the scan request 3587 */ 3588 @Override 3589 public ScanResponse scan(final RpcController controller, final ScanRequest request) 3590 throws ServiceException { 3591 if (controller != null && !(controller instanceof HBaseRpcController)) { 3592 throw new UnsupportedOperationException( 3593 "We only do " + "HBaseRpcControllers! FIX IF A PROBLEM: " + controller); 3594 } 3595 if (!request.hasScannerId() && !request.hasScan()) { 3596 throw new ServiceException( 3597 new DoNotRetryIOException("Missing required input: scannerId or scan")); 3598 } 3599 try { 3600 checkOpen(); 3601 } catch (IOException e) { 3602 if (request.hasScannerId()) { 3603 String scannerName = toScannerName(request.getScannerId()); 3604 if (LOG.isDebugEnabled()) { 3605 LOG.debug( 3606 "Server shutting down and client tried to access missing scanner " + scannerName); 3607 } 3608 final LeaseManager leaseManager = server.getLeaseManager(); 3609 if (leaseManager != null) { 3610 try { 3611 leaseManager.cancelLease(scannerName); 3612 } catch (LeaseException le) { 3613 // No problem, ignore 3614 if (LOG.isTraceEnabled()) { 3615 LOG.trace("Un-able to cancel lease of scanner. It could already be closed."); 3616 } 3617 } 3618 } 3619 } 3620 throw new ServiceException(e); 3621 } 3622 boolean trackMetrics = request.hasTrackScanMetrics() && request.getTrackScanMetrics(); 3623 ThreadLocalServerSideScanMetrics.setScanMetricsEnabled(trackMetrics); 3624 if (trackMetrics) { 3625 ThreadLocalServerSideScanMetrics.reset(); 3626 } 3627 requestCount.increment(); 3628 rpcScanRequestCount.increment(); 3629 RegionScannerContext rsx; 3630 ScanResponse.Builder builder = ScanResponse.newBuilder(); 3631 try { 3632 rsx = checkQuotaAndGetRegionScannerContext(request, builder); 3633 } catch (IOException e) { 3634 if (e == SCANNER_ALREADY_CLOSED) { 3635 // Now we will close scanner automatically if there are no more results for this region but 3636 // the old client will still send a close request to us. Just ignore it and return. 3637 return builder.build(); 3638 } 3639 throw new ServiceException(e); 3640 } 3641 String scannerName = rsx.scannerName; 3642 RegionScannerHolder rsh = rsx.holder; 3643 OperationQuota quota = rsx.quota; 3644 if (rsh.fullRegionScan) { 3645 rpcFullScanRequestCount.increment(); 3646 } 3647 HRegion region = rsh.r; 3648 LeaseManager.Lease lease; 3649 try { 3650 // Remove lease while its being processed in server; protects against case 3651 // where processing of request takes > lease expiration time. or null if none found. 3652 lease = server.getLeaseManager().removeLease(scannerName); 3653 } catch (LeaseException e) { 3654 throw new ServiceException(e); 3655 } 3656 if (request.hasRenew() && request.getRenew()) { 3657 // add back and return 3658 addScannerLeaseBack(lease); 3659 try { 3660 checkScanNextCallSeq(request, rsh); 3661 } catch (OutOfOrderScannerNextException e) { 3662 throw new ServiceException(e); 3663 } 3664 return builder.build(); 3665 } 3666 try { 3667 checkScanNextCallSeq(request, rsh); 3668 } catch (OutOfOrderScannerNextException e) { 3669 addScannerLeaseBack(lease); 3670 throw new ServiceException(e); 3671 } 3672 // Now we have increased the next call sequence. If we give client an error, the retry will 3673 // never success. So we'd better close the scanner and return a DoNotRetryIOException to client 3674 // and then client will try to open a new scanner. 3675 boolean closeScanner = request.hasCloseScanner() ? request.getCloseScanner() : false; 3676 int rows; // this is scan.getCaching 3677 if (request.hasNumberOfRows()) { 3678 rows = request.getNumberOfRows(); 3679 } else { 3680 rows = closeScanner ? 0 : 1; 3681 } 3682 RpcCall rpcCall = RpcServer.getCurrentCall().orElse(null); 3683 // now let's do the real scan. 3684 long maxQuotaResultSize = Math.min(maxScannerResultSize, quota.getMaxResultSize()); 3685 RegionScanner scanner = rsh.s; 3686 // this is the limit of rows for this scan, if we the number of rows reach this value, we will 3687 // close the scanner. 3688 int limitOfRows; 3689 if (request.hasLimitOfRows()) { 3690 limitOfRows = request.getLimitOfRows(); 3691 } else { 3692 limitOfRows = -1; 3693 } 3694 boolean scannerClosed = false; 3695 try { 3696 List<Result> results = new ArrayList<>(Math.min(rows, 512)); 3697 ServerSideScanMetrics scanMetrics = trackMetrics ? new ServerSideScanMetrics() : null; 3698 if (rows > 0) { 3699 boolean done = false; 3700 // Call coprocessor. Get region info from scanner. 3701 if (region.getCoprocessorHost() != null) { 3702 Boolean bypass = region.getCoprocessorHost().preScannerNext(scanner, results, rows); 3703 if (!results.isEmpty()) { 3704 for (Result r : results) { 3705 // add cell size from CP results so we can track response size and update limits 3706 // when calling scan below if !done. We'll also have tracked block size if the CP 3707 // got results from hbase, since StoreScanner tracks that for all calls automatically. 3708 addSize(rpcCall, r); 3709 } 3710 } 3711 if (bypass != null && bypass.booleanValue()) { 3712 done = true; 3713 } 3714 } 3715 if (!done) { 3716 scan((HBaseRpcController) controller, request, rsh, maxQuotaResultSize, rows, limitOfRows, 3717 results, builder, rpcCall, scanMetrics); 3718 } else { 3719 builder.setMoreResultsInRegion(!results.isEmpty()); 3720 } 3721 } else { 3722 // This is a open scanner call with numberOfRow = 0, so set more results in region to true. 3723 builder.setMoreResultsInRegion(true); 3724 } 3725 3726 quota.addScanResult(results); 3727 addResults(builder, results, (HBaseRpcController) controller, 3728 RegionReplicaUtil.isDefaultReplica(region.getRegionInfo()), 3729 isClientCellBlockSupport(rpcCall)); 3730 if (scanner.isFilterDone() && results.isEmpty()) { 3731 // If the scanner's filter - if any - is done with the scan 3732 // only set moreResults to false if the results is empty. This is used to keep compatible 3733 // with the old scan implementation where we just ignore the returned results if moreResults 3734 // is false. Can remove the isEmpty check after we get rid of the old implementation. 3735 builder.setMoreResults(false); 3736 } 3737 // Later we may close the scanner depending on this flag so here we need to make sure that we 3738 // have already set this flag. 3739 assert builder.hasMoreResultsInRegion(); 3740 // we only set moreResults to false in the above code, so set it to true if we haven't set it 3741 // yet. 3742 if (!builder.hasMoreResults()) { 3743 builder.setMoreResults(true); 3744 } 3745 if (builder.getMoreResults() && builder.getMoreResultsInRegion() && !results.isEmpty()) { 3746 // Record the last cell of the last result if it is a partial result 3747 // We need this to calculate the complete rows we have returned to client as the 3748 // mayHaveMoreCellsInRow is true does not mean that there will be extra cells for the 3749 // current row. We may filter out all the remaining cells for the current row and just 3750 // return the cells of the nextRow when calling RegionScanner.nextRaw. So here we need to 3751 // check for row change. 3752 Result lastResult = results.get(results.size() - 1); 3753 if (lastResult.mayHaveMoreCellsInRow()) { 3754 rsh.rowOfLastPartialResult = lastResult.getRow(); 3755 } else { 3756 rsh.rowOfLastPartialResult = null; 3757 } 3758 } 3759 if (!builder.getMoreResults() || !builder.getMoreResultsInRegion() || closeScanner) { 3760 scannerClosed = true; 3761 closeScanner(region, scanner, scannerName, rpcCall, false); 3762 } 3763 3764 // There's no point returning to a timed out client. Throwing ensures scanner is closed 3765 if (rpcCall != null && EnvironmentEdgeManager.currentTime() > rpcCall.getDeadline()) { 3766 throw new TimeoutIOException("Client deadline exceeded, cannot return results"); 3767 } 3768 3769 if (scanMetrics != null) { 3770 if (rpcCall != null) { 3771 long rpcScanTime = EnvironmentEdgeManager.currentTime() - rpcCall.getStartTime(); 3772 long rpcQueueWaitTime = rpcCall.getStartTime() - rpcCall.getReceiveTime(); 3773 scanMetrics.addToCounter(ServerSideScanMetrics.RPC_SCAN_PROCESSING_TIME_METRIC_NAME, 3774 rpcScanTime); 3775 scanMetrics.addToCounter(ServerSideScanMetrics.RPC_SCAN_QUEUE_WAIT_TIME_METRIC_NAME, 3776 rpcQueueWaitTime); 3777 } 3778 ThreadLocalServerSideScanMetrics.populateServerSideScanMetrics(scanMetrics); 3779 Map<String, Long> metrics = scanMetrics.getMetricsMap(); 3780 ScanMetrics.Builder metricBuilder = ScanMetrics.newBuilder(); 3781 NameInt64Pair.Builder pairBuilder = NameInt64Pair.newBuilder(); 3782 3783 for (Entry<String, Long> entry : metrics.entrySet()) { 3784 pairBuilder.setName(entry.getKey()); 3785 pairBuilder.setValue(entry.getValue()); 3786 metricBuilder.addMetrics(pairBuilder.build()); 3787 } 3788 3789 builder.setScanMetrics(metricBuilder.build()); 3790 } 3791 3792 return builder.build(); 3793 } catch (IOException e) { 3794 try { 3795 // scanner is closed here 3796 scannerClosed = true; 3797 // The scanner state might be left in a dirty state, so we will tell the Client to 3798 // fail this RPC and close the scanner while opening up another one from the start of 3799 // row that the client has last seen. 3800 closeScanner(region, scanner, scannerName, rpcCall, true); 3801 3802 // If it is a DoNotRetryIOException already, throw as it is. Unfortunately, DNRIOE is 3803 // used in two different semantics. 3804 // (1) The first is to close the client scanner and bubble up the exception all the way 3805 // to the application. This is preferred when the exception is really un-recoverable 3806 // (like CorruptHFileException, etc). Plain DoNotRetryIOException also falls into this 3807 // bucket usually. 3808 // (2) Second semantics is to close the current region scanner only, but continue the 3809 // client scanner by overriding the exception. This is usually UnknownScannerException, 3810 // OutOfOrderScannerNextException, etc where the region scanner has to be closed, but the 3811 // application-level ClientScanner has to continue without bubbling up the exception to 3812 // the client. See ClientScanner code to see how it deals with these special exceptions. 3813 if (e instanceof DoNotRetryIOException) { 3814 throw e; 3815 } 3816 3817 // If it is a FileNotFoundException, wrap as a 3818 // DoNotRetryIOException. This can avoid the retry in ClientScanner. 3819 if (e instanceof FileNotFoundException) { 3820 throw new DoNotRetryIOException(e); 3821 } 3822 3823 // We closed the scanner already. Instead of throwing the IOException, and client 3824 // retrying with the same scannerId only to get USE on the next RPC, we directly throw 3825 // a special exception to save an RPC. 3826 if (VersionInfoUtil.hasMinimumVersion(rpcCall.getClientVersionInfo(), 1, 4)) { 3827 // 1.4.0+ clients know how to handle 3828 throw new ScannerResetException("Scanner is closed on the server-side", e); 3829 } else { 3830 // older clients do not know about SRE. Just throw USE, which they will handle 3831 throw new UnknownScannerException("Throwing UnknownScannerException to reset the client" 3832 + " scanner state for clients older than 1.3.", e); 3833 } 3834 } catch (IOException ioe) { 3835 throw new ServiceException(ioe); 3836 } 3837 } finally { 3838 if (!scannerClosed) { 3839 // Adding resets expiration time on lease. 3840 // the closeCallBack will be set in closeScanner so here we only care about shippedCallback 3841 if (rpcCall != null) { 3842 rpcCall.setCallBack(rsh.shippedCallback); 3843 } else { 3844 // If context is null,here we call rsh.shippedCallback directly to reuse the logic in 3845 // rsh.shippedCallback to release the internal resources in rsh,and lease is also added 3846 // back to regionserver's LeaseManager in rsh.shippedCallback. 3847 runShippedCallback(rsh); 3848 } 3849 } 3850 quota.close(); 3851 } 3852 } 3853 3854 private void runShippedCallback(RegionScannerHolder rsh) throws ServiceException { 3855 assert rsh.shippedCallback != null; 3856 try { 3857 rsh.shippedCallback.run(); 3858 } catch (IOException ioe) { 3859 throw new ServiceException(ioe); 3860 } 3861 } 3862 3863 private void closeScanner(HRegion region, RegionScanner scanner, String scannerName, 3864 RpcCallContext context, boolean isError) throws IOException { 3865 if (region.getCoprocessorHost() != null) { 3866 if (region.getCoprocessorHost().preScannerClose(scanner)) { 3867 // bypass the actual close. 3868 return; 3869 } 3870 } 3871 RegionScannerHolder rsh = scanners.remove(scannerName); 3872 if (rsh != null) { 3873 if (context != null) { 3874 context.setCallBack(rsh.closeCallBack); 3875 } else { 3876 rsh.s.close(); 3877 } 3878 if (region.getCoprocessorHost() != null) { 3879 region.getCoprocessorHost().postScannerClose(scanner); 3880 } 3881 if (!isError) { 3882 closedScanners.put(scannerName, rsh.getNextCallSeq()); 3883 } 3884 } 3885 } 3886 3887 @Override 3888 public CoprocessorServiceResponse execRegionServerService(RpcController controller, 3889 CoprocessorServiceRequest request) throws ServiceException { 3890 rpcPreCheck("execRegionServerService"); 3891 return server.execRegionServerService(controller, request); 3892 } 3893 3894 @Override 3895 public GetSpaceQuotaSnapshotsResponse getSpaceQuotaSnapshots(RpcController controller, 3896 GetSpaceQuotaSnapshotsRequest request) throws ServiceException { 3897 try { 3898 final RegionServerSpaceQuotaManager manager = server.getRegionServerSpaceQuotaManager(); 3899 final GetSpaceQuotaSnapshotsResponse.Builder builder = 3900 GetSpaceQuotaSnapshotsResponse.newBuilder(); 3901 if (manager != null) { 3902 final Map<TableName, SpaceQuotaSnapshot> snapshots = manager.copyQuotaSnapshots(); 3903 for (Entry<TableName, SpaceQuotaSnapshot> snapshot : snapshots.entrySet()) { 3904 builder.addSnapshots(TableQuotaSnapshot.newBuilder() 3905 .setTableName(ProtobufUtil.toProtoTableName(snapshot.getKey())) 3906 .setSnapshot(SpaceQuotaSnapshot.toProtoSnapshot(snapshot.getValue())).build()); 3907 } 3908 } 3909 return builder.build(); 3910 } catch (Exception e) { 3911 throw new ServiceException(e); 3912 } 3913 } 3914 3915 @Override 3916 public ClearRegionBlockCacheResponse clearRegionBlockCache(RpcController controller, 3917 ClearRegionBlockCacheRequest request) throws ServiceException { 3918 try { 3919 rpcPreCheck("clearRegionBlockCache"); 3920 ClearRegionBlockCacheResponse.Builder builder = ClearRegionBlockCacheResponse.newBuilder(); 3921 CacheEvictionStatsBuilder stats = CacheEvictionStats.builder(); 3922 server.getRegionServerCoprocessorHost().preClearRegionBlockCache(); 3923 List<HRegion> regions = getRegions(request.getRegionList(), stats); 3924 for (HRegion region : regions) { 3925 try { 3926 stats = stats.append(this.server.clearRegionBlockCache(region)); 3927 } catch (Exception e) { 3928 stats.addException(region.getRegionInfo().getRegionName(), e); 3929 } 3930 } 3931 stats.withMaxCacheSize(server.getBlockCache().map(BlockCache::getMaxSize).orElse(0L)); 3932 server.getRegionServerCoprocessorHost().postClearRegionBlockCache(stats.build()); 3933 return builder.setStats(ProtobufUtil.toCacheEvictionStats(stats.build())).build(); 3934 } catch (IOException e) { 3935 throw new ServiceException(e); 3936 } 3937 } 3938 3939 private void executeOpenRegionProcedures(OpenRegionRequest request, 3940 Map<TableName, TableDescriptor> tdCache) { 3941 long masterSystemTime = request.hasMasterSystemTime() ? request.getMasterSystemTime() : -1; 3942 long initiatingMasterActiveTime = 3943 request.hasInitiatingMasterActiveTime() ? request.getInitiatingMasterActiveTime() : -1; 3944 for (RegionOpenInfo regionOpenInfo : request.getOpenInfoList()) { 3945 RegionInfo regionInfo = ProtobufUtil.toRegionInfo(regionOpenInfo.getRegion()); 3946 TableName tableName = regionInfo.getTable(); 3947 TableDescriptor tableDesc = tdCache.get(tableName); 3948 if (tableDesc == null) { 3949 try { 3950 tableDesc = server.getTableDescriptors().get(regionInfo.getTable()); 3951 } catch (IOException e) { 3952 // Here we do not fail the whole method since we also need deal with other 3953 // procedures, and we can not ignore this one, so we still schedule a 3954 // AssignRegionHandler and it will report back to master if we still can not get the 3955 // TableDescriptor. 3956 LOG.warn("Failed to get TableDescriptor of {}, will try again in the handler", 3957 regionInfo.getTable(), e); 3958 } 3959 if (tableDesc != null) { 3960 tdCache.put(tableName, tableDesc); 3961 } 3962 } 3963 if (regionOpenInfo.getFavoredNodesCount() > 0) { 3964 server.updateRegionFavoredNodesMapping(regionInfo.getEncodedName(), 3965 regionOpenInfo.getFavoredNodesList()); 3966 } 3967 long procId = regionOpenInfo.getOpenProcId(); 3968 if (server.submitRegionProcedure(procId)) { 3969 server.getExecutorService().submit(AssignRegionHandler.create(server, regionInfo, procId, 3970 tableDesc, masterSystemTime, initiatingMasterActiveTime)); 3971 } 3972 } 3973 } 3974 3975 private void executeCloseRegionProcedures(CloseRegionRequest request) { 3976 String encodedName; 3977 long initiatingMasterActiveTime = 3978 request.hasInitiatingMasterActiveTime() ? request.getInitiatingMasterActiveTime() : -1; 3979 try { 3980 encodedName = ProtobufUtil.getRegionEncodedName(request.getRegion()); 3981 } catch (DoNotRetryIOException e) { 3982 throw new UncheckedIOException("Should not happen", e); 3983 } 3984 ServerName destination = request.hasDestinationServer() 3985 ? ProtobufUtil.toServerName(request.getDestinationServer()) 3986 : null; 3987 long procId = request.getCloseProcId(); 3988 boolean evictCache = request.getEvictCache(); 3989 if (server.submitRegionProcedure(procId)) { 3990 server.getExecutorService().submit(UnassignRegionHandler.create(server, encodedName, procId, 3991 false, destination, evictCache, initiatingMasterActiveTime)); 3992 } 3993 } 3994 3995 private void executeProcedures(RemoteProcedureRequest request) { 3996 RSProcedureCallable callable; 3997 try { 3998 callable = Class.forName(request.getProcClass()).asSubclass(RSProcedureCallable.class) 3999 .getDeclaredConstructor().newInstance(); 4000 } catch (Exception e) { 4001 LOG.warn("Failed to instantiating remote procedure {}, pid={}", request.getProcClass(), 4002 request.getProcId(), e); 4003 server.remoteProcedureComplete(request.getProcId(), request.getInitiatingMasterActiveTime(), 4004 e, null); 4005 return; 4006 } 4007 callable.init(request.getProcData().toByteArray(), server); 4008 LOG.debug("Executing remote procedure {}, pid={}", callable.getClass(), request.getProcId()); 4009 server.executeProcedure(request.getProcId(), request.getInitiatingMasterActiveTime(), callable); 4010 } 4011 4012 @Override 4013 @QosPriority(priority = HConstants.ADMIN_QOS) 4014 public ExecuteProceduresResponse executeProcedures(RpcController controller, 4015 ExecuteProceduresRequest request) throws ServiceException { 4016 try { 4017 checkOpen(); 4018 throwOnWrongStartCode(request); 4019 server.getRegionServerCoprocessorHost().preExecuteProcedures(); 4020 if (request.getOpenRegionCount() > 0) { 4021 // Avoid reading from the TableDescritor every time(usually it will read from the file 4022 // system) 4023 Map<TableName, TableDescriptor> tdCache = new HashMap<>(); 4024 request.getOpenRegionList().forEach(req -> executeOpenRegionProcedures(req, tdCache)); 4025 } 4026 if (request.getCloseRegionCount() > 0) { 4027 request.getCloseRegionList().forEach(this::executeCloseRegionProcedures); 4028 } 4029 if (request.getProcCount() > 0) { 4030 request.getProcList().forEach(this::executeProcedures); 4031 } 4032 server.getRegionServerCoprocessorHost().postExecuteProcedures(); 4033 return ExecuteProceduresResponse.getDefaultInstance(); 4034 } catch (IOException e) { 4035 throw new ServiceException(e); 4036 } 4037 } 4038 4039 @Override 4040 public GetAllBootstrapNodesResponse getAllBootstrapNodes(RpcController controller, 4041 GetAllBootstrapNodesRequest request) throws ServiceException { 4042 GetAllBootstrapNodesResponse.Builder builder = GetAllBootstrapNodesResponse.newBuilder(); 4043 server.getBootstrapNodes() 4044 .forEachRemaining(server -> builder.addNode(ProtobufUtil.toServerName(server))); 4045 return builder.build(); 4046 } 4047 4048 private void setReloadableGuardrails(Configuration conf) { 4049 rowSizeWarnThreshold = 4050 conf.getInt(HConstants.BATCH_ROWS_THRESHOLD_NAME, HConstants.BATCH_ROWS_THRESHOLD_DEFAULT); 4051 rejectRowsWithSizeOverThreshold = 4052 conf.getBoolean(REJECT_BATCH_ROWS_OVER_THRESHOLD, DEFAULT_REJECT_BATCH_ROWS_OVER_THRESHOLD); 4053 maxScannerResultSize = conf.getLong(HConstants.HBASE_SERVER_SCANNER_MAX_RESULT_SIZE_KEY, 4054 HConstants.DEFAULT_HBASE_SERVER_SCANNER_MAX_RESULT_SIZE); 4055 } 4056 4057 @Override 4058 public void onConfigurationChange(Configuration conf) { 4059 super.onConfigurationChange(conf); 4060 setReloadableGuardrails(conf); 4061 } 4062 4063 @Override 4064 public GetCachedFilesListResponse getCachedFilesList(RpcController controller, 4065 GetCachedFilesListRequest request) throws ServiceException { 4066 GetCachedFilesListResponse.Builder responseBuilder = GetCachedFilesListResponse.newBuilder(); 4067 List<String> fullyCachedFiles = new ArrayList<>(); 4068 server.getBlockCache().flatMap(BlockCache::getFullyCachedFiles).ifPresent(fcf -> { 4069 fullyCachedFiles.addAll(fcf.keySet()); 4070 }); 4071 return responseBuilder.addAllCachedFiles(fullyCachedFiles).build(); 4072 } 4073 4074 /** 4075 * STUB - Refreshes the system key cache on the region server. Feature not yet implemented in 4076 * precursor PR. 4077 */ 4078 @Override 4079 @QosPriority(priority = HConstants.ADMIN_QOS) 4080 public EmptyMsg refreshSystemKeyCache(final RpcController controller, final EmptyMsg request) 4081 throws ServiceException { 4082 throw new ServiceException( 4083 new UnsupportedOperationException("Key management feature not yet implemented")); 4084 } 4085 4086 /** 4087 * STUB - Ejects a specific managed key entry from the cache. Feature not yet implemented in 4088 * precursor PR. 4089 */ 4090 @Override 4091 @QosPriority(priority = HConstants.ADMIN_QOS) 4092 public BooleanMsg ejectManagedKeyDataCacheEntry(final RpcController controller, 4093 final ManagedKeyEntryRequest request) throws ServiceException { 4094 throw new ServiceException( 4095 new UnsupportedOperationException("Key management feature not yet implemented")); 4096 } 4097 4098 /** 4099 * STUB - Clears all entries in the managed key data cache. Feature not yet implemented in 4100 * precursor PR. 4101 */ 4102 @Override 4103 @QosPriority(priority = HConstants.ADMIN_QOS) 4104 public EmptyMsg clearManagedKeyDataCache(final RpcController controller, final EmptyMsg request) 4105 throws ServiceException { 4106 throw new ServiceException( 4107 new UnsupportedOperationException("Key management feature not yet implemented")); 4108 } 4109 4110 RegionScannerContext checkQuotaAndGetRegionScannerContext(ScanRequest request, 4111 ScanResponse.Builder builder) throws IOException { 4112 if (request.hasScannerId()) { 4113 // The downstream projects such as AsyncHBase in OpenTSDB need this value. See HBASE-18000 4114 // for more details. 4115 long scannerId = request.getScannerId(); 4116 builder.setScannerId(scannerId); 4117 String scannerName = toScannerName(scannerId); 4118 RegionScannerHolder rsh = getRegionScanner(request); 4119 OperationQuota quota = 4120 getRpcQuotaManager().checkScanQuota(rsh.r, request, maxScannerResultSize, 4121 rsh.getMaxBlockBytesScanned(), rsh.getPrevBlockBytesScannedDifference()); 4122 return new RegionScannerContext(scannerName, rsh, quota); 4123 } 4124 4125 HRegion region = getRegion(request.getRegion()); 4126 OperationQuota quota = 4127 getRpcQuotaManager().checkScanQuota(region, request, maxScannerResultSize, 0L, 0L); 4128 Pair<String, RegionScannerHolder> pair = newRegionScanner(request, region, builder); 4129 return new RegionScannerContext(pair.getFirst(), pair.getSecond(), quota); 4130 } 4131 4132 private List<SlowLogPayload> getSlowLogPayloads(SlowLogResponseRequest request, 4133 NamedQueueRecorder namedQueueRecorder) { 4134 if (namedQueueRecorder == null) { 4135 return Collections.emptyList(); 4136 } 4137 List<SlowLogPayload> slowLogPayloads; 4138 NamedQueueGetRequest namedQueueGetRequest = new NamedQueueGetRequest(); 4139 namedQueueGetRequest.setNamedQueueEvent(RpcLogDetails.SLOW_LOG_EVENT); 4140 namedQueueGetRequest.setSlowLogResponseRequest(request); 4141 NamedQueueGetResponse namedQueueGetResponse = 4142 namedQueueRecorder.getNamedQueueRecords(namedQueueGetRequest); 4143 slowLogPayloads = namedQueueGetResponse != null 4144 ? namedQueueGetResponse.getSlowLogPayloads() 4145 : Collections.emptyList(); 4146 return slowLogPayloads; 4147 } 4148 4149 @Override 4150 @QosPriority(priority = HConstants.ADMIN_QOS) 4151 public HBaseProtos.LogEntry getLogEntries(RpcController controller, 4152 HBaseProtos.LogRequest request) throws ServiceException { 4153 try { 4154 requirePermission("getLogEntries", Permission.Action.ADMIN); 4155 final String logClassName = request.getLogClassName(); 4156 Class<?> logClass = Class.forName(logClassName).asSubclass(Message.class); 4157 Method method = logClass.getMethod("parseFrom", ByteString.class); 4158 if (logClassName.contains("SlowLogResponseRequest")) { 4159 SlowLogResponseRequest slowLogResponseRequest = 4160 (SlowLogResponseRequest) method.invoke(null, request.getLogMessage()); 4161 final NamedQueueRecorder namedQueueRecorder = this.server.getNamedQueueRecorder(); 4162 final List<SlowLogPayload> slowLogPayloads = 4163 getSlowLogPayloads(slowLogResponseRequest, namedQueueRecorder); 4164 SlowLogResponses slowLogResponses = 4165 SlowLogResponses.newBuilder().addAllSlowLogPayloads(slowLogPayloads).build(); 4166 return HBaseProtos.LogEntry.newBuilder() 4167 .setLogClassName(slowLogResponses.getClass().getName()) 4168 .setLogMessage(slowLogResponses.toByteString()).build(); 4169 } 4170 } catch (ClassNotFoundException | NoSuchMethodException | IllegalAccessException 4171 | InvocationTargetException | IOException e) { 4172 LOG.error("Error while retrieving log entries.", e); 4173 throw new ServiceException(e); 4174 } 4175 throw new ServiceException("Invalid request params"); 4176 } 4177}