diff options
| author | Jan Tuomi <jan@jantuomi.fi> | 2025-01-14 09:22:26 +0200 |
|---|---|---|
| committer | Jan Tuomi <jan@jantuomi.fi> | 2025-01-14 09:47:20 +0200 |
| commit | bcb70560318a76e4a78e0c9e9118684841203c13 (patch) | |
| tree | 79091a396b4477ec88c82a189781a6df267576f5 /log_db/src | |
| parent | 3182bcbabd610579b29d33eacb1839800a0cd545 (diff) | |
Add batch_find_by_records to prepare for batch operation support
Diffstat (limited to 'log_db/src')
| -rw-r--r-- | log_db/src/lib.rs | 188 |
1 files changed, 144 insertions, 44 deletions
diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index dabc4a9..924d5b9 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -428,78 +428,173 @@ impl<R: Recordable> DB<R> { /// Get a record by its primary index value. /// E.g. `db.get(Value::Int(10))`. pub fn get(&mut self, value: &Value) -> Result<Option<R>, DBError> { - let records = self.find_by_records(&self.config.primary_key.clone(), value)?; + let value_batch = std::iter::once(value); + let records = self.batch_find_by_records(&self.config.primary_key.clone(), value_batch)?; assert!(records.len() <= 1); Ok(records .into_iter() .next() - .map(|rec| R::from_record(rec.values))) + .map(|(_, rec)| R::from_record(rec.values))) } /// Get a collection of records based on a field value. /// Indexes will be used if they are applicable. pub fn find_by(&mut self, field: &R::Field, value: &Value) -> Result<Vec<R>, DBError> { + let value_batch = std::iter::once(value); Ok(self - .find_by_records(field, value)? + .batch_find_by_records(field, value_batch)? .into_iter() - .map(|rec| R::from_record(rec.values)) + .map(|(_, rec)| R::from_record(rec.values)) .collect()) } - fn find_by_records(&mut self, field: &R::Field, value: &Value) -> Result<Vec<Record>, DBError> { + fn batch_find_by_records<'a>( + &mut self, + field: &R::Field, + values: impl Iterator<Item = &'a Value>, + ) -> Result<Vec<(usize, Record)>, DBError> { let field_type = self.get_field_type(field).ok_or(DBError::ValidationError( "Field not found in schema".to_owned(), ))?; - if !type_check(&value, &field_type) { - return Err(DBError::ValidationError(format!( - "Queried value does not match key type: {:?}", - field_type - ))); - } + let indexables = values + .map(|value| { + if type_check(&value, &field_type) { + value.as_indexable().ok_or(DBError::ValidationError( + "Queried value must be indexable".to_owned(), + )) + } else { + Err(DBError::ValidationError(format!( + "Queried value {:?} does not match key type: {:?}", + value, field_type + ))) + } + }) + .collect::<Result<Vec<IndexableValue>, DBError>>()?; // Otherwise, continue with querying secondary indexes. - debug!("Finding all records with field {:?} = {:?}", field, value); - let query_key = value.as_indexable().ok_or(DBError::ValidationError( - "Queried value must be indexable".to_owned(), - ))?; + debug!( + "Finding all records with fields {:?} = {:?}", + field, indexables + ); if self.config.read_consistency == ReadConsistency::Strong { self.refresh_indexes()?; } - let log_keys = if field == &self.config.primary_key { - let opt = self.primary_memtable.get(&query_key); - match opt { - Some(log_key) => vec![log_key.clone()], - None => vec![], + let log_key_batches = indexables + .into_iter() + .map(|query_key| { + if field == &self.config.primary_key { + let opt = self.primary_memtable.get(&query_key); + let log_keys = match opt { + Some(log_key) => vec![log_key.clone()], + None => vec![], + }; + Ok(log_keys) + } else { + let smemtable_index = match get_secondary_memtable_index_by_field( + &self.config.secondary_keys, + field, + ) { + Some(index) => index, + None => { + return Err(DBError::ValidationError( + "Cannot find_by by non-indexed key".to_owned(), + )) + } + }; + + let log_keys = self.secondary_memtables[smemtable_index] + .find_by(&query_key) + .into_iter() + .cloned() + .collect(); + Ok(log_keys) + } + }) + .collect::<Result<Vec<Vec<LogKey>>, DBError>>()?; + + debug!("Found log keys in memtable: {:?}", log_key_batches); + + let tagged = log_key_batches + .into_iter() + .enumerate() + .flat_map(|(tag, log_keys)| log_keys.into_iter().map(move |log_key| (tag, log_key))); + + let tagged_records = self.read_tagged_log_keys(tagged)?; + + debug!("Read {} records", tagged_records.len()); + + Ok(tagged_records) + } + + /// Read records from segment files based on log keys. + /// The log keys are accompanied by an integer tag that can be used to identify and group them later. + fn read_tagged_log_keys( + &mut self, + log_keys: impl Iterator<Item = (usize, LogKey)>, + ) -> Result<Vec<(usize, Record)>, DBError> { + let mut records = vec![]; + let mut log_keys_map = BTreeMap::new(); + + for (tag, log_key) in log_keys { + if !log_keys_map.contains_key(&log_key.segment_num()) { + log_keys_map.insert(log_key.segment_num(), vec![(tag, log_key.index())]); + } else { + log_keys_map + .get_mut(&log_key.segment_num()) + .unwrap() + .push((tag, log_key.index())); } - } else { - let smemtable_index = - match get_secondary_memtable_index_by_field(&self.config.secondary_keys, field) { - Some(index) => index, - None => { - return Err(DBError::ValidationError( - "Cannot find_by by non-indexed key".to_owned(), - )) - } - }; + } - self.secondary_memtables[smemtable_index] - .find_by(&query_key) - .into_iter() - .cloned() - .collect() - }; + for (segment_num, mut segment_indexes) in log_keys_map { + segment_indexes.sort_unstable(); - debug!("Found log keys in memtable: {:?}", log_keys); + let metadata_path = &self.data_dir.join(metadata_filename(segment_num)); + let mut metadata_file = READ_MODE.open(&metadata_path)?; - let ret = self.read_log_keys(log_keys.into_iter())?; + request_shared_lock(&self.data_dir, &mut metadata_file)?; + + let metadata_header = read_metadata_header(&mut metadata_file)?; + + let data_path = &self.data_dir.join(metadata_header.uuid.to_string()); + let mut data_file = READ_MODE.open(&data_path)?; + + data_file.lock_shared()?; + + let header_size = METADATA_FILE_HEADER_SIZE as i64; + let row_length = METADATA_ROW_LENGTH as i64; + let mut current_metadata_offset = header_size; + for (tag, segment_index) in segment_indexes { + let new_metadata_offset = header_size + segment_index as i64 * row_length; + metadata_file.seek_relative(new_metadata_offset - current_metadata_offset)?; + + let mut metadata_buf = [0; METADATA_ROW_LENGTH]; + metadata_file.read_exact(&mut metadata_buf)?; + + let data_offset = u64::from_be_bytes(metadata_buf[0..8].try_into().unwrap()); + let data_length = u64::from_be_bytes(metadata_buf[8..16].try_into().unwrap()); + assert!(data_length > 0); + + data_file.seek(SeekFrom::Start(data_offset))?; + + let mut data_buf = vec![0; data_length as usize]; + data_file.read_exact(&mut data_buf)?; - debug!("Read {} records", ret.len()); + let record = Record::deserialize(&data_buf); + records.push((tag, record)); - Ok(ret) + current_metadata_offset = new_metadata_offset + row_length; + } + + metadata_file.unlock()?; + data_file.unlock()?; + } + + Ok(records) } // log_keys is an iterator of LogKeys @@ -691,10 +786,15 @@ impl<R: Recordable> DB<R> { } fn delete_by_field(&mut self, field: &R::Field, value: &Value) -> Result<Vec<Record>, DBError> { - let mut recs = self.find_by_records(field, value)?; - for rec in &mut recs { - rec.tombstone = true; - } + let value_batch = std::iter::once(value); + let recs: Vec<Record> = self + .batch_find_by_records(field, value_batch)? + .into_iter() + .map(|(_, mut rec)| { + rec.tombstone = true; + rec + }) + .collect(); request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?; self.active_data_file.lock_exclusive()?; |
