Skip to content

Commit 4f08edd

Browse files
committed
refactor: split JSTable into summary and data files
Separated JSTable storage into .summary (header + filter) and .data (records) files. Updated JSTable write, read, and iteration logic. Updated DB to handle new file extensions during load, flush, and compaction. Updated specs.
1 parent 282fdba commit 4f08edd

3 files changed

Lines changed: 89 additions & 52 deletions

File tree

specs/storage.md

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,18 +3,24 @@ A JSTable is like an SSTable, but stores semi-structured JSON data with an assoc
33

44
# Disk format
55

6-
Each JSTable is stored in a single binary file using the [JSONB](https://github.com/databendlabs/jsonb) format.
7-
The file consists of a sequence of entries. Each entry is encoded as:
6+
Each JSTable is stored as two binary files using the [JSONB](https://github.com/databendlabs/jsonb) format: a summary file (`.summary`) and a data file (`.data`).
7+
8+
The files consist of a sequence of entries. Each entry is encoded as:
89
1. **Length**: A 4-byte unsigned integer (little-endian) indicating the size of the following JSONB blob.
910
2. **Data**: A binary blob encoded using the JSONB format.
1011

1112
## Structure
1213

14+
### Summary File
15+
1316
1. **Header Entry**: The first entry in the file. It is a JSONB-encoded object containing:
1417
* `timestamp`: The time the table was created (Unix timestamp in milliseconds).
1518
* `schema`: The JSON Schema for the documents.
1619
2. **Filter Entry**: The second entry in the file. It is a [Binary Fuse8](https://github.com/ayazhafiz/xorf) filter of the record IDs in the table, serialized as a JSON byte vector.
17-
3. **Record Entries**: All subsequent entries. Each is a JSONB-encoded array `[id, document]`:
20+
21+
### Data File
22+
23+
1. **Record Entries**: All entries. Each is a JSONB-encoded array `[id, document]`:
1824
* `id`: String.
1925
* `document`: The document object (or `null` for tombstone).
2026

src/db.rs

Lines changed: 14 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -109,16 +109,15 @@ impl Collection {
109109
// Count existing JSTables and load filters
110110
let mut jstable_count = 0;
111111
let mut filters = Vec::new();
112-
while dir.join(format!("jstable-{}", jstable_count)).exists() {
112+
// Check for .summary file to confirm JSTable existence
113+
while dir
114+
.join(format!("jstable-{}.summary", jstable_count))
115+
.exists()
116+
{
113117
let path = dir.join(format!("jstable-{}", jstable_count));
114118
if let Ok(filter) = jstable::read_filter(path.to_str().unwrap()) {
115119
filters.push(filter);
116120
} else {
117-
// Should not happen if file exists and is valid, but handle gracefully?
118-
// For now, if we can't read the filter, we might just panic or log error.
119-
// Given this is a simple DB, let's assume valid state or fail.
120-
// However, we need to push *something* or fail the whole load.
121-
// Let's panic to signal corruption.
122121
panic!("Failed to read filter for jstable-{}", jstable_count);
123122
}
124123
jstable_count += 1;
@@ -200,8 +199,11 @@ impl Collection {
200199
let merged_table = jstable::merge_jstables(&tables);
201200

202201
for i in 0..self.jstable_count {
203-
let path = self.dir.join(format!("jstable-{}", i));
204-
fs::remove_file(path).unwrap();
202+
let base_path = self.dir.join(format!("jstable-{}", i));
203+
let summary_path = format!("{}.summary", base_path.to_str().unwrap());
204+
let data_path = format!("{}.data", base_path.to_str().unwrap());
205+
fs::remove_file(summary_path).unwrap();
206+
fs::remove_file(data_path).unwrap();
205207
}
206208

207209
let new_path = self.dir.join("jstable-0");
@@ -315,10 +317,11 @@ impl DB {
315317
let dir_path = entry.path();
316318

317319
// Try to find collection name from JSTable-0
318-
let jstable_path = dir_path.join("jstable-0");
319-
let col_name = if jstable_path.exists() {
320+
let jstable_base_path = dir_path.join("jstable-0");
321+
let jstable_summary_path = dir_path.join("jstable-0.summary");
322+
let col_name = if jstable_summary_path.exists() {
320323
if let Ok(iter) =
321-
jstable::JSTableIterator::new(jstable_path.to_str().unwrap())
324+
jstable::JSTableIterator::new(jstable_base_path.to_str().unwrap())
322325
{
323326
Some(iter.collection)
324327
} else {

src/jstable.rs

Lines changed: 66 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -36,9 +36,13 @@ impl JSTable {
3636
}
3737

3838
pub fn write(&self, path: &str) -> io::Result<()> {
39-
let mut file = File::create(path)?;
39+
let summary_path = format!("{}.summary", path);
40+
let data_path = format!("{}.data", path);
4041

41-
// Write Header
42+
let mut summary_file = File::create(summary_path)?;
43+
let mut data_file = File::create(data_path)?;
44+
45+
// Write Header to summary
4246
let header = JSTableHeader {
4347
timestamp: self.timestamp,
4448
collection: self.collection.clone(),
@@ -49,10 +53,10 @@ impl JSTable {
4953
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
5054
let header_bytes = header_blob.to_vec();
5155
let header_len = header_bytes.len() as u32;
52-
file.write_all(&header_len.to_le_bytes())?;
53-
file.write_all(&header_bytes)?;
56+
summary_file.write_all(&header_len.to_le_bytes())?;
57+
summary_file.write_all(&header_bytes)?;
5458

55-
// Write Filter
59+
// Write Filter to summary
5660
let keys: Vec<u64> = self
5761
.documents
5862
.keys()
@@ -66,22 +70,22 @@ impl JSTable {
6670
let filter = BinaryFuse8::try_from(&keys).map_err(|_| {
6771
io::Error::new(io::ErrorKind::InvalidData, "Failed to create XOR filter")
6872
})?;
69-
// Use serde_json for filter serialization as fallback
73+
// Use serde_json for filter serialization
7074
let filter_bytes = serde_json::to_vec(&filter)
7175
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
7276
let filter_len = filter_bytes.len() as u32;
73-
file.write_all(&filter_len.to_le_bytes())?;
74-
file.write_all(&filter_bytes)?;
77+
summary_file.write_all(&filter_len.to_le_bytes())?;
78+
summary_file.write_all(&filter_bytes)?;
7579

76-
// Write Documents
80+
// Write Documents to data
7781
for (id, doc) in &self.documents {
7882
let record: (String, &Value) = (id.clone(), doc);
7983
let record_blob = jsonb::to_owned_jsonb(&record)
8084
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
8185
let record_bytes = record_blob.to_vec();
8286
let record_len = record_bytes.len() as u32;
83-
file.write_all(&record_len.to_le_bytes())?;
84-
file.write_all(&record_bytes)?;
87+
data_file.write_all(&record_len.to_le_bytes())?;
88+
data_file.write_all(&record_bytes)?;
8589
}
8690
Ok(())
8791
}
@@ -96,17 +100,20 @@ pub struct JSTableIterator {
96100

97101
impl JSTableIterator {
98102
pub fn new(path: &str) -> io::Result<Self> {
99-
let file = File::open(path)?;
100-
let mut reader = BufReader::new(file);
103+
let summary_path = format!("{}.summary", path);
104+
let data_path = format!("{}.data", path);
105+
106+
let summary_file = File::open(summary_path)?;
107+
let mut summary_reader = BufReader::new(summary_file);
101108

102-
// Read Header Length
109+
// Read Header Length from summary
103110
let mut len_buf = [0u8; 4];
104-
reader.read_exact(&mut len_buf)?;
111+
summary_reader.read_exact(&mut len_buf)?;
105112
let header_len = u32::from_le_bytes(len_buf) as usize;
106113

107-
// Read Header Blob
114+
// Read Header Blob from summary
108115
let mut header_blob = vec![0u8; header_len];
109-
reader.read_exact(&mut header_blob)?;
116+
summary_reader.read_exact(&mut header_blob)?;
110117

111118
let header_val = jsonb::from_slice(&header_blob)
112119
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
@@ -115,18 +122,13 @@ impl JSTableIterator {
115122
let header: JSTableHeader = serde_json::from_str(&header_str)
116123
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
117124

118-
// Read/Skip Filter
119-
let mut len_buf = [0u8; 4];
120-
reader.read_exact(&mut len_buf)?;
121-
let filter_len = u32::from_le_bytes(len_buf) as usize;
122-
// Skip filter bytes
123-
io::copy(
124-
&mut reader.by_ref().take(filter_len as u64),
125-
&mut io::sink(),
126-
)?;
125+
// We don't need to read the filter here, so we are done with summary file
126+
127+
let data_file = File::open(data_path)?;
128+
let data_reader = BufReader::new(data_file);
127129

128130
Ok(Self {
129-
reader,
131+
reader: data_reader,
130132
timestamp: header.timestamp,
131133
collection: header.collection,
132134
schema: header.schema,
@@ -190,7 +192,8 @@ pub fn read_jstable(path: &str) -> io::Result<JSTable> {
190192
}
191193

192194
pub fn read_filter(path: &str) -> io::Result<BinaryFuse8> {
193-
let file = File::open(path)?;
195+
let summary_path = format!("{}.summary", path);
196+
let file = File::open(summary_path)?;
194197
let mut reader = BufReader::new(file);
195198

196199
// Read Header Length
@@ -254,7 +257,7 @@ mod tests {
254257
use super::*;
255258
use crate::schema::SchemaType;
256259
use serde_json::json;
257-
use tempfile::NamedTempFile;
260+
use tempfile::tempdir;
258261
use xorf::Filter;
259262

260263
#[test]
@@ -277,10 +280,11 @@ mod tests {
277280
documents.clone(),
278281
);
279282

280-
let file = NamedTempFile::new().unwrap();
281-
jstable.write(file.path().to_str().unwrap()).unwrap();
283+
let dir = tempdir()?;
284+
let file_path = dir.path().join("test_table");
285+
jstable.write(file_path.to_str().unwrap()).unwrap();
282286

283-
let read_table = read_jstable(file.path().to_str().unwrap()).unwrap();
287+
let read_table = read_jstable(file_path.to_str().unwrap()).unwrap();
284288

285289
assert_eq!(read_table.timestamp, 12345);
286290
assert_eq!(read_table.collection, "test_col");
@@ -311,10 +315,11 @@ mod tests {
311315
documents.clone(),
312316
);
313317

314-
let file = NamedTempFile::new().unwrap();
315-
jstable.write(file.path().to_str().unwrap()).unwrap();
318+
let dir = tempdir()?;
319+
let file_path = dir.path().join("test_table");
320+
jstable.write(file_path.to_str().unwrap()).unwrap();
316321

317-
let iterator = JSTableIterator::new(file.path().to_str().unwrap())?;
322+
let iterator = JSTableIterator::new(file_path.to_str().unwrap())?;
318323
assert_eq!(iterator.timestamp, 12345);
319324
assert_eq!(iterator.collection, "test_col");
320325

@@ -353,10 +358,11 @@ mod tests {
353358
documents.clone(),
354359
);
355360

356-
let file = NamedTempFile::new().unwrap();
357-
jstable.write(file.path().to_str().unwrap()).unwrap();
361+
let dir = tempdir()?;
362+
let file_path = dir.path().join("test_table");
363+
jstable.write(file_path.to_str().unwrap()).unwrap();
358364

359-
let filter = read_filter(file.path().to_str().unwrap())?;
365+
let filter = read_filter(file_path.to_str().unwrap())?;
360366

361367
// Helper to hash string for filter check
362368
let hash = |s: &str| {
@@ -407,4 +413,26 @@ mod tests {
407413
);
408414
assert_eq!(merged_reverse.timestamp, 200);
409415
}
416+
417+
#[test]
418+
fn test_jstable_keys_sorted_on_disk() -> Result<(), Box<dyn std::error::Error>> {
419+
let schema = Schema::new(SchemaType::Object);
420+
let mut documents = BTreeMap::new();
421+
// Insert keys in non-sorted order (BTreeMap will sort them)
422+
documents.insert("c".to_string(), json!(3));
423+
documents.insert("a".to_string(), json!(1));
424+
documents.insert("b".to_string(), json!(2));
425+
426+
let jstable = JSTable::new(123, "sorted_test".to_string(), schema, documents);
427+
428+
let dir = tempdir()?;
429+
let file_path = dir.path().join("test_table");
430+
jstable.write(file_path.to_str().unwrap())?;
431+
432+
let iterator = JSTableIterator::new(file_path.to_str().unwrap())?;
433+
let keys: Vec<String> = iterator.map(|r| r.unwrap().0).collect();
434+
435+
assert_eq!(keys, vec!["a", "b", "c"]);
436+
Ok(())
437+
}
410438
}

0 commit comments

Comments
 (0)