Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion rust-toolchain
Original file line number Diff line number Diff line change
@@ -1 +1 @@
nightly-2021-07-28
nightly-2021-12-07
4 changes: 2 additions & 2 deletions src/consistency.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,8 +28,8 @@ impl ReplayMachine for ConsistencyChecker {
for item in item_batch.iter() {
if let LogItemContent::EntryIndexes(ents) = &item.content {
if !ents.0.is_empty() {
let incoming_first_index = ents.0.first().unwrap().index;
let incoming_last_index = ents.0.last().unwrap().index;
let incoming_first_index = ents.0.front().unwrap().index;
let incoming_last_index = ents.0.back().unwrap().index;
let index_range = self
.raft_groups
.entry(item.raft_group_id)
Expand Down
1 change: 0 additions & 1 deletion src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@
// See the License for the specific language governing permissions and
// limitations under the License.

#![feature(shrink_to)]
#![feature(map_first_last)]
#![feature(generic_associated_types)]
#![feature(test)]
Expand Down
29 changes: 15 additions & 14 deletions src/log_batch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ use byteorder::{BigEndian, ReadBytesExt, WriteBytesExt};
use log::error;
use protobuf::parse_from_bytes;
use protobuf::Message;
use std::collections::VecDeque;

use crate::codec::{self, NumberEncoder};
use crate::memtable::EntryIndex;
Expand Down Expand Up @@ -63,12 +64,12 @@ type SliceReader<'a> = &'a [u8];
// Format:
// { count | first index | [ tail offsets ] }
#[derive(Clone, Debug, PartialEq)]
pub struct EntryIndexes(pub Vec<EntryIndex>);
pub struct EntryIndexes(pub VecDeque<EntryIndex>);

impl EntryIndexes {
pub fn decode(buf: &mut SliceReader, entries_size: &mut usize) -> Result<Self> {
let mut count = codec::decode_var_u64(buf)?;
let mut entry_indexes = Vec::with_capacity(count as usize);
let mut entry_indexes = VecDeque::with_capacity(count as usize);
let mut index = 0;
if count > 0 {
index = codec::decode_var_u64(buf)?;
Expand All @@ -83,7 +84,7 @@ impl EntryIndexes {
..Default::default()
};
*entries_size += entry_len;
entry_indexes.push(entry_index);
entry_indexes.push_back(entry_index);
index += 1;
count -= 1;
}
Expand Down Expand Up @@ -247,7 +248,7 @@ pub enum LogItemContent {
}

impl LogItem {
pub fn new_entry_indexes(raft_group_id: u64, entry_indexes: Vec<EntryIndex>) -> LogItem {
pub fn new_entry_indexes(raft_group_id: u64, entry_indexes: VecDeque<EntryIndex>) -> LogItem {
LogItem {
raft_group_id,
content: LogItemContent::EntryIndexes(EntryIndexes(entry_indexes)),
Expand Down Expand Up @@ -423,7 +424,7 @@ impl LogItemBatch {
}
}

pub fn add_entry_indexes(&mut self, region_id: u64, mut entry_indexes: Vec<EntryIndex>) {
pub fn add_entry_indexes(&mut self, region_id: u64, mut entry_indexes: VecDeque<EntryIndex>) {
for ei in entry_indexes.iter_mut() {
ei.entry_offset = self.entries_size as u64;
self.entries_size += ei.entry_len;
Expand Down Expand Up @@ -474,7 +475,7 @@ impl LogItemBatch {
compression_type: CompressionType,
) -> Result<LogItemBatch> {
verify_checksum(buf)?;
*buf = &mut &buf[..buf.len() - LOG_BATCH_CHECKSUM_LEN];
*buf = &buf[..buf.len() - LOG_BATCH_CHECKSUM_LEN];
let count = codec::decode_var_u64(buf)?;
let mut items = LogItemBatch::with_capacity(count as usize);
let mut entries_size = 0;
Expand Down Expand Up @@ -574,12 +575,12 @@ impl LogBatch {
) -> Result<()> {
debug_assert!(self.buf_state == BufState::Open);

let mut entry_indexes = Vec::with_capacity(entries.len());
let mut entry_indexes = VecDeque::with_capacity(entries.len());
self.buf_state = BufState::Incomplete;
for e in entries {
let buf_offset = self.buf.len();
e.write_to_vec(&mut self.buf)?;
entry_indexes.push(EntryIndex {
entry_indexes.push_back(EntryIndex {
index: M::index(e),
entry_len: self.buf.len() - buf_offset,
..Default::default()
Expand All @@ -593,7 +594,7 @@ impl LogBatch {
pub(crate) fn add_raw_entries(
&mut self,
region_id: u64,
mut entry_indexes: Vec<EntryIndex>,
mut entry_indexes: VecDeque<EntryIndex>,
entries: Vec<Vec<u8>>,
) -> Result<()> {
debug_assert!(entry_indexes.len() == entries.len());
Expand Down Expand Up @@ -836,7 +837,7 @@ mod tests {

fn decode_entries_from_bytes<M: MessageExt>(
buf: &[u8],
entry_indexes: &[EntryIndex],
entry_indexes: &VecDeque<EntryIndex>,
_encoded: bool,
) -> Vec<M::Entry> {
let mut entries = Vec::with_capacity(entry_indexes.len());
Expand All @@ -848,7 +849,7 @@ mod tests {

#[test]
fn test_entry_indexes_enc_dec() {
fn encode_and_decode(entry_indexes: &mut Vec<EntryIndex>) -> EntryIndexes {
fn encode_and_decode(entry_indexes: &mut VecDeque<EntryIndex>) -> EntryIndexes {
let mut entries_size = 0;
for idx in entry_indexes.iter_mut() {
idx.entry_offset = entries_size as u64;
Expand All @@ -868,7 +869,7 @@ mod tests {
decoded_indexes
}

let entry_indexes = vec![Vec::new(), generate_entry_indexes_opt(7, 17, None)];
let entry_indexes = vec![VecDeque::new(), generate_entry_indexes_opt(7, 17, None)];
for mut idxs in entry_indexes.into_iter() {
let decoded = encode_and_decode(&mut idxs);
assert_eq!(idxs, decoded.0);
Expand Down Expand Up @@ -1062,8 +1063,8 @@ mod tests {
if let LogItemContent::EntryIndexes(entry_indexes) = &item.content {
if !entry_indexes.0.is_empty() {
let (begin, end) = (
entry_indexes.0.first().unwrap().index,
entry_indexes.0.last().unwrap().index + 1,
entry_indexes.0.front().unwrap().index,
entry_indexes.0.back().unwrap().index + 1,
);
let origin_entries = generate_entries(begin, end, Some(entry_data));
let decoded_entries =
Expand Down
32 changes: 23 additions & 9 deletions src/memtable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -186,29 +186,28 @@ impl MemTable {
}
}

pub fn append(&mut self, entry_indexes: Vec<EntryIndex>) {
pub fn append(&mut self, mut entry_indexes: VecDeque<EntryIndex>) {
let len = entry_indexes.len();
if len > 0 {
self.prepare_append(entry_indexes[0].index, false, false);
self.global_stats.add(LogQueue::Append, len);
// TODO: Optimize this.
self.entry_indexes.extend(entry_indexes);
self.entry_indexes.append(&mut entry_indexes);
}
}

// This will only be called during recovery.
pub fn append_rewrite(&mut self, entry_indexes: Vec<EntryIndex>) {
pub fn append_rewrite(&mut self, mut entry_indexes: VecDeque<EntryIndex>) {
let len = entry_indexes.len();
if len > 0 {
debug_assert_eq!(self.rewrite_count, self.entry_indexes.len());
self.prepare_append(entry_indexes[0].index, true, true);
self.global_stats.add(LogQueue::Rewrite, len);
self.entry_indexes.extend(entry_indexes);
self.entry_indexes.append(&mut entry_indexes);
self.rewrite_count = self.entry_indexes.len();
}
}

pub fn rewrite(&mut self, rewrite_indexes: Vec<EntryIndex>, gate: Option<FileSeq>) {
pub fn rewrite(&mut self, rewrite_indexes: VecDeque<EntryIndex>, gate: Option<FileSeq>) {
if rewrite_indexes.is_empty() {
return;
}
Expand Down Expand Up @@ -240,8 +239,8 @@ impl MemTable {
);
let rewrite_pos = (rewrite_first - rewrite_indexes[0].index) as usize;

for (i, rindex) in rewrite_indexes[rewrite_pos..rewrite_pos + rewrite_len]
.iter()
for (i, rindex) in rewrite_indexes
.range(rewrite_pos..rewrite_pos + rewrite_len)
.enumerate()
{
let index = &mut self.entry_indexes[i + pos];
Expand Down Expand Up @@ -844,7 +843,7 @@ mod tests {
memtable.consistency_check();

// Empty.
memtable.append(Vec::new());
memtable.append(VecDeque::new());

// Hole.
assert!(
Expand Down Expand Up @@ -1739,4 +1738,19 @@ mod tests {
memtable.put(key2.clone(), value.clone(), FileId::dummy(LogQueue::Append));
});
}

#[bench]
fn bench_memtable_append(b: &mut test::Bencher) {
let mut memtable = MemTable::new(0, Arc::new(GlobalStats::default()));
b.iter(move || {
let last_index = memtable.last_index().unwrap_or(0);
let entry_indexes = generate_entry_indexes(
last_index + 1,
last_index + 11,
FileId::dummy(LogQueue::Append),
);
memtable.append(entry_indexes);
memtable.compact_to(last_index + 1);
});
}
}
8 changes: 6 additions & 2 deletions src/purge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -291,7 +291,11 @@ where
std::mem::swap(&mut take_entries, &mut entries);
let mut take_entry_indexes = entry_indexes.split_off(cursor + 1);
std::mem::swap(&mut take_entry_indexes, &mut entry_indexes);
log_batch.add_raw_entries(region_id, take_entry_indexes, take_entries)?;
log_batch.add_raw_entries(
region_id,
take_entry_indexes.into(),
take_entries,
)?;
self.rewrite_impl(&mut log_batch, rewrite, false)?;
total_size = 0;
cursor = 0;
Expand All @@ -300,7 +304,7 @@ where
}
}
if !entries.is_empty() {
log_batch.add_raw_entries(region_id, entry_indexes, entries)?;
log_batch.add_raw_entries(region_id, entry_indexes.into(), entries)?;
}
for (k, v) in kvs {
log_batch.put(region_id, k, v);
Expand Down
11 changes: 8 additions & 3 deletions src/test_util.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
// Copyright (c) 2017-present, PingCAP, Inc. Licensed under Apache-2.0.

use std::collections::VecDeque;
use std::panic::{self, AssertUnwindSafe};

use raft::eraftpb::Entry;
Expand All @@ -22,15 +23,19 @@ pub fn generate_entries(begin_index: u64, end_index: u64, data: Option<&[u8]>) -
v
}

pub fn generate_entry_indexes(begin_idx: u64, end_idx: u64, file_id: FileId) -> Vec<EntryIndex> {
pub fn generate_entry_indexes(
begin_idx: u64,
end_idx: u64,
file_id: FileId,
) -> VecDeque<EntryIndex> {
generate_entry_indexes_opt(begin_idx, end_idx, Some(file_id))
}

pub fn generate_entry_indexes_opt(
begin_idx: u64,
end_idx: u64,
file_id: Option<FileId>,
) -> Vec<EntryIndex> {
) -> VecDeque<EntryIndex> {
assert!(end_idx >= begin_idx);
let mut ents_idx = vec![];
for idx in begin_idx..end_idx {
Expand All @@ -43,7 +48,7 @@ pub fn generate_entry_indexes_opt(

ents_idx.push(ent_idx);
}
ents_idx
ents_idx.into()
}

/// Catch panic while suppressing default panic hook.
Expand Down
Loading