Skip to content

Commit e2c2b5a

Browse files
committed
fix: avoid MemTable sort
1 parent a8e7b34 commit e2c2b5a

1 file changed

Lines changed: 84 additions & 15 deletions

File tree

src/db.rs

Lines changed: 84 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,76 @@ struct MergedIterator<'a> {
1919
projections: Option<Vec<Expression<'a>>>,
2020
}
2121

22+
struct HybridIterator<'a> {
23+
mem_iter: std::collections::hash_map::Iter<'a, String, Value>,
24+
disk_iter: MergedIterator<'a>,
25+
memtable: &'a HashMap<String, Value>,
26+
phase: ScanPhase,
27+
predicate: Option<Expression<'a>>,
28+
projections: Option<Vec<Expression<'a>>>,
29+
}
30+
31+
enum ScanPhase {
32+
MemTable,
33+
Disk,
34+
}
35+
36+
impl<'a> Iterator for HybridIterator<'a> {
37+
type Item = ExecutionResult;
38+
39+
fn next(&mut self) -> Option<Self::Item> {
40+
loop {
41+
match self.phase {
42+
ScanPhase::MemTable => {
43+
if let Some((id, val)) = self.mem_iter.next() {
44+
use jsonb_schema::Value as JsonbValue;
45+
if matches!(val, JsonbValue::Null) {
46+
continue; // Tombstone
47+
}
48+
49+
if let Some(pred) = &self.predicate {
50+
if evaluate_expression(pred, val) != Value::Bool(true) {
51+
continue;
52+
}
53+
}
54+
55+
if let Some(projs) = &self.projections {
56+
let mut new_doc = BTreeMap::new();
57+
for expr in projs {
58+
let v = evaluate_expression(expr, val);
59+
let key = match expr {
60+
Expression::FieldReference(_, raw) => raw.to_string(),
61+
Expression::JsonPath(_, raw) => raw.to_string(),
62+
_ => "col".to_string(),
63+
};
64+
new_doc.insert(key, v);
65+
}
66+
return Some(ExecutionResult::Value(
67+
id.clone(),
68+
Value::Object(new_doc),
69+
));
70+
}
71+
72+
return Some(ExecutionResult::Value(id.clone(), val.clone()));
73+
} else {
74+
self.phase = ScanPhase::Disk;
75+
}
76+
}
77+
ScanPhase::Disk => {
78+
if let Some(res) = self.disk_iter.next() {
79+
if self.memtable.contains_key(res.id()) {
80+
continue;
81+
}
82+
return Some(res);
83+
} else {
84+
return None;
85+
}
86+
}
87+
}
88+
}
89+
}
90+
}
91+
2292
use crate::expression::evaluate_expression_lazy;
2393

2494
impl<'a> Iterator for MergedIterator<'a> {
@@ -309,33 +379,32 @@ impl Collection {
309379
predicate: Option<Expression<'a>>,
310380
projections: Option<Vec<Expression<'a>>>,
311381
) -> impl Iterator<Item = ExecutionResult> + 'a {
312-
let mut sources: Vec<SourceIterator> = Vec::new();
382+
let mut disk_sources: Vec<SourceIterator> = Vec::new();
313383

314-
// 1. MemTable Iterator (Priority 0 - Highest)
315-
// Must sort because MemTable now uses HashMap
316-
let mut mem_docs: Vec<_> = self.memtable.documents.iter().collect();
317-
mem_docs.sort_by_key(|(k, _)| *k);
318-
let mem_iter = mem_docs
319-
.into_iter()
320-
.map(|(k, v)| ExecutionResult::Value(k.clone(), v.clone()));
321-
322-
sources.push((Box::new(mem_iter) as Box<dyn Iterator<Item = ExecutionResult>>).peekable());
323-
324-
// 2. JSTable Iterators (Newer to Older)
384+
// JSTable Iterators (Newer to Older)
325385
for i in (0..self.jstable_count).rev() {
326386
let path = self.dir.join(format!("jstable-{}", i));
327387
if let Ok(iter) = jstable::JSTableLazyIterator::new(path.to_str().unwrap()) {
328388
let iter = iter.map(|r| {
329389
let doc = r.unwrap();
330390
ExecutionResult::Lazy(doc)
331391
});
332-
sources
392+
disk_sources
333393
.push((Box::new(iter) as Box<dyn Iterator<Item = ExecutionResult>>).peekable());
334394
}
335395
}
336396

337-
MergedIterator {
338-
sources,
397+
let disk_iter = MergedIterator {
398+
sources: disk_sources,
399+
predicate: predicate.clone(),
400+
projections: projections.clone(),
401+
};
402+
403+
HybridIterator {
404+
mem_iter: self.memtable.documents.iter(),
405+
disk_iter,
406+
memtable: &self.memtable.documents,
407+
phase: ScanPhase::MemTable,
339408
predicate,
340409
projections,
341410
}

0 commit comments

Comments
 (0)