2011-04-05 00:08:24 -04:00
#
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#
# Moves regions. Will confirm region access in current location and will
# not move a new region until successful confirm of region loading in new
# location. Presumes balancer is disabled when we run (not harmful if its
# on but this script and balancer will end up fighting each other).
require 'optparse'
2014-02-26 18:05:50 -05:00
require File . join ( File . dirname ( __FILE__ ) , 'thread-pool' )
2011-04-05 00:08:24 -04:00
include Java
import org . apache . hadoop . hbase . HConstants
import org . apache . hadoop . hbase . HBaseConfiguration
import org . apache . hadoop . hbase . client . HBaseAdmin
import org . apache . hadoop . hbase . client . Get
import org . apache . hadoop . hbase . client . Scan
import org . apache . hadoop . hbase . client . HTable
import org . apache . hadoop . hbase . client . HConnectionManager
import org . apache . hadoop . hbase . filter . FirstKeyOnlyFilter ;
import org . apache . hadoop . hbase . util . Bytes
import org . apache . hadoop . hbase . util . Writables
import org . apache . hadoop . conf . Configuration
import org . apache . commons . logging . Log
import org . apache . commons . logging . LogFactory
2013-03-26 11:03:37 -04:00
import org . apache . hadoop . hbase . protobuf . ProtobufUtil
import org . apache . hadoop . hbase . ServerName
import org . apache . hadoop . hbase . HRegionInfo
2011-04-05 00:08:24 -04:00
# Name of this script
NAME = " region_mover "
# Get meta table reference
def getMetaTable ( config )
# Keep meta reference in ruby global
if not $META
$META = HTable . new ( config , HConstants :: META_TABLE_NAME )
end
return $META
end
# Get table instance.
# Maintains cache of table instances.
def getTable ( config , name )
# Keep dictionary of tables in ruby global
if not $TABLES
$TABLES = { }
end
2013-08-15 21:20:59 -04:00
key = name . toString ( )
2011-04-05 00:08:24 -04:00
if not $TABLES [ key ]
$TABLES [ key ] = HTable . new ( config , name )
end
return $TABLES [ key ]
end
2013-11-28 11:23:26 -05:00
def closeTables ( )
if not $TABLES
return
end
$LOG . info ( " Close all tables " )
$TABLES . each do | name , table |
$TABLES . delete ( name )
table . close ( )
end
end
2011-04-05 00:08:24 -04:00
2014-05-27 13:36:55 -04:00
# Returns true if passed region is still on 'original' when we look at hbase:meta.
2011-04-05 00:08:24 -04:00
def isSameServer ( admin , r , original )
server = getServerNameForRegion ( admin , r )
2014-02-07 12:57:44 -05:00
return false unless server and original
2011-04-05 00:08:24 -04:00
return server == original
end
class RubyAbortable
include org . apache . hadoop . hbase . Abortable
def abort ( why , e )
puts " ABORTED! why= " + why + " , e= " + e . to_s
end
end
2014-05-27 13:36:55 -04:00
# Get servername that is up in hbase:meta; this is hostname + port + startcode comma-delimited.
2011-04-05 00:08:24 -04:00
# Can return nil
def getServerNameForRegion ( admin , r )
2014-02-07 12:57:44 -05:00
return nil unless admin . isTableEnabled ( r . getTableName )
2013-03-26 11:03:37 -04:00
if r . isMetaRegion ( )
2011-04-05 00:08:24 -04:00
# Hack
2013-10-09 16:26:03 -04:00
zkw = org . apache . hadoop . hbase . zookeeper . ZooKeeperWatcher . new ( admin . getConfiguration ( ) , " region_mover " , nil )
2014-07-02 22:01:04 -04:00
mtl = org . apache . hadoop . hbase . zookeeper . MetaTableLocator . new ( )
2013-10-09 16:26:03 -04:00
begin
2014-07-02 22:01:04 -04:00
while not mtl . isLocationAvailable ( zkw )
2013-10-09 16:26:03 -04:00
sleep 0 . 1
end
# Make a fake servername by appending ','
2014-07-02 22:01:04 -04:00
metaServer = mtl . getMetaRegionLocation ( zkw ) . toString ( ) + " , "
2013-10-09 16:26:03 -04:00
return metaServer
ensure
zkw . close ( )
2011-04-05 00:08:24 -04:00
end
end
table = nil
2013-03-26 11:03:37 -04:00
table = getMetaTable ( admin . getConfiguration ( ) )
2011-04-05 00:08:24 -04:00
g = Get . new ( r . getRegionName ( ) )
g . addColumn ( HConstants :: CATALOG_FAMILY , HConstants :: SERVER_QUALIFIER )
g . addColumn ( HConstants :: CATALOG_FAMILY , HConstants :: STARTCODE_QUALIFIER )
result = table . get ( g )
2014-02-07 12:57:44 -05:00
return nil unless result
2011-04-05 00:08:24 -04:00
server = result . getValue ( HConstants :: CATALOG_FAMILY , HConstants :: SERVER_QUALIFIER )
startcode = result . getValue ( HConstants :: CATALOG_FAMILY , HConstants :: STARTCODE_QUALIFIER )
return nil unless server
return java . lang . String . new ( Bytes . toString ( server ) ) . replaceFirst ( " : " , " , " ) + " , " + Bytes . toLong ( startcode ) . to_s
end
# Trys to scan a row from passed region
# Throws exception if can't
def isSuccessfulScan ( admin , r )
scan = Scan . new ( r . getStartKey ( ) )
scan . setBatch ( 1 )
scan . setCaching ( 1 )
scan . setFilter ( FirstKeyOnlyFilter . new ( ) )
2014-02-07 12:57:44 -05:00
begin
table = getTable ( admin . getConfiguration ( ) , r . getTableName ( ) )
scanner = table . getScanner ( scan )
rescue org . apache . hadoop . hbase . TableNotFoundException ,
org . apache . hadoop . hbase . TableNotEnabledException = > e
$LOG . warn ( " Region " + r . getEncodedName ( ) + " belongs to recently " +
" deleted/disabled table. Skipping... " + e . message )
return
end
2011-04-05 00:08:24 -04:00
begin
results = scanner . next ( )
# We might scan into next region, this might be an empty table.
# But if no exception, presume scanning is working.
ensure
scanner . close ( )
2013-11-28 11:23:26 -05:00
# Do not close the htable. It is cached in $TABLES and
# may be reused in moving another region of same table.
# table.close()
2011-04-05 00:08:24 -04:00
end
end
# Check region has moved successful and is indeed hosted on another server
# Wait until that is the case.
def move ( admin , r , newServer , original )
# Now move it. Do it in a loop so can retry if fail. Have seen issue where
# we tried move region but failed and retry put it back on old location;
# retry in this case.
2014-02-26 18:05:50 -05:00
2011-04-05 00:08:24 -04:00
retries = admin . getConfiguration . getInt ( " hbase.move.retries.max " , 5 )
count = 0
same = true
2013-11-28 11:23:26 -05:00
start = Time . now
2011-04-05 00:08:24 -04:00
while count < retries and same
if count > 0
$LOG . info ( " Retry " + count . to_s + " of maximum " + retries . to_s )
end
count = count + 1
begin
admin . move ( Bytes . toBytes ( r . getEncodedName ( ) ) , Bytes . toBytes ( newServer ) )
2014-02-07 12:57:44 -05:00
rescue java . lang . reflect . UndeclaredThrowableException ,
org . apache . hadoop . hbase . UnknownRegionException = > e
2011-04-05 00:08:24 -04:00
$LOG . info ( " Exception moving " + r . getEncodedName ( ) +
" ; split/moved? Continuing: " + e )
return
end
# Wait till its up on new server before moving on
maxWaitInSeconds = admin . getConfiguration . getInt ( " hbase.move.wait.max " , 60 )
maxWait = Time . now + maxWaitInSeconds
while Time . now < maxWait
same = isSameServer ( admin , r , original )
break unless same
sleep 0 . 1
end
end
raise RuntimeError , " Region stuck on #{ original } , newserver= #{ newServer } " if same
# Assert can Scan from new location.
isSuccessfulScan ( admin , r )
2013-11-28 11:23:26 -05:00
$LOG . info ( " Moved region " + r . getRegionNameAsString ( ) + " cost: " +
java . lang . String . format ( " %.3f " , ( Time . now - start ) ) )
2011-04-05 00:08:24 -04:00
end
2014-09-17 19:24:53 -04:00
# Return the hostname:port out of a servername (all up to second ',')
def getHostPortFromServerName ( serverName )
return serverName . split ( ',' ) [ 0 .. 1 ]
2011-04-05 00:08:24 -04:00
end
# Return array of servernames where servername is hostname+port+startcode
# comma-delimited
def getServers ( admin )
serverInfos = admin . getClusterStatus ( ) . getServerInfo ( )
servers = [ ]
for server in serverInfos
servers << server . getServerName ( )
end
return servers
end
# Remove the servername whose hostname portion matches from the passed
# array of servers. Returns as side-effect the servername removed.
2014-09-17 19:24:53 -04:00
def stripServer ( servers , hostname , port )
2011-04-05 00:08:24 -04:00
count = servers . length
servername = nil
for server in servers
2014-09-17 19:24:53 -04:00
hostFromServerName , portFromServerName = getHostPortFromServerName ( server )
if hostFromServerName == hostname and portFromServerName == port
2011-04-05 00:08:24 -04:00
servername = servers . delete ( server )
end
end
# Check server to exclude is actually present
2014-09-17 19:24:53 -04:00
raise RuntimeError , " Server %s:%d not online " % [ hostname , port ] unless servers . length < count
2011-04-05 00:08:24 -04:00
return servername
end
2012-07-06 18:12:16 -04:00
# Returns a new serverlist that excludes the servername whose hostname portion
# matches from the passed array of servers.
def stripExcludes ( servers , excludefile )
excludes = readExcludes ( excludefile )
2014-09-17 19:24:53 -04:00
servers = servers . find_all { | server |
! excludes . contains ( getHostPortFromServerName ( server ) . join ( " : " ) )
}
2012-07-06 18:12:16 -04:00
# return updated servers list
return servers
end
2014-09-17 19:24:53 -04:00
# Return servername that matches passed hostname and port
def getServerName ( servers , hostname , port )
2011-04-05 00:08:24 -04:00
servername = nil
for server in servers
2014-09-17 19:24:53 -04:00
hostFromServerName , portFromServerName = getHostPortFromServerName ( server )
if hostFromServerName == hostname and portFromServerName == port
2011-04-05 00:08:24 -04:00
servername = server
break
end
end
2014-09-17 19:24:53 -04:00
raise ArgumentError , " Server %s:%d not online " % [ hostname , port ] unless servername
2011-04-05 00:08:24 -04:00
return servername
end
# Create a logger and disable the DEBUG-level annoying client logging
def configureLogging ( options )
apacheLogger = LogFactory . getLog ( NAME )
# Configure log4j to not spew so much
unless ( options [ :debug ] )
logger = org . apache . log4j . Logger . getLogger ( " org.apache.hadoop.hbase.client " )
logger . setLevel ( org . apache . log4j . Level :: INFO )
end
return apacheLogger
end
# Get configuration instance
def getConfiguration ( )
config = HBaseConfiguration . create ( )
2014-05-27 13:36:55 -04:00
# No prefetching on hbase:meta This is for versions pre 0.99. Newer versions do not prefetch.
2011-04-05 00:08:24 -04:00
config . setInt ( " hbase.client.prefetch.limit " , 1 )
# Make a config that retries at short intervals many times
config . setInt ( " hbase.client.pause " , 500 )
config . setInt ( " hbase.client.retries.number " , 100 )
return config
end
# Now get list of regions on targetServer
def getRegions ( config , servername )
2013-03-26 11:03:37 -04:00
connection = HConnectionManager :: getConnection ( config ) ;
2013-12-17 08:05:10 -05:00
return ProtobufUtil :: getOnlineRegions ( connection . getAdmin ( ServerName . valueOf ( servername ) ) ) ;
2011-04-05 00:08:24 -04:00
end
def deleteFile ( filename )
f = java . io . File . new ( filename )
f . delete ( ) if f . exists ( )
end
# Write HRegionInfo to file
# Need to serialize in case non-printable characters.
# Format is count of regionnames followed by serialized regionnames.
def writeFile ( filename , regions )
fos = java . io . FileOutputStream . new ( filename )
dos = java . io . DataOutputStream . new ( fos )
# Write out a count of region names
dos . writeInt ( regions . size ( ) )
# Write actual region names.
for r in regions
2013-03-26 11:03:37 -04:00
Bytes . writeByteArray ( dos , r . toByteArray ( ) )
2011-04-05 00:08:24 -04:00
end
dos . close ( )
end
# See writeFile above.
# Returns array of HRegionInfos
def readFile ( filename )
f = java . io . File . new ( filename )
return java . util . ArrayList . new ( ) unless f . exists ( )
fis = java . io . FileInputStream . new ( f )
dis = java . io . DataInputStream . new ( fis )
# Read count of regions
count = dis . readInt ( )
regions = java . util . ArrayList . new ( count )
index = 0
while index < count
2013-03-26 11:03:37 -04:00
regions . add ( HRegionInfo . parseFromOrNull ( Bytes . readByteArray ( dis ) ) )
2011-04-05 00:08:24 -04:00
index = index + 1
end
dis . close ( )
return regions
end
2014-09-17 19:24:53 -04:00
# Move regions off the passed hostname:port
def unloadRegions ( options , hostname , port )
2011-04-05 00:08:24 -04:00
# Get configuration
config = getConfiguration ( )
# Clean up any old files.
2014-09-17 19:24:53 -04:00
filename = getFilename ( options , hostname , port )
2011-04-05 00:08:24 -04:00
deleteFile ( filename )
# Get an admin instance
admin = HBaseAdmin . new ( config )
2014-09-17 19:24:53 -04:00
port = config . getInt ( HConstants :: REGIONSERVER_PORT , HConstants :: DEFAULT_REGIONSERVER_PORT ) \
unless port
2011-04-05 00:08:24 -04:00
servers = getServers ( admin )
# Remove the server we are unloading from from list of servers.
# Side-effect is the servername that matches this hostname
2014-09-17 19:24:53 -04:00
servername = stripServer ( servers , hostname , port )
2012-07-06 18:12:16 -04:00
# Remove the servers in our exclude list from list of servers.
servers = stripExcludes ( servers , options [ :excludesFile ] )
puts " Valid region move targets: " , servers
2014-02-28 12:30:23 -05:00
if servers . length == 0
puts " No regions were moved - there was no server available "
exit 4
end
2011-04-05 00:08:24 -04:00
movedRegions = java . util . ArrayList . new ( )
while true
rs = getRegions ( config , servername )
2014-02-07 12:57:44 -05:00
# Remove those already tried to move
rs . removeAll ( movedRegions )
2011-04-05 00:08:24 -04:00
break if rs . length == 0
count = 0
$LOG . info ( " Moving " + rs . length . to_s + " region(s) from " + servername +
2014-02-26 18:05:50 -05:00
" on " + servers . length . to_s + " servers using " + options [ :maxthreads ] . to_s + " threads. " )
counter = 0
pool = ThreadPool . new ( options [ :maxthreads ] )
server_index = 0
while counter < rs . length do
pool . launch ( rs , counter , server_index ) do | _rs , _counter , _server_index |
2014-03-03 14:06:51 -05:00
$LOG . info ( " Moving region " + _rs [ _counter ] . getEncodedName ( ) + " ( " + ( _counter + 1 ) . to_s +
2014-02-26 18:05:50 -05:00
" of " + _rs . length . to_s + " ) to server= " + servers [ _server_index ] + " for " + servername )
# Assert we can scan region in its current location
isSuccessfulScan ( admin , _rs [ _counter ] )
# Now move it.
move ( admin , _rs [ _counter ] , servers [ _server_index ] , servername )
movedRegions . add ( _rs [ _counter ] )
end
counter += 1
server_index = ( server_index + 1 ) % servers . length
2011-04-05 00:08:24 -04:00
end
2014-02-26 18:05:50 -05:00
$LOG . info ( " Waiting for the pool to complete " )
pool . stop
$LOG . info ( " Pool completed " )
2011-04-05 00:08:24 -04:00
end
if movedRegions . size ( ) > 0
# Write out file of regions moved
writeFile ( filename , movedRegions )
$LOG . info ( " Wrote list of moved regions to " + filename )
end
end
# Move regions to the passed hostname
2014-09-17 19:24:53 -04:00
def loadRegions ( options , hostname , port )
2011-04-05 00:08:24 -04:00
# Get configuration
config = getConfiguration ( )
# Get an admin instance
admin = HBaseAdmin . new ( config )
2014-09-17 19:24:53 -04:00
port = config . getInt ( HConstants :: REGIONSERVER_PORT , HConstants :: DEFAULT_REGIONSERVER_PORT ) \
unless port
filename = getFilename ( options , hostname , port )
2011-04-05 00:08:24 -04:00
regions = readFile ( filename )
return if regions . isEmpty ( )
servername = nil
# Wait till server is up
maxWaitInSeconds = admin . getConfiguration . getInt ( " hbase.serverstart.wait.max " , 180 )
maxWait = Time . now + maxWaitInSeconds
while Time . now < maxWait
servers = getServers ( admin )
begin
2014-09-17 19:24:53 -04:00
servername = getServerName ( servers , hostname , port )
2011-04-05 00:08:24 -04:00
rescue ArgumentError = > e
2014-09-17 19:24:53 -04:00
$LOG . info ( " hostname= " + hostname . to_s + " : " + port . to_s + " is not up yet, waiting " ) ;
2011-04-05 00:08:24 -04:00
end
break if servername
sleep 0 . 5
end
$LOG . info ( " Moving " + regions . size ( ) . to_s + " regions to " + servername )
count = 0
2013-11-28 11:23:26 -05:00
# sleep 20s to make sure the rs finished initialization.
sleep 20
2014-02-26 18:05:50 -05:00
counter = 0
pool = ThreadPool . new ( options [ :maxthreads ] )
while counter < regions . length do
r = regions [ counter ]
2011-04-05 00:08:24 -04:00
exists = false
begin
2012-03-01 12:53:03 -05:00
isSuccessfulScan ( admin , r )
exists = true
2013-08-21 16:21:14 -04:00
rescue org . apache . hadoop . hbase . NotServingRegionException = > e
2011-04-05 00:08:24 -04:00
$LOG . info ( " Failed scan of " + e . message )
end
next unless exists
currentServer = getServerNameForRegion ( admin , r )
if currentServer and currentServer == servername
$LOG . info ( " Region " + r . getRegionNameAsString ( ) + " ( " + count . to_s +
2014-02-26 18:05:50 -05:00
" of " + regions . length . to_s + " ) already on target server= " + servername )
2014-08-28 17:33:22 -04:00
counter = counter + 1
2011-04-05 00:08:24 -04:00
next
end
2014-02-26 18:05:50 -05:00
pool . launch ( r , currentServer , count ) do | _r , _currentServer , _count |
2014-03-03 14:06:51 -05:00
$LOG . info ( " Moving region " + _r . getRegionNameAsString ( ) + " ( " + ( _count + 1 ) . to_s +
2014-03-06 05:05:40 -05:00
" of " + regions . length . to_s + " ) from " + _currentServer . to_s + " to server= " +
2014-02-26 18:05:50 -05:00
servername ) ;
move ( admin , _r , servername , _currentServer )
end
counter = counter + 1
2011-04-05 00:08:24 -04:00
end
2014-02-26 18:05:50 -05:00
pool . stop
2011-04-05 00:08:24 -04:00
end
2012-07-06 18:12:16 -04:00
# Returns an array of hosts to exclude as region move targets
def readExcludes ( filename )
if filename == nil
return java . util . ArrayList . new ( )
end
if ! File . exist? ( filename )
puts " Error: Unable to read host exclude file: " , filename
raise RuntimeError
end
f = File . new ( filename , " r " )
# Read excluded hosts list
excludes = java . util . ArrayList . new ( )
while ( line = f . gets )
line . strip! # do an inplace drop of pre and post whitespaces
excludes . add ( line ) unless line . empty? # exclude empty lines
end
puts " Excluding hosts as region move targets: " , excludes
f . close
return excludes
end
2014-09-17 19:24:53 -04:00
def getFilename ( options , targetServer , port )
2011-04-05 00:08:24 -04:00
filename = options [ :file ]
if not filename
2014-09-17 19:24:53 -04:00
filename = " /tmp/ " + ENV [ 'USER' ] + targetServer + " : " + port
2011-04-05 00:08:24 -04:00
end
return filename
end
# Do command-line parsing
options = { }
optparse = OptionParser . new do | opts |
2014-09-17 19:24:53 -04:00
opts . banner = " Usage: #{ NAME } .rb [options] load|unload [<hostname>|<hostname:port>] "
2011-04-05 00:08:24 -04:00
opts . separator 'Load or unload regions by moving one at a time'
options [ :file ] = nil
2014-02-26 18:05:50 -05:00
options [ :maxthreads ] = 1
2014-09-17 19:24:53 -04:00
opts . on ( '-f' , '--filename=FILE' , 'File to save regions list into unloading, or read from loading; default /tmp/<hostname:port>' ) do | file |
2011-04-05 00:08:24 -04:00
options [ :file ] = file
end
opts . on ( '-h' , '--help' , 'Display usage information' ) do
puts opts
exit
end
options [ :debug ] = false
opts . on ( '-d' , '--debug' , 'Display extra debug logging' ) do
options [ :debug ] = true
end
2012-07-06 18:12:16 -04:00
opts . on ( '-x' , '--excludefile=FILE' , 'File with hosts-per-line to exclude as unload targets; default excludes only target host; useful for rack decommisioning.' ) do | file |
options [ :excludesFile ] = file
end
2014-02-26 18:05:50 -05:00
opts . on ( '-m' , '--maxthreads=XX' , 'Define the maximum number of threads to use to unload and reload the regions' ) do | number |
options [ :maxthreads ] = number . to_i
end
2011-04-05 00:08:24 -04:00
end
optparse . parse!
# Check ARGVs
if ARGV . length < 2
puts optparse
exit 1
end
2014-09-17 19:24:53 -04:00
hostname , port = ARGV [ 1 ] . split ( " : " )
2011-04-05 00:08:24 -04:00
if not hostname
opts optparse
exit 2
end
# Create a logger and save it to ruby global
$LOG = configureLogging ( options )
case ARGV [ 0 ]
when 'load'
2014-09-17 19:24:53 -04:00
loadRegions ( options , hostname , port )
2011-04-05 00:08:24 -04:00
when 'unload'
2014-09-17 19:24:53 -04:00
unloadRegions ( options , hostname , port )
2011-04-05 00:08:24 -04:00
else
puts optparse
exit 3
end
2013-11-28 11:23:26 -05:00
closeTables ( )