#![allow(dead_code, unused_imports)]
pub mod key_based;
pub mod loader;
pub mod position_based;
pub mod record_buffer;
pub mod record_positions;
pub mod row_extraction;
pub mod spillable_map;
pub use record_buffer::FileGroupRecordBuffer;
pub use spillable_map::{MergeMap, SpillableRecordMap};
use crate::Result;
use crate::file_group::log_file::log_block::LogBlock;
use crate::file_group::reader_v2::buffered_record::{BufferedRecord, DeleteRecord};
use crate::file_group::reader_v2::update_processor::UpdateStats;
use arrow_array::RecordBatch;
use arrow_schema::SchemaRef;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BufferType {
KeyBasedMerge,
PositionBasedMerge,
Unmerged,
}
pub trait HoodieFileGroupRecordBuffer: Send + std::fmt::Debug {
fn get_buffer_type(&self) -> BufferType;
fn process_data_block(&mut self, block: &mut LogBlock) -> Result<()>;
fn process_next_data_record(&mut self, record: BufferedRecord, key: &str) -> Result<()>;
fn process_delete_block(&mut self, block: &mut LogBlock) -> Result<()>;
fn process_next_deleted_record(&mut self, delete_record: DeleteRecord, key: &str)
-> Result<()>;
fn contains_log_record(&self, record_key: &str) -> bool;
fn size(&self) -> usize;
fn get_total_log_records(&self) -> u64;
fn merge_map_peak_entries(&self) -> u64 {
0
}
fn merge_map_spilled(&self) -> bool {
false
}
fn merge_map_peak_in_memory_bytes(&self) -> u64 {
0
}
fn current_in_memory_bytes(&self) -> u64 {
0
}
fn update_stats_snapshot(&self) -> UpdateStats {
UpdateStats::default()
}
fn get_log_records(&self) -> &MergeMap;
fn set_reader_schema(&mut self, schema: SchemaRef);
fn set_base_file_source(&mut self, source: Box<dyn arrow_array::RecordBatchReader + Send>);
fn set_base_file_iterator(&mut self, batches: Vec<RecordBatch>) {
let schema = batches
.first()
.map(|b| b.schema())
.unwrap_or_else(|| std::sync::Arc::new(arrow_schema::Schema::empty()));
let iter = arrow_array::RecordBatchIterator::new(batches.into_iter().map(Ok), schema);
self.set_base_file_source(Box::new(iter));
}
fn compact_pinned_batches(&mut self) -> Result<()> {
Ok(())
}
fn has_next(&mut self) -> Result<bool>;
fn next(&mut self) -> Option<BufferedRecord>;
fn merge_and_collect_with_stats(self: Box<Self>) -> Result<(RecordBatch, UpdateStats)>;
fn merge_and_collect(self: Box<Self>) -> Result<RecordBatch> {
self.merge_and_collect_with_stats().map(|(batch, _)| batch)
}
fn next_merged_base_batch(&mut self, target_schema: &SchemaRef) -> Result<Option<RecordBatch>>;
fn merge_base_batch(
&mut self,
base: &RecordBatch,
target_schema: &SchemaRef,
) -> Result<Option<RecordBatch>>;
fn drain_log_only_inserts(&mut self, target_schema: &SchemaRef) -> Result<Option<RecordBatch>>;
}