@ -2273,26 +2273,20 @@ Status DBImpl::Write(const WriteOptions& options, WriteBatch* my_batch) {
StopWatch sw ( env_ , options_ . statistics , DB_WRITE ) ;
StopWatch sw ( env_ , options_ . statistics , DB_WRITE ) ;
MutexLock l ( & mutex_ ) ;
MutexLock l ( & mutex_ ) ;
writers_ . push_back ( & w ) ;
// If WAL is disabled, we avoid any queueing.
while ( ! w . done & & & w ! = writers_ . front ( ) ) {
if ( ! options . disableWAL ) {
w . cv . Wait ( ) ;
writers_ . push_back ( & w ) ;
}
while ( ! w . done & & & w ! = writers_ . front ( ) ) {
if ( w . done ) {
w . cv . Wait ( ) ;
return w . status ;
}
if ( w . done ) {
return w . status ;
}
}
}
// May temporarily unlock and wait.
// May temporarily unlock and wait.
Status status = MakeRoomForWrite ( my_batch = = nullptr ) ;
Status status = MakeRoomForWrite ( my_batch = = nullptr ) ;
uint64_t last_sequence = versions_ - > LastSequence ( ) ;
uint64_t last_sequence = versions_ - > LastSequence ( ) ;
Writer * last_writer = & w ;
Writer * last_writer = & w ;
if ( status . ok ( ) & & my_batch ! = nullptr ) { // nullptr batch is for compactions
if ( status . ok ( ) & & my_batch ! = nullptr ) { // nullptr batch is for compactions
WriteBatch * updates = options . disableWAL ? my_batch :
WriteBatch * updates = BuildBatchGroup ( & last_writer ) ;
BuildBatchGroup ( & last_writer ) ;
const SequenceNumber current_sequence = last_sequence + 1 ;
const SequenceNumber current_sequence = last_sequence + 1 ;
WriteBatchInternal : : SetSequence ( updates , current_sequence ) ;
WriteBatchInternal : : SetSequence ( updates , current_sequence ) ;
int my_batch_count = WriteBatchInternal : : Count ( updates ) ;
int my_batch_count = WriteBatchInternal : : Count ( updates ) ;
@ -2307,12 +2301,12 @@ Status DBImpl::Write(const WriteOptions& options, WriteBatch* my_batch) {
// and protects against concurrent loggers and concurrent writes
// and protects against concurrent loggers and concurrent writes
// into mem_.
// into mem_.
{
{
mutex_ . Unlock ( ) ;
if ( options . disableWAL ) {
if ( options . disableWAL ) {
// If WAL is disabled, then we do not drop the mutex. We keep the
// mutex to protect concurrent insertions into the memtable.
flush_on_destroy_ = true ;
flush_on_destroy_ = true ;
} else {
}
mutex_ . Unlock ( ) ;
if ( ! options . disableWAL ) {
status = log_ - > AddRecord ( WriteBatchInternal : : Contents ( updates ) ) ;
status = log_ - > AddRecord ( WriteBatchInternal : : Contents ( updates ) ) ;
if ( status . ok ( ) & & options . sync ) {
if ( status . ok ( ) & & options . sync ) {
if ( options_ . use_fsync ) {
if ( options_ . use_fsync ) {
@ -2337,29 +2331,25 @@ Status DBImpl::Write(const WriteOptions& options, WriteBatch* my_batch) {
versions_ - > SetLastSequence ( last_sequence ) ;
versions_ - > SetLastSequence ( last_sequence ) ;
last_flushed_sequence_ = current_sequence ;
last_flushed_sequence_ = current_sequence ;
}
}
if ( ! options . disableWAL ) {
mutex_ . Lock ( ) ;
mutex_ . Lock ( ) ;
}
}
}
if ( updates = = & tmp_batch_ ) tmp_batch_ . Clear ( ) ;
if ( updates = = & tmp_batch_ ) tmp_batch_ . Clear ( ) ;
}
}
if ( ! options . disableWAL ) {
while ( true ) {
while ( true ) {
Writer * ready = writers_ . front ( ) ;
Writer * ready = writers_ . front ( ) ;
writers_ . pop_front ( ) ;
writers_ . pop_front ( ) ;
if ( ready ! = & w ) {
if ( ready ! = & w ) {
ready - > status = status ;
ready - > status = status ;
ready - > done = true ;
ready - > done = true ;
ready - > cv . Signal ( ) ;
ready - > cv . Signal ( ) ;
}
if ( ready = = last_writer ) break ;
}
}
if ( ready = = last_writer ) break ;
}
// Notify new head of write queue
// Notify new head of write queue
if ( ! writers_ . empty ( ) ) {
if ( ! writers_ . empty ( ) ) {
writers_ . front ( ) - > cv . Signal ( ) ;
writers_ . front ( ) - > cv . Signal ( ) ;
}
}
}
return status ;
return status ;
}
}
@ -2423,6 +2413,7 @@ WriteBatch* DBImpl::BuildBatchGroup(Writer** last_writer) {
// REQUIRES: this thread is currently at the front of the writer queue
// REQUIRES: this thread is currently at the front of the writer queue
Status DBImpl : : MakeRoomForWrite ( bool force ) {
Status DBImpl : : MakeRoomForWrite ( bool force ) {
mutex_ . AssertHeld ( ) ;
mutex_ . AssertHeld ( ) ;
assert ( ! writers_ . empty ( ) ) ;
bool allow_delay = ! force ;
bool allow_delay = ! force ;
bool allow_rate_limit_delay = ! force ;
bool allow_rate_limit_delay = ! force ;
uint64_t rate_limit_delay_millis = 0 ;
uint64_t rate_limit_delay_millis = 0 ;