use crate::types::{ColumnType, RowId, Value};
use crate::{Result, StorageError};
use std::fs::{File, OpenOptions};
use std::io::{BufWriter, Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use memmap2::Mmap;
const COLUMNAR_MAGIC: u32 = 0x434D5442; const COLUMNAR_VERSION: u32 = 2; const HEADER_SIZE: usize = 144; const FOOTER_SIZE: usize = 20;
pub const MAX_COLUMNS: usize = 128;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum ColumnTypeTag {
Integer = 0,
Float = 1,
Bool = 2,
Timestamp = 3,
Text = 4,
Vector = 5,
Spatial = 6,
}
impl ColumnTypeTag {
fn from_column_type(ct: &ColumnType) -> Self {
match ct {
ColumnType::Integer => Self::Integer,
ColumnType::Float => Self::Float,
ColumnType::Boolean => Self::Bool,
ColumnType::Timestamp => Self::Timestamp,
ColumnType::Text => Self::Text,
ColumnType::Tensor(_) => Self::Vector,
ColumnType::Spatial => Self::Spatial,
}
}
#[allow(dead_code)]
pub(crate) fn to_column_type(&self) -> ColumnType {
match self {
Self::Integer => ColumnType::Integer,
Self::Float => ColumnType::Float,
Self::Bool => ColumnType::Boolean,
Self::Timestamp => ColumnType::Timestamp,
Self::Text => ColumnType::Text,
Self::Vector => ColumnType::Tensor(0), Self::Spatial => ColumnType::Spatial,
}
}
pub(crate) fn is_fixed(&self) -> bool {
matches!(
self,
Self::Integer | Self::Float | Self::Bool | Self::Timestamp
)
}
fn fixed_size(&self) -> usize {
match self {
Self::Integer | Self::Float | Self::Timestamp => 8,
Self::Bool => 1,
_ => 0,
}
}
}
#[derive(Clone, Debug)]
pub(crate) struct ColumnarHeader {
num_rows: u32,
num_columns: u16,
column_tags: [u8; MAX_COLUMNS],
}
impl ColumnarHeader {
fn serialize(&self) -> [u8; HEADER_SIZE] {
let mut buf = [0u8; HEADER_SIZE];
buf[0..4].copy_from_slice(&COLUMNAR_MAGIC.to_le_bytes());
buf[4..8].copy_from_slice(&COLUMNAR_VERSION.to_le_bytes());
buf[8..12].copy_from_slice(&self.num_rows.to_le_bytes());
buf[12..14].copy_from_slice(&self.num_columns.to_le_bytes());
buf[14..14 + MAX_COLUMNS].copy_from_slice(&self.column_tags);
buf
}
fn deserialize(data: &[u8]) -> Result<Self> {
if data.len() < HEADER_SIZE {
return Err(StorageError::InvalidData(
"Columnar header too short".into(),
));
}
let magic = u32::from_le_bytes([data[0], data[1], data[2], data[3]]);
if magic != COLUMNAR_MAGIC {
return Err(StorageError::InvalidData(format!(
"Bad columnar magic: 0x{:08X}",
magic
)));
}
let version = u32::from_le_bytes([data[4], data[5], data[6], data[7]]);
if version != COLUMNAR_VERSION {
return Err(StorageError::InvalidData(format!(
"Unsupported columnar version: {}",
version
)));
}
let num_rows = u32::from_le_bytes([data[8], data[9], data[10], data[11]]);
let num_columns = u16::from_le_bytes([data[12], data[13]]);
let mut column_tags = [0u8; MAX_COLUMNS];
column_tags.copy_from_slice(&data[14..14 + MAX_COLUMNS]);
Ok(Self {
num_rows,
num_columns,
column_tags,
})
}
}
#[derive(Clone, Debug)]
pub struct ColumnIndexEntry {
pub offset: u64,
pub size: u64,
}
const COLUMN_INDEX_ENTRY_SIZE: usize = 16;
#[derive(Clone, Debug)]
pub struct RowMap {
pub num_rows: usize,
data: SegData, keys_offset: usize,
timestamps_offset: usize,
deleted_offset: usize,
#[allow(dead_code)]
deleted_len: usize,
}
impl RowMap {
fn compute_sizes(num_rows: usize) -> (usize, usize, usize, usize) {
let keys_size = num_rows * 8;
let timestamps_size = num_rows * 8;
let deleted_len = num_rows.div_ceil(8);
(
keys_size + timestamps_size + deleted_len,
keys_size,
timestamps_size,
deleted_len,
)
}
#[allow(dead_code)]
pub(crate) fn from_bytes(data: Vec<u8>, num_rows: usize) -> Self {
let (_, keys_size, timestamps_size, deleted_len) = Self::compute_sizes(num_rows);
Self {
num_rows,
keys_offset: 0,
timestamps_offset: keys_size,
deleted_offset: keys_size + timestamps_size,
deleted_len,
data: SegData::Owned(data),
}
}
#[allow(dead_code)]
pub(crate) fn from_mmap(mmap: Arc<Mmap>, offset: usize, num_rows: usize) -> Result<Self> {
let (_total, keys_size, timestamps_size, deleted_len) = Self::compute_sizes(num_rows);
Ok(Self {
num_rows,
keys_offset: offset,
timestamps_offset: offset + keys_size,
deleted_offset: offset + keys_size + timestamps_size,
deleted_len,
data: SegData::Mmap { mmap, offset },
})
}
#[inline]
pub fn key(&self, row_idx: usize) -> u64 {
let off = self.keys_offset + row_idx * 8;
let s = self.data.slice(off, 8);
u64::from_le_bytes([s[0], s[1], s[2], s[3], s[4], s[5], s[6], s[7]])
}
#[inline]
pub fn timestamp(&self, row_idx: usize) -> u64 {
let off = self.timestamps_offset + row_idx * 8;
let s = self.data.slice(off, 8);
u64::from_le_bytes([s[0], s[1], s[2], s[3], s[4], s[5], s[6], s[7]])
}
pub fn find_key(&self, target: u64) -> Option<usize> {
let mut lo = 0usize;
let mut hi = self.num_rows;
while lo < hi {
let mid = (lo + hi) / 2;
let k = self.key(mid);
if k < target {
lo = mid + 1;
} else if k > target {
hi = mid;
} else {
return Some(mid);
}
}
None
}
#[inline]
pub fn is_deleted(&self, row_idx: usize) -> bool {
let byte = self.data.get(self.deleted_offset + row_idx / 8);
(byte >> (row_idx % 8)) & 1 != 0
}
pub fn has_any_deleted(&self) -> bool {
let n = self.num_rows;
let nb = n.div_ceil(8);
for i in 0..nb {
if self.data.get(self.deleted_offset + i) != 0 {
return true;
}
}
false
}
}
#[derive(Clone, Debug)]
enum SegData {
#[allow(dead_code)]
Owned(Vec<u8>),
#[allow(dead_code)]
Mmap { mmap: Arc<Mmap>, offset: usize },
}
impl SegData {
#[inline]
fn get(&self, idx: usize) -> u8 {
match self {
SegData::Owned(v) => v[idx],
SegData::Mmap { mmap, offset } => mmap[offset + idx],
}
}
fn slice(&self, start: usize, len: usize) -> &[u8] {
match self {
SegData::Owned(v) => &v[start..start + len],
SegData::Mmap { mmap, offset } => &mmap[*offset + start..*offset + start + len],
}
}
fn len(&self) -> usize {
match self {
SegData::Owned(v) => v.len(),
SegData::Mmap { mmap, offset } => mmap.len().saturating_sub(*offset),
}
}
fn as_bytes(&self) -> &[u8] {
match self {
SegData::Owned(v) => v.as_slice(),
SegData::Mmap { mmap, offset } => &mmap[*offset..],
}
}
}
#[derive(Clone)]
pub struct FixedSegment {
pub num_rows: usize,
null_bitmap: SegData,
data: SegData,
#[allow(dead_code)]
elem_size: usize,
#[allow(dead_code)]
tag: ColumnTypeTag,
}
impl FixedSegment {
#[allow(dead_code)]
pub(crate) fn from_bytes(data: &[u8], num_rows: usize, tag: ColumnTypeTag) -> Result<Self> {
let null_bytes = num_rows.div_ceil(8);
let elem_size = tag.fixed_size();
let data_size = num_rows * elem_size;
let expected = null_bytes + data_size;
if data.len() < expected {
return Err(StorageError::InvalidData(format!(
"Fixed segment too short: {} < {}",
data.len(),
expected
)));
}
Ok(Self {
num_rows,
null_bitmap: SegData::Owned(data[..null_bytes].to_vec()),
data: SegData::Owned(data[null_bytes..null_bytes + data_size].to_vec()),
elem_size,
tag,
})
}
#[allow(dead_code)]
pub(crate) fn from_mmap(
mmap: Arc<Mmap>,
offset: usize,
num_rows: usize,
tag: ColumnTypeTag,
) -> Self {
let null_bytes = num_rows.div_ceil(8);
Self {
num_rows,
null_bitmap: SegData::Mmap {
mmap: mmap.clone(),
offset,
},
data: SegData::Mmap {
mmap,
offset: offset + null_bytes,
},
elem_size: tag.fixed_size(),
tag,
}
}
#[inline]
pub fn is_null(&self, row_idx: usize) -> bool {
(self.null_bitmap.get(row_idx / 8) >> (row_idx % 8)) & 1 != 0
}
pub fn has_nulls(&self) -> bool {
let nb = self.null_bitmap.len();
for i in 0..nb {
if self.null_bitmap.get(i) != 0 {
return true;
}
}
false
}
pub fn raw_f64_slice(&self) -> &[u8] {
self.data.as_bytes()
}
#[inline]
pub fn get_i64(&self, row_idx: usize) -> Option<i64> {
if self.is_null(row_idx) {
return None;
}
let off = row_idx * 8;
let s = self.data.slice(off, 8);
Some(i64::from_le_bytes([
s[0], s[1], s[2], s[3], s[4], s[5], s[6], s[7],
]))
}
#[inline]
pub fn get_f64(&self, row_idx: usize) -> Option<f64> {
if self.is_null(row_idx) {
return None;
}
let off = row_idx * 8;
let s = self.data.slice(off, 8);
Some(f64::from_le_bytes([
s[0], s[1], s[2], s[3], s[4], s[5], s[6], s[7],
]))
}
#[inline]
pub fn get_bool(&self, row_idx: usize) -> Option<bool> {
if self.is_null(row_idx) {
return None;
}
Some(self.data.get(row_idx) != 0)
}
}
#[derive(Clone)]
pub struct TextSegment {
pub num_rows: usize,
null_bitmap: SegData,
offsets_data: SegData,
string_data: SegData,
pub trust_utf8: bool,
#[allow(dead_code)]
offsets_start: usize,
}
impl TextSegment {
pub(crate) fn from_bytes(data: &[u8], num_rows: usize) -> Result<Self> {
let null_bytes = num_rows.div_ceil(8);
let offsets_size = (num_rows + 1) * 4;
if data.len() < null_bytes + offsets_size {
return Err(StorageError::InvalidData("Text segment too short".into()));
}
Ok(Self {
num_rows,
null_bitmap: SegData::Owned(data[..null_bytes].to_vec()),
offsets_data: SegData::Owned(data[null_bytes..null_bytes + offsets_size].to_vec()),
string_data: SegData::Owned(data[null_bytes + offsets_size..].to_vec()),
trust_utf8: false,
offsets_start: 0,
})
}
#[allow(dead_code)]
pub(crate) fn from_mmap(mmap: Arc<Mmap>, offset: usize, num_rows: usize) -> Self {
let null_bytes = num_rows.div_ceil(8);
let offsets_size = (num_rows + 1) * 4;
Self {
num_rows,
null_bitmap: SegData::Mmap {
mmap: mmap.clone(),
offset,
},
offsets_data: SegData::Mmap {
mmap: mmap.clone(),
offset: offset + null_bytes,
},
string_data: SegData::Mmap {
mmap,
offset: offset + null_bytes + offsets_size,
},
trust_utf8: true, offsets_start: 0,
}
}
#[inline]
pub fn is_null(&self, row_idx: usize) -> bool {
(self.null_bitmap.get(row_idx / 8) >> (row_idx % 8)) & 1 != 0
}
pub fn prefix_match_indices(&self, prefix: &[u8]) -> Vec<usize> {
let n = self.num_rows;
let plen = prefix.len();
if plen == 0 || n == 0 {
return (0..n).collect();
}
let mut result = Vec::with_capacity(n / 4);
let off_bytes = self.offsets_data.as_bytes();
let str_bytes = self.string_data.as_bytes();
let has_nulls = self.has_any_null();
if !has_nulls && off_bytes.len() >= (n + 1) * 4 {
for i in 0..n {
let ob = i * 4;
let start = u32::from_le_bytes([
off_bytes[ob],
off_bytes[ob + 1],
off_bytes[ob + 2],
off_bytes[ob + 3],
]) as usize;
let end = u32::from_le_bytes([
off_bytes[ob + 4],
off_bytes[ob + 5],
off_bytes[ob + 6],
off_bytes[ob + 7],
]) as usize;
if end - start >= plen {
let candidate = &str_bytes[start..start + plen];
let matched = if !prefix.contains(&b'_') {
candidate == prefix
} else {
candidate
.iter()
.zip(prefix.iter())
.all(|(&c, &p)| p == b'_' || p == c)
};
if matched {
result.push(i);
}
}
}
} else {
for i in 0..n {
if has_nulls && self.is_null(i) {
continue;
}
let s = self.get_str_fast(i);
if s.len() >= plen && &s.as_bytes()[..plen] == prefix {
result.push(i);
}
}
}
result
}
pub fn eq_match_indices(&self, target: &[u8]) -> Vec<usize> {
let n = self.num_rows;
let tlen = target.len();
if n == 0 {
return Vec::new();
}
let mut result = Vec::with_capacity(n / 8);
let off_bytes = self.offsets_data.as_bytes();
let str_bytes = self.string_data.as_bytes();
let has_nulls = self.has_any_null();
if !has_nulls && off_bytes.len() >= (n + 1) * 4 {
for i in 0..n {
let ob = i * 4;
let start = u32::from_le_bytes([
off_bytes[ob],
off_bytes[ob + 1],
off_bytes[ob + 2],
off_bytes[ob + 3],
]) as usize;
let end = u32::from_le_bytes([
off_bytes[ob + 4],
off_bytes[ob + 5],
off_bytes[ob + 6],
off_bytes[ob + 7],
]) as usize;
if end - start == tlen && &str_bytes[start..end] == target {
result.push(i);
}
}
} else {
for i in 0..n {
if has_nulls && self.is_null(i) {
continue;
}
let s = self.get_str_fast(i);
if s.as_bytes() == target {
result.push(i);
}
}
}
result
}
pub fn in_set_match_indices(&self, targets: &std::collections::HashSet<&[u8]>) -> Vec<usize> {
let n = self.num_rows;
if n == 0 || targets.is_empty() {
return Vec::new();
}
let mut result = Vec::with_capacity(n / 8);
let off_bytes = self.offsets_data.as_bytes();
let str_bytes = self.string_data.as_bytes();
let has_nulls = self.has_any_null();
if !has_nulls && off_bytes.len() >= (n + 1) * 4 {
for i in 0..n {
let ob = i * 4;
let start = u32::from_le_bytes([
off_bytes[ob],
off_bytes[ob + 1],
off_bytes[ob + 2],
off_bytes[ob + 3],
]) as usize;
let end = u32::from_le_bytes([
off_bytes[ob + 4],
off_bytes[ob + 5],
off_bytes[ob + 6],
off_bytes[ob + 7],
]) as usize;
let slice = &str_bytes[start..end];
if targets.contains(slice) {
result.push(i);
}
}
} else {
for i in 0..n {
if has_nulls && self.is_null(i) {
continue;
}
let s = self.get_str_fast(i);
if targets.contains(s.as_bytes()) {
result.push(i);
}
}
}
result
}
pub fn for_each_str<F: FnMut(&str)>(&self, mut f: F) {
let n = self.num_rows;
let has_nulls = self.has_any_null();
for i in 0..n {
if has_nulls && self.is_null(i) {
continue;
}
let s = self.get_str_fast(i);
f(s);
}
}
pub fn has_any_null(&self) -> bool {
let nb = self.null_bitmap.len();
for i in 0..nb {
if self.null_bitmap.get(i) != 0 {
return true;
}
}
false
}
#[inline]
fn get_offset(&self, idx: usize) -> u32 {
let s = self.offsets_data.slice(idx * 4, 4);
u32::from_le_bytes([s[0], s[1], s[2], s[3]])
}
#[inline]
pub fn get_str(&self, row_idx: usize) -> Option<&str> {
if self.is_null(row_idx) {
return None;
}
let start = self.get_offset(row_idx) as usize;
let end = self.get_offset(row_idx + 1) as usize;
if start > end {
return None;
}
let bytes = self.string_data.slice(start, end - start);
if self.trust_utf8 {
unsafe { Some(std::str::from_utf8_unchecked(bytes)) }
} else {
std::str::from_utf8(bytes).ok()
}
}
#[inline]
pub fn get_str_fast(&self, row_idx: usize) -> &str {
let off_base = row_idx * 4;
let start_bytes = self.offsets_data.slice(off_base, 4);
let end_bytes = self.offsets_data.slice(off_base + 4, 4);
let start = u32::from_le_bytes([
start_bytes[0],
start_bytes[1],
start_bytes[2],
start_bytes[3],
]) as usize;
let end =
u32::from_le_bytes([end_bytes[0], end_bytes[1], end_bytes[2], end_bytes[3]]) as usize;
let bytes = self.string_data.slice(start, end - start);
if self.trust_utf8 {
unsafe { std::str::from_utf8_unchecked(bytes) }
} else {
std::str::from_utf8(bytes).unwrap_or("")
}
}
#[cfg(feature = "rayon")]
pub fn extract_all_raw_keys_par(&self) -> Vec<([u8; 64], usize)> {
use rayon::prelude::*;
let n = self.num_rows;
if self.has_any_null() || n < 50000 {
return self.extract_all_raw_keys_unchecked();
}
let offsets_len = (n + 1) * 4;
let offsets_bytes: &[u8] = self.offsets_data.slice(0, offsets_len);
let total_str_len = self.string_data.len();
let string_bytes: &[u8] = self.string_data.slice(0, total_str_len);
(0..n)
.into_par_iter()
.map(|i| {
let off_base = i * 4;
let start = u32::from_le_bytes([
offsets_bytes[off_base],
offsets_bytes[off_base + 1],
offsets_bytes[off_base + 2],
offsets_bytes[off_base + 3],
]) as usize;
let end = u32::from_le_bytes([
offsets_bytes[off_base + 4],
offsets_bytes[off_base + 5],
offsets_bytes[off_base + 6],
offsets_bytes[off_base + 7],
]) as usize;
let len = (end - start).min(64);
let mut buf = [0u8; 64];
buf[..len].copy_from_slice(&string_bytes[start..start + len]);
(buf, i)
})
.collect()
}
pub fn extract_all_raw_keys_unchecked(&self) -> Vec<([u8; 64], usize)> {
let n = self.num_rows;
let mut result: Vec<([u8; 64], usize)> = Vec::with_capacity(n);
if self.has_any_null() {
return self.bulk_extract_raw_keys();
}
let offsets_len = (n + 1) * 4;
let offsets_bytes = self.offsets_data.slice(0, offsets_len);
let string_bytes = self.string_data.slice(0, self.string_data.len());
for i in 0..n {
let off_base = i * 4;
let start = u32::from_le_bytes([
offsets_bytes[off_base],
offsets_bytes[off_base + 1],
offsets_bytes[off_base + 2],
offsets_bytes[off_base + 3],
]) as usize;
let end = u32::from_le_bytes([
offsets_bytes[off_base + 4],
offsets_bytes[off_base + 5],
offsets_bytes[off_base + 6],
offsets_bytes[off_base + 7],
]) as usize;
let len = (end - start).min(64);
let mut buf = [0u8; 64];
buf[..len].copy_from_slice(&string_bytes[start..start + len]);
result.push((buf, i));
}
result
}
pub fn bulk_extract_raw_keys(&self) -> Vec<([u8; 64], usize)> {
let n = self.num_rows;
let mut result: Vec<([u8; 64], usize)> = Vec::with_capacity(n);
if !self.has_any_null() {
for i in 0..n {
let off_base = i * 4;
let off_bytes = self.offsets_data.slice(off_base, 8);
let start =
u32::from_le_bytes([off_bytes[0], off_bytes[1], off_bytes[2], off_bytes[3]])
as usize;
let end =
u32::from_le_bytes([off_bytes[4], off_bytes[5], off_bytes[6], off_bytes[7]])
as usize;
let len = (end - start).min(64);
let mut buf = [0u8; 64];
let src = self.string_data.slice(start, len);
buf[..len].copy_from_slice(src);
result.push((buf, i));
}
} else {
for i in 0..n {
if self.is_null(i) {
continue;
}
let start = self.get_offset(i) as usize;
let end = self.get_offset(i + 1) as usize;
let len = (end - start).min(64);
let mut buf = [0u8; 64];
let src = self.string_data.slice(start, len);
buf[..len].copy_from_slice(src);
result.push((buf, i));
}
}
result
}
}
pub struct ColumnarSSTable {
pub path: PathBuf,
file_data: Vec<u8>,
#[allow(dead_code)]
mmap: Option<Arc<Mmap>>,
file: Option<parking_lot::Mutex<File>>,
#[allow(dead_code)]
header: ColumnarHeader,
pub column_index: Vec<ColumnIndexEntry>,
pub row_map: RowMap,
pub column_tags: Vec<ColumnTypeTag>,
pub num_rows: usize,
}
impl ColumnarSSTable {
#[inline]
pub fn release_pages(&self) {
if let Some(ref m) = self.mmap {
unsafe {
libc::madvise(m.as_ptr() as *mut _, m.len(), libc::MADV_DONTNEED);
}
}
}
pub fn is_columnar<P: AsRef<Path>>(path: P) -> bool {
let path = path.as_ref();
if let Ok(mut file) = OpenOptions::new().read(true).open(path) {
if let Ok(metadata) = file.metadata() {
let file_len = metadata.len();
if file_len >= FOOTER_SIZE as u64
&& file.seek(SeekFrom::End(-(FOOTER_SIZE as i64))).is_ok()
{
let mut footer = [0u8; FOOTER_SIZE];
if file.read_exact(&mut footer).is_ok() {
let magic =
u32::from_le_bytes([footer[16], footer[17], footer[18], footer[19]]);
return magic == COLUMNAR_MAGIC;
}
}
}
}
false
}
pub fn open<P: AsRef<Path>>(path: P) -> Result<Self> {
let path = path.as_ref().to_path_buf();
let mut file = OpenOptions::new().read(true).open(&path)?;
let file_len = file.metadata()?.len();
if file_len < FOOTER_SIZE as u64 {
return Err(StorageError::InvalidData(
"File too small for columnar footer".into(),
));
}
file.seek(SeekFrom::End(-(FOOTER_SIZE as i64)))?;
let mut footer_buf = [0u8; FOOTER_SIZE];
file.read_exact(&mut footer_buf)?;
let magic = u32::from_le_bytes([
footer_buf[16],
footer_buf[17],
footer_buf[18],
footer_buf[19],
]);
if magic != COLUMNAR_MAGIC {
return Err(StorageError::InvalidData("Not a columnar SSTable".into()));
}
let _column_index_offset = u64::from_le_bytes([
footer_buf[0],
footer_buf[1],
footer_buf[2],
footer_buf[3],
footer_buf[4],
footer_buf[5],
footer_buf[6],
footer_buf[7],
]);
let row_map_offset = u64::from_le_bytes([
footer_buf[8],
footer_buf[9],
footer_buf[10],
footer_buf[11],
footer_buf[12],
footer_buf[13],
footer_buf[14],
footer_buf[15],
]);
let lazy_load = file_len > 256 * 1024;
let mmap: Option<Arc<Mmap>> = None;
let mut file_data: Vec<u8> = Vec::new();
if !lazy_load {
file.seek(SeekFrom::Start(0))?;
file_data = vec![0u8; file_len as usize];
file.read_exact(&mut file_data)?;
}
let header = if !file_data.is_empty() {
ColumnarHeader::deserialize(&file_data[..HEADER_SIZE])?
} else {
file.seek(SeekFrom::Start(0))?;
let mut hb = vec![0u8; HEADER_SIZE];
file.read_exact(&mut hb)?;
ColumnarHeader::deserialize(&hb)?
};
let num_columns = header.num_columns as usize;
let num_rows = header.num_rows as usize;
let ci_size = num_columns * COLUMN_INDEX_ENTRY_SIZE;
let ci_start = HEADER_SIZE;
let ci_buf = if !file_data.is_empty() {
file_data[ci_start..ci_start + ci_size].to_vec()
} else {
let mut b = vec![0u8; ci_size];
file.seek(SeekFrom::Start(ci_start as u64))?;
file.read_exact(&mut b)?;
b
};
let ci_data: &[u8] = &ci_buf;
let column_index: Vec<ColumnIndexEntry> = (0..num_columns)
.map(|i| {
let off = i * COLUMN_INDEX_ENTRY_SIZE;
ColumnIndexEntry {
offset: u64::from_le_bytes([
ci_data[off],
ci_data[off + 1],
ci_data[off + 2],
ci_data[off + 3],
ci_data[off + 4],
ci_data[off + 5],
ci_data[off + 6],
ci_data[off + 7],
]),
size: u64::from_le_bytes([
ci_data[off + 8],
ci_data[off + 9],
ci_data[off + 10],
ci_data[off + 11],
ci_data[off + 12],
ci_data[off + 13],
ci_data[off + 14],
ci_data[off + 15],
]),
}
})
.collect();
let (rm_total, keys_size, timestamps_size, deleted_len) = RowMap::compute_sizes(num_rows);
let rm_data = if !file_data.is_empty() {
file_data[row_map_offset as usize..row_map_offset as usize + rm_total].to_vec()
} else {
let mut d = vec![0u8; rm_total];
file.seek(SeekFrom::Start(row_map_offset))?;
file.read_exact(&mut d)?;
d
};
let row_map = RowMap {
num_rows,
keys_offset: 0,
timestamps_offset: keys_size,
deleted_offset: keys_size + timestamps_size,
deleted_len,
data: SegData::Owned(rm_data),
};
let column_tags: Vec<ColumnTypeTag> = header.column_tags[..num_columns]
.iter()
.map(|&t| unsafe { std::mem::transmute(t) })
.collect();
let file = if file_data.is_empty() {
std::fs::File::open(&path).ok().map(parking_lot::Mutex::new)
} else {
None
};
Ok(Self {
path,
file_data,
mmap,
file,
header,
column_index,
row_map,
column_tags,
num_rows,
})
}
fn decompress_segment(data: &[u8]) -> std::borrow::Cow<'_, [u8]> {
if data.is_empty() {
return std::borrow::Cow::Borrowed(data);
}
match data[0] {
1 => {
match snap::raw::Decoder::new().decompress_vec(&data[1..]) {
Ok(v) => std::borrow::Cow::Owned(v),
Err(_) => std::borrow::Cow::Borrowed(&data[1..]), }
}
_ => std::borrow::Cow::Borrowed(&data[1..]), }
}
pub fn read_fixed_i64(&self, col_idx: usize) -> Result<FixedSegment> {
let tag = self.column_tags[col_idx];
if !tag.is_fixed() {
return Err(StorageError::InvalidData(
"Column is not fixed-width".into(),
));
}
let entry = &self.column_index[col_idx];
let start = entry.offset as usize;
let end = start + entry.size as usize;
let seg_bytes = self.read_segment_bytes(start, end);
FixedSegment::from_bytes(&seg_bytes, self.num_rows, tag)
}
pub fn read_fixed_f64(&self, col_idx: usize) -> Result<FixedSegment> {
self.read_fixed_i64(col_idx)
}
pub fn read_text(&self, col_idx: usize) -> Result<TextSegment> {
let entry = &self.column_index[col_idx];
let start = entry.offset as usize;
let end = start + entry.size as usize;
let seg_bytes = self.read_segment_bytes(start, end);
TextSegment::from_bytes(&seg_bytes, self.num_rows)
}
pub fn read_segment_bytes(&self, start: usize, end: usize) -> std::borrow::Cow<'_, [u8]> {
if !self.file_data.is_empty() {
return Self::decompress_segment(&self.file_data[start..end]);
}
if let Some(ref mmap) = self.mmap {
return Self::decompress_segment(&mmap[start..end]);
}
let len = end - start;
let mut buf = vec![0u8; len];
use std::io::{Read, Seek};
let ok = if let Some(ref cached) = self.file {
let mut f = cached.lock();
f.seek(SeekFrom::Start(start as u64)).is_ok() && f.read_exact(&mut buf).is_ok()
} else if let Ok(mut f) = std::fs::File::open(&self.path) {
f.seek(SeekFrom::Start(start as u64)).is_ok() && f.read_exact(&mut buf).is_ok()
} else {
false
};
if ok {
Self::decompress_segment(&buf).into_owned().into()
} else {
std::borrow::Cow::Owned(Vec::new())
}
}
pub fn read_spatial(&self, col_idx: usize) -> Result<Vec<(RowId, crate::types::Geometry)>> {
let entry = &self.column_index[col_idx];
let seg_start = entry.offset as usize;
let seg_end = seg_start + entry.size as usize;
let seg_bytes = self.read_segment_bytes(seg_start, seg_end);
let data = seg_bytes.as_ref();
let null_bytes = self.num_rows.div_ceil(8);
if null_bytes + 2 > data.len() {
return Ok(Vec::new());
}
let mut result = Vec::new();
let mut pos = null_bytes;
for i in 0..self.num_rows {
if (data[i / 8] >> (i % 8)) & 1 != 0 {
if pos + 2 <= data.len() {
let len = u16::from_le_bytes([data[pos], data[pos + 1]]) as usize;
pos += 2 + len;
}
continue;
}
if self.row_map.is_deleted(i) {
continue;
}
if pos + 2 > data.len() {
break;
}
let len = u16::from_le_bytes([data[pos], data[pos + 1]]) as usize;
pos += 2;
if len == 0 || pos + len > data.len() {
continue;
}
if let Ok(geom) = bincode::deserialize::<crate::types::Geometry>(&data[pos..pos + len])
{
let row_id = (self.row_map.key(i) & 0xFFFFFFFF) as RowId;
result.push((row_id, geom));
}
pos += len;
}
Ok(result)
}
pub fn get_row(&self, key: u64, col_types: &[ColumnType]) -> Option<Vec<Value>> {
let idx = self.row_map.find_key(key)?;
if self.row_map.is_deleted(idx) {
return None;
}
let mut row = Vec::with_capacity(col_types.len());
for ci in 0..col_types.len() {
if self.column_tags[ci].is_fixed() {
if let Ok(seg) = self.read_fixed_i64(ci) {
match &col_types[ci] {
crate::types::ColumnType::Integer => row.push(
seg.get_i64(idx)
.map(crate::types::Value::Integer)
.unwrap_or(crate::types::Value::Null),
),
crate::types::ColumnType::Float => row.push(
seg.get_f64(idx)
.map(crate::types::Value::Float)
.unwrap_or(crate::types::Value::Null),
),
crate::types::ColumnType::Boolean => row.push(
seg.get_bool(idx)
.map(crate::types::Value::Bool)
.unwrap_or(crate::types::Value::Null),
),
crate::types::ColumnType::Timestamp => row.push(
seg.get_i64(idx)
.map(|v| {
crate::types::Value::Timestamp(
crate::types::Timestamp::from_micros(v),
)
})
.unwrap_or(crate::types::Value::Null),
),
_ => row.push(crate::types::Value::Null),
}
} else {
row.push(crate::types::Value::Null);
}
} else if let Ok(seg) = self.read_text(ci) {
row.push(
seg.get_str(idx)
.map(|s| {
crate::types::Value::Text(crate::types::ArcString(
std::sync::Arc::from(s),
))
})
.unwrap_or(crate::types::Value::Null),
);
} else {
row.push(crate::types::Value::Null);
}
}
Some(row)
}
pub fn read_vectors(&self, col_idx: usize) -> Result<Vec<(RowId, Vec<f32>)>> {
let entry = &self.column_index[col_idx];
let seg_bytes =
self.read_segment_bytes(entry.offset as usize, (entry.offset + entry.size) as usize);
let data = seg_bytes.as_ref();
let null_bytes = self.num_rows.div_ceil(8);
if null_bytes + 2 > data.len() {
return Ok(Vec::new());
}
let dim = u16::from_le_bytes([data[null_bytes], data[null_bytes + 1]]) as usize;
if dim == 0 || dim > 65536 {
return Ok(Vec::new());
}
let stride = dim * 4;
let data_start = null_bytes + 2;
let n = ((data.len() - data_start) / stride).min(self.num_rows);
let mut result = Vec::with_capacity(n);
for i in 0..n {
if (data[i / 8] >> (i % 8)) & 1 != 0 {
continue;
} if self.row_map.is_deleted(i) {
continue;
}
let row_id = (self.row_map.key(i) & 0xFFFFFFFF) as RowId;
let mut v = Vec::with_capacity(dim);
let base = data_start + i * stride;
for j in 0..dim {
let off = base + j * 4;
v.push(f32::from_le_bytes([
data[off],
data[off + 1],
data[off + 2],
data[off + 3],
]));
}
result.push((row_id, v));
}
Ok(result)
}
}
pub struct ColumnarSSTableBuilder {
pub path: PathBuf,
pub column_types: Vec<ColumnType>,
pub column_tags: Vec<ColumnTypeTag>,
pub num_rows: usize,
pub(crate) keys: Vec<u64>,
timestamps: Vec<u64>,
pub(crate) deleted: Vec<bool>,
pub(crate) column_buffers: Vec<Vec<u8>>,
pub(crate) null_flags: Vec<Vec<bool>>,
finished: bool,
}
impl ColumnarSSTableBuilder {
pub fn new<P: AsRef<Path>>(path: P, column_types: Vec<ColumnType>) -> Self {
let column_tags: Vec<ColumnTypeTag> = column_types
.iter()
.map(ColumnTypeTag::from_column_type)
.collect();
let num_cols = column_types.len();
Self {
path: path.as_ref().to_path_buf(),
column_types,
column_tags,
num_rows: 0,
keys: Vec::new(),
timestamps: Vec::new(),
deleted: Vec::new(),
column_buffers: vec![Vec::new(); num_cols],
null_flags: vec![Vec::new(); num_cols],
finished: false,
}
}
pub fn add_values(
&mut self,
key: u64,
timestamp: u64,
deleted: bool,
row: &[Value],
) -> Result<()> {
self.keys.push(key);
self.timestamps.push(timestamp);
self.deleted.push(deleted);
for (col_idx, value) in row.iter().enumerate() {
if col_idx >= self.column_buffers.len() {
break;
}
let buf = &mut self.column_buffers[col_idx];
self.null_flags[col_idx].push(matches!(value, Value::Null));
match &self.column_tags[col_idx] {
ColumnTypeTag::Integer => {
let i = match value {
Value::Integer(v) => *v,
Value::Null => i64::MIN,
Value::Float(f) => f.to_bits() as i64,
_ => 0,
};
buf.extend_from_slice(&i.to_le_bytes());
}
ColumnTypeTag::Float => {
let f = match value {
Value::Float(v) => *v,
Value::Null => f64::NAN,
_ => 0.0,
};
buf.extend_from_slice(&f.to_le_bytes());
}
ColumnTypeTag::Bool => {
let b = match value {
Value::Bool(v) => {
if *v {
1
} else {
0
}
}
_ => 2,
};
buf.push(b);
}
ColumnTypeTag::Timestamp => {
let ts = match value {
Value::Timestamp(t) => t.as_micros(),
Value::Null => i64::MIN,
_ => 0,
};
buf.extend_from_slice(&ts.to_le_bytes());
}
ColumnTypeTag::Text => {
match value {
Value::Null => {
buf.extend_from_slice(&0xFFFFu16.to_le_bytes());
}
Value::Text(t) => {
let s = t.as_str();
let len = s.len().min(65534) as u16; buf.extend_from_slice(&len.to_le_bytes());
buf.extend_from_slice(&s.as_bytes()[..len as usize]);
}
_ => {
buf.extend_from_slice(&0xFFFFu16.to_le_bytes());
}
}
}
ColumnTypeTag::Vector => {
match value {
Value::Vector(v) => {
let floats: &[f32] = &v.0;
buf.extend_from_slice(&(floats.len() as u16).to_le_bytes());
for f in floats {
buf.extend_from_slice(&f.to_le_bytes());
}
}
Value::Tensor(t) => {
let floats = t.to_f32();
buf.extend_from_slice(&(floats.len() as u16).to_le_bytes());
for f in &floats {
buf.extend_from_slice(&f.to_le_bytes());
}
}
_ => buf.extend_from_slice(&0u16.to_le_bytes()),
}
}
ColumnTypeTag::Spatial => {
match value {
Value::Spatial(g) => {
let bytes = bincode::serialize(&**g).unwrap_or_default();
let len = bytes.len().min(65535) as u16;
buf.extend_from_slice(&len.to_le_bytes());
buf.extend_from_slice(&bytes[..len as usize]);
}
_ => buf.extend_from_slice(&0u16.to_le_bytes()),
}
}
}
}
self.num_rows += 1;
Ok(())
}
pub fn add_values_raw(
&mut self,
key: u64,
timestamp: u64,
deleted: bool,
col_raw: &[&[u8]],
) -> Result<()> {
self.keys.push(key);
self.timestamps.push(timestamp);
self.deleted.push(deleted);
for (col_idx, bytes) in col_raw.iter().enumerate() {
let is_null = match self.column_tags.get(col_idx) {
Some(crate::storage::lsm::columnar::ColumnTypeTag::Vector) => {
bytes.len() >= 2 && u16::from_le_bytes([bytes[0], bytes[1]]) == 0
}
Some(crate::storage::lsm::columnar::ColumnTypeTag::Spatial) => {
bytes.len() >= 2 && u16::from_le_bytes([bytes[0], bytes[1]]) == 0
}
_ => false,
};
self.null_flags[col_idx].push(is_null);
self.column_buffers[col_idx].extend_from_slice(bytes);
}
self.num_rows += 1;
Ok(())
}
pub fn add_values_raw_with_nulls(
&mut self,
key: u64,
timestamp: u64,
deleted: bool,
col_raw: &[&[u8]],
col_nulls: &[bool],
) -> Result<()> {
self.keys.push(key);
self.timestamps.push(timestamp);
self.deleted.push(deleted);
for (col_idx, bytes) in col_raw.iter().enumerate() {
let inferred = match self.column_tags.get(col_idx) {
Some(crate::storage::lsm::columnar::ColumnTypeTag::Vector) => {
bytes.len() >= 2 && u16::from_le_bytes([bytes[0], bytes[1]]) == 0
}
Some(crate::storage::lsm::columnar::ColumnTypeTag::Spatial) => {
bytes.len() >= 2 && u16::from_le_bytes([bytes[0], bytes[1]]) == 0
}
_ => false,
};
let is_null = col_nulls.get(col_idx).copied().unwrap_or(inferred);
self.null_flags[col_idx].push(is_null);
self.column_buffers[col_idx].extend_from_slice(bytes);
}
self.num_rows += 1;
Ok(())
}
pub fn add_row(
&mut self,
key: u64,
timestamp: u64,
deleted: bool,
row_data: &[u8],
) -> Result<()> {
use crate::storage::row_format;
let col_types = &self.column_types;
let row: Vec<Value> = row_format::decode(row_data, col_types)?;
self.keys.push(key);
self.timestamps.push(timestamp);
self.deleted.push(deleted);
for (col_idx, value) in row.iter().enumerate() {
if col_idx >= self.column_buffers.len() {
break;
}
let buf = &mut self.column_buffers[col_idx];
self.null_flags[col_idx].push(matches!(value, Value::Null));
match &self.column_tags[col_idx] {
ColumnTypeTag::Integer => {
let i = match value {
Value::Integer(v) => *v,
Value::Null => i64::MIN, _ => 0,
};
buf.extend_from_slice(&i.to_le_bytes());
}
ColumnTypeTag::Float => {
let f = match value {
Value::Float(v) => *v,
Value::Null => f64::NAN, _ => 0.0,
};
buf.extend_from_slice(&f.to_le_bytes());
}
ColumnTypeTag::Bool => {
let b = match value {
Value::Bool(v) => {
if *v {
1
} else {
0
}
}
_ => 2, };
buf.push(b);
}
ColumnTypeTag::Timestamp => {
let ts = match value {
Value::Timestamp(t) => t.as_micros(),
Value::Null => i64::MIN,
_ => 0,
};
buf.extend_from_slice(&ts.to_le_bytes());
}
ColumnTypeTag::Text => {
let s = match value {
Value::Text(t) => t.as_str().to_string(),
Value::Null => String::new(),
_ => String::new(),
};
let len = s.len().min(65535) as u16;
buf.extend_from_slice(&len.to_le_bytes());
buf.extend_from_slice(s.as_bytes());
}
ColumnTypeTag::Vector => {
let bytes = match value {
Value::Vector(v) => {
let floats: &[f32] = &v.0;
let mut b = Vec::with_capacity(2 + floats.len() * 4);
b.extend_from_slice(&(floats.len() as u16).to_le_bytes());
for f in floats {
b.extend_from_slice(&f.to_le_bytes());
}
b
}
_ => vec![0u8; 2],
};
buf.extend_from_slice(&bytes);
}
ColumnTypeTag::Spatial => {
let wkt = match value {
Value::Spatial(g) => {
use crate::types::Geometry;
match **g {
Geometry::Point3D(ref p) => {
format!("POINT({},{},{})", p.x, p.y, p.z)
}
_ => String::new(),
}
}
Value::Null => String::new(),
_ => String::new(),
};
let bytes = wkt.as_bytes();
let len = bytes.len().min(65535) as u16;
buf.extend_from_slice(&len.to_le_bytes());
buf.extend_from_slice(bytes);
}
}
}
self.num_rows += 1;
Ok(())
}
pub fn finish(mut self) -> Result<()> {
self.finish_and_reset()
}
pub fn check_key(&self, key: u64) -> Option<bool> {
for i in (0..self.num_rows).rev() {
if self.keys[i] == key {
return Some(self.deleted[i]);
}
}
None
}
pub fn latest_entries(&self) -> Vec<(u64, bool)> {
let mut latest: std::collections::HashMap<u64, bool> =
std::collections::HashMap::with_capacity(self.num_rows);
for i in 0..self.num_rows {
latest.insert(self.keys[i], self.deleted[i]);
}
latest.into_iter().collect()
}
fn dedup_keys_newest_wins(&mut self) {
let mut last_idx: std::collections::HashMap<u64, usize> =
std::collections::HashMap::with_capacity(self.num_rows);
for (i, &k) in self.keys.iter().enumerate() {
last_idx.insert(k, i);
}
let keep: Vec<usize> = (0..self.num_rows)
.filter(|&i| last_idx.get(&self.keys[i]) == Some(&i))
.collect();
if keep.len() == self.num_rows {
return; }
let col_types = self.column_types.clone();
let mut kept_rows: Vec<(u64, u64, bool, Vec<Value>)> = Vec::with_capacity(keep.len());
for &i in &keep {
let mut row = Vec::with_capacity(col_types.len());
for (ci, tag) in self.column_tags.iter().enumerate() {
match tag {
ColumnTypeTag::Integer | ColumnTypeTag::Timestamp => {
let buf = &self.column_buffers[ci];
let off = i * 8;
if off + 8 > buf.len() {
row.push(Value::Null);
continue;
}
if self.null_flags.get(ci).and_then(|f| f.get(i)) == Some(&true) {
row.push(Value::Null);
continue;
}
let val = i64::from_le_bytes([
buf[off],
buf[off + 1],
buf[off + 2],
buf[off + 3],
buf[off + 4],
buf[off + 5],
buf[off + 6],
buf[off + 7],
]);
if matches!(tag, ColumnTypeTag::Timestamp) {
row.push(Value::Timestamp(crate::types::Timestamp::from_micros(val)));
} else {
row.push(Value::Integer(val));
}
}
ColumnTypeTag::Float => {
let buf = &self.column_buffers[ci];
let off = i * 8;
if off + 8 > buf.len() {
row.push(Value::Null);
continue;
}
if self.null_flags.get(ci).and_then(|f| f.get(i)) == Some(&true) {
row.push(Value::Null);
continue;
}
let bits = u64::from_le_bytes([
buf[off],
buf[off + 1],
buf[off + 2],
buf[off + 3],
buf[off + 4],
buf[off + 5],
buf[off + 6],
buf[off + 7],
]);
row.push(Value::Float(f64::from_bits(bits)));
}
ColumnTypeTag::Bool => {
let buf = &self.column_buffers[ci];
if self.null_flags.get(ci).and_then(|f| f.get(i)) == Some(&true) {
row.push(Value::Null);
continue;
}
row.push(Value::Bool(buf.get(i).copied().unwrap_or(0) != 0));
}
ColumnTypeTag::Text => {
let buf = &self.column_buffers[ci];
let mut p = 0usize;
let mut r = 0usize;
let mut found = None;
while p + 2 <= buf.len() {
let len = u16::from_le_bytes([buf[p], buf[p + 1]]) as usize;
p += 2;
if r == i {
if len == 0xFFFF {
found = Some(Value::Null);
} else if p + len <= buf.len() {
found = Some(Value::text(
String::from_utf8_lossy(&buf[p..p + len]).into_owned(),
));
} else {
found = Some(Value::Null);
}
break;
}
p += if len == 0xFFFF { 0 } else { len };
r += 1;
}
row.push(found.unwrap_or(Value::Null));
}
ColumnTypeTag::Vector => {
let buf = &self.column_buffers[ci];
let mut p = 0usize;
let mut r = 0usize;
let mut found = None;
while p + 2 <= buf.len() {
let dim = u16::from_le_bytes([buf[p], buf[p + 1]]) as usize;
p += 2;
if r == i {
if dim == 0 {
found = Some(Value::Null);
} else if p + dim * 4 <= buf.len() {
let mut v = Vec::with_capacity(dim);
for j in 0..dim {
let off = p + j * 4;
v.push(f32::from_le_bytes([
buf[off],
buf[off + 1],
buf[off + 2],
buf[off + 3],
]));
}
found = Some(Value::Vector(crate::types::ArcVec(
std::sync::Arc::new(v),
)));
} else {
found = Some(Value::Null);
}
break;
}
p += dim * 4;
r += 1;
}
row.push(found.unwrap_or(Value::Null));
}
ColumnTypeTag::Spatial => {
let buf = &self.column_buffers[ci];
let mut p = 0usize;
let mut r = 0usize;
let mut found = None;
while p + 2 <= buf.len() {
let len = u16::from_le_bytes([buf[p], buf[p + 1]]) as usize;
p += 2;
if r == i {
if len == 0 || p + len > buf.len() {
found = Some(Value::Null);
} else {
match bincode::deserialize::<crate::types::Geometry>(
&buf[p..p + len],
) {
Ok(g) => {
found = Some(Value::Spatial(std::boxed::Box::new(g)))
}
Err(_) => found = Some(Value::Null),
}
}
break;
}
p += len;
r += 1;
}
row.push(found.unwrap_or(Value::Null));
}
};
}
kept_rows.push((self.keys[i], self.timestamps[i], self.deleted[i], row));
}
self.keys.clear();
self.timestamps.clear();
self.deleted.clear();
for b in self.column_buffers.iter_mut() {
b.clear();
}
for f in self.null_flags.iter_mut() {
f.clear();
}
self.num_rows = 0;
for (key, ts, deleted, row) in kept_rows {
let _ = self.add_values(key, ts, deleted, &row);
}
}
pub fn finish_and_reset(&mut self) -> Result<()> {
if self.finished {
return Ok(());
}
if self.num_rows == 0 {
return Ok(());
}
if self.num_rows > 1 {
self.dedup_keys_newest_wins();
}
let num_rows = self.num_rows;
if num_rows == 0 {
return Ok(());
}
let num_cols = self.column_tags.len();
let mut segments: Vec<Vec<u8>> = Vec::with_capacity(num_cols);
for col_idx in 0..num_cols {
let tag = &self.column_tags[col_idx];
let raw = &self.column_buffers[col_idx];
let null_bytes = num_rows.div_ceil(8);
let mut seg = Vec::with_capacity(null_bytes + raw.len());
if tag.is_fixed() {
let mut nulls = vec![0u8; null_bytes];
let elem_size = tag.fixed_size();
let _ = elem_size;
let null_flags = &self.null_flags[col_idx];
for row_idx in 0..num_rows {
if row_idx < null_flags.len() && null_flags[row_idx] {
nulls[row_idx / 8] |= 1 << (row_idx % 8);
}
}
seg.extend_from_slice(&nulls);
seg.extend_from_slice(raw);
} else if matches!(tag, ColumnTypeTag::Text) {
let mut nulls = vec![0u8; null_bytes];
let mut offsets = Vec::with_capacity((num_rows + 1) * 4);
let mut str_data = Vec::new();
let mut current_offset = 0u32;
let null_flags = &self.null_flags[col_idx];
let mut pos = 0usize;
for row_idx in 0..num_rows {
if pos + 2 > raw.len() {
break;
}
let len = u16::from_le_bytes([raw[pos], raw[pos + 1]]) as usize;
pos += 2;
let is_null =
null_flags.get(row_idx).copied().unwrap_or(false) || len == 0xFFFF; if is_null {
nulls[row_idx / 8] |= 1 << (row_idx % 8);
offsets.push(current_offset);
if len != 0xFFFF {
pos += len;
}
continue;
}
offsets.push(current_offset);
if pos + len <= raw.len() {
str_data.extend_from_slice(&raw[pos..pos + len]);
current_offset += len as u32;
}
pos += len;
}
offsets.push(current_offset);
seg.extend_from_slice(&nulls);
for off in &offsets {
seg.extend_from_slice(&off.to_le_bytes());
}
seg.extend_from_slice(&str_data);
} else if matches!(tag, ColumnTypeTag::Vector) {
let mut nulls = vec![0u8; null_bytes];
let mut col_dim: usize = 0;
let mut row_dims: Vec<usize> = Vec::with_capacity(num_rows);
{
let mut pos = 0usize;
for row_idx in 0..num_rows {
if pos + 2 > raw.len() {
row_dims.push(0);
continue;
}
let d = u16::from_le_bytes([raw[pos], raw[pos + 1]]) as usize;
if d == 0 {
nulls[row_idx / 8] |= 1 << (row_idx % 8);
row_dims.push(0);
} else {
if d > col_dim {
col_dim = d;
}
row_dims.push(d);
}
pos += 2 + d * 4;
}
}
seg.extend_from_slice(&nulls);
seg.extend_from_slice(&(col_dim as u16).to_le_bytes());
let mut pos = 0usize;
for row_idx in 0..num_rows {
let d = row_dims[row_idx];
let mut vals = vec![0f32; col_dim];
if d > 0 && pos + 2 <= raw.len() {
let base = pos + 2;
for j in 0..d.min(col_dim) {
let off = base + j * 4;
if off + 4 <= raw.len() {
vals[j] = f32::from_le_bytes([
raw[off],
raw[off + 1],
raw[off + 2],
raw[off + 3],
]);
}
}
}
for v in &vals {
seg.extend_from_slice(&v.to_le_bytes());
}
pos += 2 + d * 4;
}
} else {
let nulls = vec![0u8; null_bytes];
seg.extend_from_slice(&nulls);
seg.extend_from_slice(raw);
}
segments.push(seg);
}
let (rm_size, _, _, _deleted_len) = RowMap::compute_sizes(num_rows);
let mut row_map = vec![0u8; rm_size];
for (i, k) in self.keys.iter().enumerate() {
let off = i * 8;
row_map[off..off + 8].copy_from_slice(&k.to_le_bytes());
}
let ts_off = num_rows * 8;
for (i, ts) in self.timestamps.iter().enumerate() {
let off = ts_off + i * 8;
row_map[off..off + 8].copy_from_slice(&ts.to_le_bytes());
}
let del_off = num_rows * 16;
for (i, d) in self.deleted.iter().enumerate() {
if *d {
row_map[del_off + i / 8] |= 1 << (i % 8);
}
}
let ci_offset = HEADER_SIZE as u64;
let ci_size = num_cols * COLUMN_INDEX_ENTRY_SIZE;
let segments_start = HEADER_SIZE + ci_size;
let mut compressed_segs: Vec<Vec<u8>> = Vec::with_capacity(num_cols);
let mut column_entries = Vec::with_capacity(num_cols);
let mut current_offset = segments_start as u64;
for seg in &segments {
let compressed = snap::raw::Encoder::new()
.compress_vec(seg)
.unwrap_or_else(|_| seg.clone());
let seg_data: Vec<u8> = if compressed.len() + 1 < seg.len() {
let mut out = Vec::with_capacity(1 + compressed.len());
out.push(1u8); out.extend_from_slice(&compressed);
out
} else {
let mut out = Vec::with_capacity(1 + seg.len());
out.push(0u8); out.extend_from_slice(seg);
out
};
let size = seg_data.len() as u64;
column_entries.push(ColumnIndexEntry {
offset: current_offset,
size,
});
current_offset += size;
compressed_segs.push(seg_data);
}
let row_map_offset = current_offset;
let total_size = row_map_offset as usize + row_map.len() + FOOTER_SIZE;
let mut buf = Vec::with_capacity(total_size);
let mut header_tags = [0u8; MAX_COLUMNS];
for (i, tag) in self.column_tags.iter().enumerate() {
header_tags[i] = *tag as u8;
}
let header = ColumnarHeader {
num_rows: num_rows as u32,
num_columns: num_cols as u16,
column_tags: header_tags,
};
buf.extend_from_slice(&header.serialize());
for entry in &column_entries {
buf.extend_from_slice(&entry.offset.to_le_bytes());
buf.extend_from_slice(&entry.size.to_le_bytes());
}
for seg in &compressed_segs {
buf.extend_from_slice(seg);
}
buf.extend_from_slice(&row_map);
let mut footer = [0u8; FOOTER_SIZE];
footer[0..8].copy_from_slice(&ci_offset.to_le_bytes());
footer[8..16].copy_from_slice(&row_map_offset.to_le_bytes());
footer[16..20].copy_from_slice(&COLUMNAR_MAGIC.to_le_bytes());
buf.extend_from_slice(&footer);
let final_path = self.path.clone();
let dir = final_path
.parent()
.unwrap_or_else(|| std::path::Path::new("."));
let tmp_path = dir.join(format!(
".{}.tmp",
final_path
.file_name()
.and_then(|n| n.to_str())
.unwrap_or("col.tmp")
));
let file = OpenOptions::new()
.write(true)
.create(true)
.truncate(true)
.open(&tmp_path)?;
let mut writer = BufWriter::new(file);
writer.write_all(&buf)?;
writer.flush()?;
writer.get_ref().sync_all()?;
drop(writer);
std::fs::rename(&tmp_path, &final_path)?;
self.finished = false;
self.num_rows = 0;
self.keys.clear();
self.timestamps.clear();
self.deleted.clear();
self.column_buffers = vec![Vec::new(); num_cols];
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::types::Value;
use tempfile::TempDir;
fn make_test_row(id: i64, name: &str, amount: f64, region: &str) -> Vec<Value> {
vec![
Value::Integer(id),
Value::Text(crate::types::ArcString(std::sync::Arc::from(name))),
Value::Float(amount),
Value::Text(crate::types::ArcString(std::sync::Arc::from(region))),
]
}
#[test]
#[cfg_attr(
target_os = "macos",
ignore = "macOS mmap coherence issue with files < page size"
)]
fn test_columnar_build_and_read() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("test.col.sst");
let col_types = vec![
ColumnType::Integer,
ColumnType::Text,
ColumnType::Float,
ColumnType::Text,
];
let mut builder = ColumnarSSTableBuilder::new(&path, col_types.clone());
let rows = vec![
(1u64, 100u64, false, make_test_row(1, "Alice", 99.5, "US")),
(2, 101, false, make_test_row(2, "Bob", 50.0, "EU")),
(3, 102, false, make_test_row(3, "Carol", 75.0, "US")),
(4, 103, true, make_test_row(4, "Dave", 0.0, "EU")), ];
for (key, ts, del, row) in &rows {
let encoded = crate::storage::row_format::encode(row, &col_types).unwrap();
builder.add_row(*key, *ts, *del, &encoded).unwrap();
}
builder.finish().unwrap();
assert!(ColumnarSSTable::is_columnar(&path));
let col_sst = ColumnarSSTable::open(&path).unwrap();
assert_eq!(col_sst.num_rows, 4);
assert_eq!(col_sst.column_tags.len(), 4);
assert_eq!(col_sst.row_map.key(0), 1);
assert_eq!(col_sst.row_map.key(3), 4);
assert!(!col_sst.row_map.is_deleted(0));
assert!(col_sst.row_map.is_deleted(3));
let id_seg = col_sst.read_fixed_i64(0).unwrap();
assert_eq!(id_seg.get_i64(0), Some(1));
assert_eq!(id_seg.get_i64(1), Some(2));
assert_eq!(id_seg.get_i64(2), Some(3));
assert_eq!(id_seg.get_i64(3), Some(4));
let amt_seg = col_sst.read_fixed_f64(2).unwrap();
assert_eq!(amt_seg.get_f64(0), Some(99.5));
assert_eq!(amt_seg.get_f64(1), Some(50.0));
let reg_seg = col_sst.read_text(3).unwrap();
assert_eq!(reg_seg.get_str(0), Some("US"));
assert_eq!(reg_seg.get_str(1), Some("EU"));
assert_eq!(reg_seg.get_str(2), Some("US"));
assert_eq!(reg_seg.get_str(3), Some("EU"));
}
#[test]
fn test_columnar_roundtrip_300k() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("large.col.sst");
let col_types = vec![
ColumnType::Integer,
ColumnType::Text,
ColumnType::Float,
ColumnType::Text,
];
let n = 10000;
let mut builder = ColumnarSSTableBuilder::new(&path, col_types.clone());
for i in 0..n {
let region = if i % 3 == 0 { "US" } else { "EU" };
let row = vec![
Value::Integer(i as i64),
Value::Text(crate::types::ArcString(std::sync::Arc::from(format!(
"cust_{}",
i % 100
)))),
Value::Float(i as f64 * 1.5),
Value::Text(crate::types::ArcString(std::sync::Arc::from(region))),
];
let encoded = crate::storage::row_format::encode(&row, &col_types).unwrap();
builder
.add_row(i as u64, i as u64 + 1000, i % 7 == 0, &encoded)
.unwrap();
}
builder.finish().unwrap();
let col_sst = ColumnarSSTable::open(&path).unwrap();
assert_eq!(col_sst.num_rows, n as usize);
let id_seg = col_sst.read_fixed_i64(0).unwrap();
let amt_seg = col_sst.read_fixed_f64(2).unwrap();
let reg_seg = col_sst.read_text(3).unwrap();
for i in 0..n {
assert_eq!(id_seg.get_i64(i), Some(i as i64), "id mismatch at {}", i);
let expected_amt = i as f64 * 1.5;
let got_amt = amt_seg.get_f64(i).unwrap();
assert!(
(got_amt - expected_amt).abs() < 0.001,
"amount mismatch at {}",
i
);
let expected_reg = if i % 3 == 0 { "US" } else { "EU" };
assert_eq!(
reg_seg.get_str(i),
Some(expected_reg),
"region mismatch at {}",
i
);
}
}
}