@ -22,7 +22,6 @@
# include "rocksdb/slice.h"
# include "rocksdb/slice.h"
# include "rocksdb/utilities/transaction_db_mutex.h"
# include "rocksdb/utilities/transaction_db_mutex.h"
# include "util/autovector.h"
# include "util/murmurhash.h"
# include "util/murmurhash.h"
# include "util/sync_point.h"
# include "util/sync_point.h"
# include "util/thread_local.h"
# include "util/thread_local.h"
@ -31,15 +30,20 @@
namespace rocksdb {
namespace rocksdb {
struct LockInfo {
struct LockInfo {
TransactionID txn_id ;
bool exclusive ;
autovector < TransactionID > txn_ids ;
// Transaction locks are not valid after this time in us
// Transaction locks are not valid after this time in us
uint64_t expiration_time ;
uint64_t expiration_time ;
LockInfo ( TransactionID id , uint64_t time )
LockInfo ( TransactionID id , uint64_t time , bool ex )
: txn_id ( id ) , expiration_time ( time ) { }
: exclusive ( ex ) , expiration_time ( time ) {
txn_ids . push_back ( id ) ;
}
LockInfo ( const LockInfo & lock_info )
LockInfo ( const LockInfo & lock_info )
: txn_id ( lock_info . txn_id ) , expiration_time ( lock_info . expiration_time ) { }
: exclusive ( lock_info . exclusive ) ,
txn_ids ( lock_info . txn_ids ) ,
expiration_time ( lock_info . expiration_time ) { }
} ;
} ;
struct LockMapStripe {
struct LockMapStripe {
@ -192,7 +196,8 @@ std::shared_ptr<LockMap> TransactionLockMgr::GetLockMap(
// transaction.
// transaction.
// If false, sets *expire_time to the expiration time of the lock according
// If false, sets *expire_time to the expiration time of the lock according
// to Env->GetMicros() or 0 if no expiration.
// to Env->GetMicros() or 0 if no expiration.
bool TransactionLockMgr : : IsLockExpired ( const LockInfo & lock_info , Env * env ,
bool TransactionLockMgr : : IsLockExpired ( TransactionID txn_id ,
const LockInfo & lock_info , Env * env ,
uint64_t * expire_time ) {
uint64_t * expire_time ) {
auto now = env - > NowMicros ( ) ;
auto now = env - > NowMicros ( ) ;
@ -203,12 +208,18 @@ bool TransactionLockMgr::IsLockExpired(const LockInfo& lock_info, Env* env,
// return how many microseconds until lock will be expired
// return how many microseconds until lock will be expired
* expire_time = lock_info . expiration_time ;
* expire_time = lock_info . expiration_time ;
} else {
} else {
bool success =
for ( auto id : lock_info . txn_ids ) {
txn_db_impl_ - > TryStealingExpiredTransactionLocks ( lock_info . txn_id ) ;
if ( txn_id = = id ) {
if ( ! success ) {
continue ;
expired = false ;
}
bool success = txn_db_impl_ - > TryStealingExpiredTransactionLocks ( id ) ;
if ( ! success ) {
expired = false ;
break ;
}
* expire_time = 0 ;
}
}
* expire_time = 0 ;
}
}
return expired ;
return expired ;
@ -216,7 +227,8 @@ bool TransactionLockMgr::IsLockExpired(const LockInfo& lock_info, Env* env,
Status TransactionLockMgr : : TryLock ( TransactionImpl * txn ,
Status TransactionLockMgr : : TryLock ( TransactionImpl * txn ,
uint32_t column_family_id ,
uint32_t column_family_id ,
const std : : string & key , Env * env ) {
const std : : string & key , Env * env ,
bool exclusive ) {
// Lookup lock map for this column family id
// Lookup lock map for this column family id
std : : shared_ptr < LockMap > lock_map_ptr = GetLockMap ( column_family_id ) ;
std : : shared_ptr < LockMap > lock_map_ptr = GetLockMap ( column_family_id ) ;
LockMap * lock_map = lock_map_ptr . get ( ) ;
LockMap * lock_map = lock_map_ptr . get ( ) ;
@ -233,7 +245,7 @@ Status TransactionLockMgr::TryLock(TransactionImpl* txn,
assert ( lock_map - > lock_map_stripes_ . size ( ) > stripe_num ) ;
assert ( lock_map - > lock_map_stripes_ . size ( ) > stripe_num ) ;
LockMapStripe * stripe = lock_map - > lock_map_stripes_ . at ( stripe_num ) ;
LockMapStripe * stripe = lock_map - > lock_map_stripes_ . at ( stripe_num ) ;
LockInfo lock_info ( txn - > GetID ( ) , txn - > GetExpirationTime ( ) ) ;
LockInfo lock_info ( txn - > GetID ( ) , txn - > GetExpirationTime ( ) , exclusive ) ;
int64_t timeout = txn - > GetLockTimeout ( ) ;
int64_t timeout = txn - > GetLockTimeout ( ) ;
return AcquireWithTimeout ( txn , lock_map , stripe , column_family_id , key , env ,
return AcquireWithTimeout ( txn , lock_map , stripe , column_family_id , key , env ,
@ -268,9 +280,9 @@ Status TransactionLockMgr::AcquireWithTimeout(
// Acquire lock if we are able to
// Acquire lock if we are able to
uint64_t expire_time_hint = 0 ;
uint64_t expire_time_hint = 0 ;
TransactionID wait_id = 0 ;
autovector < TransactionID > wait_ids ;
result = AcquireLocked ( lock_map , stripe , key , env , lock_info ,
result = AcquireLocked ( lock_map , stripe , key , env , lock_info ,
& expire_time_hint , & wait_id ) ;
& expire_time_hint , & wait_ids ) ;
if ( ! result . ok ( ) & & timeout ! = 0 ) {
if ( ! result . ok ( ) & & timeout ! = 0 ) {
// If we weren't able to acquire the lock, we will keep retrying as long
// If we weren't able to acquire the lock, we will keep retrying as long
@ -289,19 +301,19 @@ Status TransactionLockMgr::AcquireWithTimeout(
cv_end_time = end_time ;
cv_end_time = end_time ;
}
}
assert ( result . IsBusy ( ) | | wait_id ! = 0 ) ;
assert ( result . IsBusy ( ) | | wait_ids . size ( ) ! = 0 ) ;
// We are dependent on a transaction to finish, so perform deadlock
// We are dependent on a transaction to finish, so perform deadlock
// detection.
// detection.
if ( wait_id ! = 0 ) {
if ( wait_ids . size ( ) ! = 0 ) {
if ( txn - > IsDeadlockDetect ( ) ) {
if ( txn - > IsDeadlockDetect ( ) ) {
if ( IncrementWaiters ( txn , wait_id ) ) {
if ( IncrementWaiters ( txn , wait_ids ) ) {
result = Status : : Busy ( Status : : SubCode : : kDeadlock ) ;
result = Status : : Busy ( Status : : SubCode : : kDeadlock ) ;
stripe - > stripe_mutex - > UnLock ( ) ;
stripe - > stripe_mutex - > UnLock ( ) ;
return result ;
return result ;
}
}
}
}
txn - > SetWaitingTxn ( wait_id , column_family_id , & key ) ;
txn - > SetWaitingTxn ( wait_ids , column_family_id , & key ) ;
}
}
TEST_SYNC_POINT ( " TransactionLockMgr::AcquireWithTimeout:WaitingTxn " ) ;
TEST_SYNC_POINT ( " TransactionLockMgr::AcquireWithTimeout:WaitingTxn " ) ;
@ -316,10 +328,10 @@ Status TransactionLockMgr::AcquireWithTimeout(
}
}
}
}
if ( wait_id ! = 0 ) {
if ( wait_ids . size ( ) ! = 0 ) {
txn - > SetWaitingTxn ( 0 , 0 , nullptr ) ;
txn - > ClearWaitingTxn ( ) ;
if ( txn - > IsDeadlockDetect ( ) ) {
if ( txn - > IsDeadlockDetect ( ) ) {
DecrementWaiters ( txn , wait_id ) ;
DecrementWaiters ( txn , wait_ids ) ;
}
}
}
}
@ -332,7 +344,7 @@ Status TransactionLockMgr::AcquireWithTimeout(
if ( result . ok ( ) | | result . IsTimedOut ( ) ) {
if ( result . ok ( ) | | result . IsTimedOut ( ) ) {
result = AcquireLocked ( lock_map , stripe , key , env , lock_info ,
result = AcquireLocked ( lock_map , stripe , key , env , lock_info ,
& expire_time_hint , & wait_id ) ;
& expire_time_hint , & wait_ids ) ;
}
}
} while ( ! result . ok ( ) & & ! timed_out ) ;
} while ( ! result . ok ( ) & & ! timed_out ) ;
}
}
@ -342,35 +354,40 @@ Status TransactionLockMgr::AcquireWithTimeout(
return result ;
return result ;
}
}
void TransactionLockMgr : : DecrementWaiters ( const TransactionImpl * txn ,
void TransactionLockMgr : : DecrementWaiters (
TransactionID wait_id ) {
const TransactionImpl * txn , const autovector < TransactionID > & wait_ids ) {
std : : lock_guard < std : : mutex > lock ( wait_txn_map_mutex_ ) ;
std : : lock_guard < std : : mutex > lock ( wait_txn_map_mutex_ ) ;
DecrementWaitersImpl ( txn , wait_id ) ;
DecrementWaitersImpl ( txn , wait_ids ) ;
}
}
void TransactionLockMgr : : DecrementWaitersImpl ( const TransactionImpl * txn ,
void TransactionLockMgr : : DecrementWaitersImpl (
TransactionID wait_id ) {
const TransactionImpl * txn , const autovector < TransactionID > & wait_ids ) {
auto id = txn - > GetID ( ) ;
auto id = txn - > GetID ( ) ;
assert ( wait_txn_map_ . Contains ( id ) ) ;
assert ( wait_txn_map_ . Contains ( id ) ) ;
wait_txn_map_ . Delete ( id ) ;
wait_txn_map_ . Delete ( id ) ;
rev_wait_txn_map_ . Get ( wait_id ) - - ;
for ( auto wait_id : wait_ids ) {
if ( rev_wait_txn_map_ . Get ( wait_id ) = = 0 ) {
rev_wait_txn_map_ . Get ( wait_id ) - - ;
rev_wait_txn_map_ . Delete ( wait_id ) ;
if ( rev_wait_txn_map_ . Get ( wait_id ) = = 0 ) {
rev_wait_txn_map_ . Delete ( wait_id ) ;
}
}
}
}
}
bool TransactionLockMgr : : IncrementWaiters ( const TransactionImpl * txn ,
bool TransactionLockMgr : : IncrementWaiters (
TransactionID wait_id ) {
const TransactionImpl * txn , const autovector < TransactionID > & wait_ids ) {
auto id = txn - > GetID ( ) ;
auto id = txn - > GetID ( ) ;
std : : vector < TransactionID > queue ( txn - > GetDeadlockDetectDepth ( ) ) ;
std : : lock_guard < std : : mutex > lock ( wait_txn_map_mutex_ ) ;
std : : lock_guard < std : : mutex > lock ( wait_txn_map_mutex_ ) ;
assert ( ! wait_txn_map_ . Contains ( id ) ) ;
assert ( ! wait_txn_map_ . Contains ( id ) ) ;
wait_txn_map_ . Insert ( id , wait_id ) ;
wait_txn_map_ . Insert ( id , wait_ids ) ;
if ( rev_wait_txn_map_ . Contains ( wait_id ) ) {
for ( auto wait_id : wait_ids ) {
rev_wait_txn_map_ . Get ( wait_id ) + + ;
if ( rev_wait_txn_map_ . Contains ( wait_id ) ) {
} else {
rev_wait_txn_map_ . Get ( wait_id ) + + ;
rev_wait_txn_map_ . Insert ( wait_id , 1 ) ;
} else {
rev_wait_txn_map_ . Insert ( wait_id , 1 ) ;
}
}
}
// No deadlock if nobody is waiting on self.
// No deadlock if nobody is waiting on self.
@ -378,20 +395,36 @@ bool TransactionLockMgr::IncrementWaiters(const TransactionImpl* txn,
return false ;
return false ;
}
}
TransactionID next = wait_id ;
const auto * next_ids = & wait_ids ;
for ( int i = 0 ; i < txn - > GetDeadlockDetectDepth ( ) ; i + + ) {
for ( int tail = 0 , head = 0 ; head < txn - > GetDeadlockDetectDepth ( ) ; head + + ) {
uint i = 0 ;
if ( next_ids ) {
for ( ; i < next_ids - > size ( ) & & tail + i < txn - > GetDeadlockDetectDepth ( ) ;
i + + ) {
queue [ tail + i ] = ( * next_ids ) [ i ] ;
}
tail + = i ;
}
// No more items in the list, meaning no deadlock.
if ( tail = = head ) {
return false ;
}
auto next = queue [ head ] ;
if ( next = = id ) {
if ( next = = id ) {
DecrementWaitersImpl ( txn , wait_id ) ;
DecrementWaitersImpl ( txn , wait_ids ) ;
return true ;
return true ;
} else if ( ! wait_txn_map_ . Contains ( next ) ) {
} else if ( ! wait_txn_map_ . Contains ( next ) ) {
return false ;
next_ids = nullptr ;
continue ;
} else {
} else {
next = wait_txn_map_ . Get ( next ) ;
next_ids = & wait_txn_map_ . Get ( next ) ;
}
}
}
}
// Wait cycle too big, just assume deadlock.
// Wait cycle too big, just assume deadlock.
DecrementWaitersImpl ( txn , wait_id ) ;
DecrementWaitersImpl ( txn , wait_ids ) ;
return true ;
return true ;
}
}
@ -404,24 +437,47 @@ Status TransactionLockMgr::AcquireLocked(LockMap* lock_map,
const std : : string & key , Env * env ,
const std : : string & key , Env * env ,
const LockInfo & txn_lock_info ,
const LockInfo & txn_lock_info ,
uint64_t * expire_time ,
uint64_t * expire_time ,
TransactionID * txn_id ) {
autovector < TransactionID > * txn_ids ) {
assert ( txn_lock_info . txn_ids . size ( ) = = 1 ) ;
Status result ;
Status result ;
// Check if this key is already locked
// Check if this key is already locked
if ( stripe - > keys . find ( key ) ! = stripe - > keys . end ( ) ) {
if ( stripe - > keys . find ( key ) ! = stripe - > keys . end ( ) ) {
// Lock already held
// Lock already held
LockInfo & lock_info = stripe - > keys . at ( key ) ;
LockInfo & lock_info = stripe - > keys . at ( key ) ;
if ( lock_info . txn_id ! = txn_lock_info . txn_id ) {
assert ( lock_info . txn_ids . size ( ) = = 1 | | ! lock_info . exclusive ) ;
// locked by another txn. Check if it's expired
if ( IsLockExpired ( lock_info , env , expire_time ) ) {
if ( lock_info . exclusive | | txn_lock_info . exclusive ) {
// lock is expired, can steal it
if ( lock_info . txn_ids . size ( ) = = 1 & &
lock_info . txn_id = txn_lock_info . txn_id ;
lock_info . txn_ids [ 0 ] = = txn_lock_info . txn_ids [ 0 ] ) {
// The list contains one txn and we're it, so just take it.
lock_info . exclusive = txn_lock_info . exclusive ;
lock_info . expiration_time = txn_lock_info . expiration_time ;
lock_info . expiration_time = txn_lock_info . expiration_time ;
// lock_cnt does not change
} else {
} else {
result = Status : : TimedOut ( Status : : SubCode : : kLockTimeout ) ;
// Check if it's expired. Skips over txn_lock_info.txn_ids[0] in case
* txn_id = lock_info . txn_id ;
// it's there for a shared lock with multiple holders which was not
// caught in the first case.
if ( IsLockExpired ( txn_lock_info . txn_ids [ 0 ] , lock_info , env ,
expire_time ) ) {
// lock is expired, can steal it
lock_info . txn_ids = txn_lock_info . txn_ids ;
lock_info . exclusive = txn_lock_info . exclusive ;
lock_info . expiration_time = txn_lock_info . expiration_time ;
// lock_cnt does not change
} else {
result = Status : : TimedOut ( Status : : SubCode : : kLockTimeout ) ;
* txn_ids = lock_info . txn_ids ;
}
}
}
} else {
// We are requesting shared access to a shared lock, so just grant it.
lock_info . txn_ids . push_back ( txn_lock_info . txn_ids [ 0 ] ) ;
// Using std::max means that expiration time never goes down even when
// a transaction is removed from the list. The correct solution would be
// to track expiry for every transaction, but this would also work for
// now.
lock_info . expiration_time =
std : : max ( lock_info . expiration_time , txn_lock_info . expiration_time ) ;
}
}
} else { // Lock not held.
} else { // Lock not held.
// Check lock limit
// Check lock limit
@ -442,6 +498,42 @@ Status TransactionLockMgr::AcquireLocked(LockMap* lock_map,
return result ;
return result ;
}
}
void TransactionLockMgr : : UnLockKey ( const TransactionImpl * txn ,
const std : : string & key ,
LockMapStripe * stripe , LockMap * lock_map ,
Env * env ) {
TransactionID txn_id = txn - > GetID ( ) ;
auto stripe_iter = stripe - > keys . find ( key ) ;
if ( stripe_iter ! = stripe - > keys . end ( ) ) {
auto & txns = stripe_iter - > second . txn_ids ;
auto txn_it = std : : find ( txns . begin ( ) , txns . end ( ) , txn_id ) ;
// Found the key we locked. unlock it.
if ( txn_it ! = txns . end ( ) ) {
if ( txns . size ( ) = = 1 ) {
stripe - > keys . erase ( stripe_iter ) ;
} else {
auto last_it = txns . end ( ) - 1 ;
if ( txn_it ! = last_it ) {
* txn_it = * last_it ;
}
txns . pop_back ( ) ;
}
if ( max_num_locks_ > 0 ) {
// Maintain lock count if there is a limit on the number of locks.
assert ( lock_map - > lock_cnt . load ( std : : memory_order_relaxed ) > 0 ) ;
lock_map - > lock_cnt - - ;
}
}
} else {
// This key is either not locked or locked by someone else. This should
// only happen if the unlocking transaction has expired.
assert ( txn - > GetExpirationTime ( ) > 0 & &
txn - > GetExpirationTime ( ) < env - > NowMicros ( ) ) ;
}
}
void TransactionLockMgr : : UnLock ( TransactionImpl * txn , uint32_t column_family_id ,
void TransactionLockMgr : : UnLock ( TransactionImpl * txn , uint32_t column_family_id ,
const std : : string & key , Env * env ) {
const std : : string & key , Env * env ) {
std : : shared_ptr < LockMap > lock_map_ptr = GetLockMap ( column_family_id ) ;
std : : shared_ptr < LockMap > lock_map_ptr = GetLockMap ( column_family_id ) ;
@ -456,26 +548,8 @@ void TransactionLockMgr::UnLock(TransactionImpl* txn, uint32_t column_family_id,
assert ( lock_map - > lock_map_stripes_ . size ( ) > stripe_num ) ;
assert ( lock_map - > lock_map_stripes_ . size ( ) > stripe_num ) ;
LockMapStripe * stripe = lock_map - > lock_map_stripes_ . at ( stripe_num ) ;
LockMapStripe * stripe = lock_map - > lock_map_stripes_ . at ( stripe_num ) ;
TransactionID txn_id = txn - > GetID ( ) ;
stripe - > stripe_mutex - > Lock ( ) ;
stripe - > stripe_mutex - > Lock ( ) ;
UnLockKey ( txn , key , stripe , lock_map , env ) ;
const auto & iter = stripe - > keys . find ( key ) ;
if ( iter ! = stripe - > keys . end ( ) & & iter - > second . txn_id = = txn_id ) {
// Found the key we locked. unlock it.
stripe - > keys . erase ( iter ) ;
if ( max_num_locks_ > 0 ) {
// Maintain lock count if there is a limit on the number of locks.
assert ( lock_map - > lock_cnt . load ( std : : memory_order_relaxed ) > 0 ) ;
lock_map - > lock_cnt - - ;
}
} else {
// This key is either not locked or locked by someone else. This should
// only happen if the unlocking transaction has expired.
assert ( txn - > GetExpirationTime ( ) > 0 & &
txn - > GetExpirationTime ( ) < env - > NowMicros ( ) ) ;
}
stripe - > stripe_mutex - > UnLock ( ) ;
stripe - > stripe_mutex - > UnLock ( ) ;
// Signal waiting threads to retry locking
// Signal waiting threads to retry locking
@ -484,8 +558,6 @@ void TransactionLockMgr::UnLock(TransactionImpl* txn, uint32_t column_family_id,
void TransactionLockMgr : : UnLock ( const TransactionImpl * txn ,
void TransactionLockMgr : : UnLock ( const TransactionImpl * txn ,
const TransactionKeyMap * key_map , Env * env ) {
const TransactionKeyMap * key_map , Env * env ) {
TransactionID txn_id = txn - > GetID ( ) ;
for ( auto & key_map_iter : * key_map ) {
for ( auto & key_map_iter : * key_map ) {
uint32_t column_family_id = key_map_iter . first ;
uint32_t column_family_id = key_map_iter . first ;
auto & keys = key_map_iter . second ;
auto & keys = key_map_iter . second ;
@ -520,22 +592,7 @@ void TransactionLockMgr::UnLock(const TransactionImpl* txn,
stripe - > stripe_mutex - > Lock ( ) ;
stripe - > stripe_mutex - > Lock ( ) ;
for ( const std : : string * key : stripe_keys ) {
for ( const std : : string * key : stripe_keys ) {
const auto & iter = stripe - > keys . find ( * key ) ;
UnLockKey ( txn , * key , stripe , lock_map , env ) ;
if ( iter ! = stripe - > keys . end ( ) & & iter - > second . txn_id = = txn_id ) {
// Found the key we locked. unlock it.
stripe - > keys . erase ( iter ) ;
if ( max_num_locks_ > 0 ) {
// Maintain lock count if there is a limit on the number of locks.
assert ( lock_map - > lock_cnt . load ( std : : memory_order_relaxed ) > 0 ) ;
lock_map - > lock_cnt - - ;
}
} else {
// This key is either not locked or locked by someone else. This
// should only
// happen if the unlocking transaction has expired.
assert ( txn - > GetExpirationTime ( ) > 0 & &
txn - > GetExpirationTime ( ) < env - > NowMicros ( ) ) ;
}
}
}
stripe - > stripe_mutex - > UnLock ( ) ;
stripe - > stripe_mutex - > UnLock ( ) ;
@ -565,7 +622,13 @@ TransactionLockMgr::LockStatusData TransactionLockMgr::GetLockStatusData() {
for ( const auto & j : stripes ) {
for ( const auto & j : stripes ) {
j - > stripe_mutex - > Lock ( ) ;
j - > stripe_mutex - > Lock ( ) ;
for ( const auto & it : j - > keys ) {
for ( const auto & it : j - > keys ) {
data . insert ( { i , { it . first , it . second . txn_id } } ) ;
struct KeyLockInfo info ;
info . exclusive = it . second . exclusive ;
info . key = it . first ;
for ( const auto & id : it . second . txn_ids ) {
info . ids . push_back ( id ) ;
}
data . insert ( { i , info } ) ;
}
}
}
}
}
}