use crate::Result;
use crate::file_group::reader_v2::buffer::HoodieFileGroupRecordBuffer;
use crate::file_group::reader_v2::buffer::key_based::KeyBasedFileGroupRecordBuffer;
use crate::file_group::reader_v2::buffer::position_based::PositionBasedFileGroupRecordBuffer;
use crate::file_group::reader_v2::input_split::InputSplit;
use crate::file_group::reader_v2::merged_log_record_reader::HoodieMergedLogRecordReader;
use crate::file_group::reader_v2::metadata_merger::resolve_custom_merger;
use crate::file_group::reader_v2::read_stats::HoodieReadStats;
use crate::file_group::reader_v2::reader_context::ReaderContext;
use crate::file_group::reader_v2::reader_parameters::ReaderParameters;
use crate::storage::Storage;
use std::sync::Arc;
pub struct RecordBufferLoadResult {
pub record_buffer: Box<dyn HoodieFileGroupRecordBuffer>,
pub valid_block_instants: Vec<String>,
}
pub trait FileGroupRecordBufferLoader: Send + Sync + std::fmt::Debug {
fn get_record_buffer(
&self,
reader_context: Arc<ReaderContext>,
storage: Arc<Storage>,
input_split: &InputSplit,
reader_parameters: &ReaderParameters,
read_stats: &mut HoodieReadStats,
) -> impl std::future::Future<Output = Result<RecordBufferLoadResult>> + Send;
}
#[derive(Debug)]
pub struct DefaultFileGroupRecordBufferLoader;
impl DefaultFileGroupRecordBufferLoader {
pub fn new() -> Self {
Self
}
}
impl Default for DefaultFileGroupRecordBufferLoader {
fn default() -> Self {
Self::new()
}
}
impl FileGroupRecordBufferLoader for DefaultFileGroupRecordBufferLoader {
async fn get_record_buffer(
&self,
reader_context: Arc<ReaderContext>,
storage: Arc<Storage>,
input_split: &InputSplit,
reader_parameters: &ReaderParameters,
read_stats: &mut HoodieReadStats,
) -> Result<RecordBufferLoadResult> {
let merge_mode = if reader_context.merge_mode.is_empty() {
"COMMIT_TIME_ORDERING".to_string()
} else {
reader_context.merge_mode.to_uppercase()
};
match merge_mode.as_str() {
"COMMIT_TIME_ORDERING" | "EVENT_TIME_ORDERING" => {}
"CUSTOM" if resolve_custom_merger(&reader_context.table_config).is_some() => {}
unsupported => {
return Err(crate::error::CoreError::ReadFileSliceError(format!(
"Unsupported merge mode: '{unsupported}'. Only COMMIT_TIME_ORDERING, \
EVENT_TIME_ORDERING, and CUSTOM with a supported payload class are \
supported (MOR scan path)."
)));
}
}
log::debug!(
"[DefaultFileGroupRecordBufferLoader] getRecordBuffer: merge_mode={merge_mode} \
record_key_field={} ordering_fields={:?} \
log_files={} latest_commit_time={}",
reader_context.record_key_field(),
reader_context.ordering_field_names(),
input_split.log_file_paths.len(),
reader_context.latest_commit_time,
);
let is_skip_merge = reader_context
.hoodie_reader_config
.get(crate::file_group::reader_v2::reader_context::CONFIG_MERGE_TYPE)
.map(|v| v.eq_ignore_ascii_case("skip_merge"))
.unwrap_or(false);
let record_buffer: Box<dyn HoodieFileGroupRecordBuffer> = if is_skip_merge {
return Err(crate::error::CoreError::Unsupported(
"UnmergedFileGroupRecordBuffer (skip_merge mode) is not yet implemented"
.to_string(),
));
} else if reader_parameters.sort_output {
return Err(crate::error::CoreError::Unsupported(
"SortedKeyBasedFileGroupRecordBuffer (sort_output mode) is not yet implemented"
.to_string(),
));
} else if reader_parameters.use_record_position
&& input_split.base_file_path.is_some()
&& input_split.base_file_commit_time.is_some()
&& base_file_is_parquet(&reader_context.base_file_format)
{
let base_file_instant_time =
input_split.base_file_commit_time.clone().ok_or_else(|| {
crate::error::CoreError::ReadFileSliceError(
"internal: position-merge branch entered without a base-file commit time"
.to_string(),
)
})?;
Box::new(PositionBasedFileGroupRecordBuffer::new(
reader_context.clone(),
merge_mode,
reader_parameters.emit_delete,
base_file_instant_time,
)?)
} else if reader_parameters.emit_delete {
return Err(crate::error::CoreError::Unsupported(
"emit_delete=true (emitting delete records into the output) is not yet \
implemented; the supported read path drops deletes from the merged output"
.to_string(),
));
} else {
Box::new(KeyBasedFileGroupRecordBuffer::new(
reader_context.clone(),
merge_mode,
reader_parameters.emit_delete,
)?)
};
let (mut populated_buffer, valid_block_instants, stats) = scan_log_files(
reader_context,
storage,
input_split,
record_buffer,
reader_parameters,
)
.await?;
populated_buffer.compact_pinned_batches()?;
read_stats.total_log_read_time_us = stats.total_time_taken_to_read_and_merge_blocks_us;
read_stats.total_log_records = stats.total_log_records;
read_stats.total_log_blocks = stats.total_log_blocks;
read_stats.total_log_files_compacted = stats.total_log_files;
read_stats.total_corrupt_log_blocks = stats.total_corrupt_blocks;
read_stats.total_rollback_blocks = stats.total_rollbacks;
read_stats.log_block_read_us = stats.log_block_read_us;
read_stats.merge_insert_us = stats.merge_insert_us;
read_stats.log_block_fetch_us = stats.log_block_fetch_us;
read_stats.log_block_decode_us = stats.log_block_decode_us;
read_stats.merge_upsert_us = stats.merge_upsert_us;
read_stats.merge_map_peak_entries = stats.merge_map_peak_entries;
read_stats.merge_map_spilled = stats.merge_map_spilled;
read_stats.merge_map_peak_in_memory_bytes = stats.merge_map_peak_in_memory_bytes;
Ok(RecordBufferLoadResult {
record_buffer: populated_buffer,
valid_block_instants,
})
}
}
async fn scan_log_files(
reader_context: Arc<ReaderContext>,
storage: Arc<Storage>,
input_split: &InputSplit,
record_buffer: Box<dyn HoodieFileGroupRecordBuffer>,
reader_parameters: &ReaderParameters,
) -> Result<(
Box<dyn HoodieFileGroupRecordBuffer>,
Vec<String>,
crate::file_group::reader_v2::merged_log_record_reader::ScanStats,
)> {
if !input_split.has_log_files() {
let stats = crate::file_group::reader_v2::merged_log_record_reader::ScanStats::default();
return Ok((record_buffer, Vec::new(), stats));
}
let latest_instant_time = if reader_context.latest_commit_time.is_empty() {
crate::file_group::reader_v2::MAX_INSTANT_TIME.to_string()
} else {
reader_context.latest_commit_time.clone()
};
let completion_gate_inputs = reader_context.completion_gate_inputs.clone();
let reader = HoodieMergedLogRecordReader::new_builder()
.with_reader_context(reader_context)
.with_storage(storage)
.with_log_files(input_split.log_file_paths.clone())
.with_latest_instant_time(latest_instant_time)
.with_record_buffer(record_buffer)
.with_allow_inflight_instants(reader_parameters.allow_inflight_instants)
.with_completion_gate_inputs(completion_gate_inputs)
.with_force_full_scan(true)
.build()
.await?;
Ok(reader.into_parts())
}
pub(crate) fn base_file_is_parquet(base_file_format: &str) -> bool {
base_file_format.is_empty() || base_file_format.eq_ignore_ascii_case("parquet")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::HudiConfigs;
use crate::error::CoreError;
use crate::file_group::reader_v2::read_stats::HoodieReadStats;
use crate::file_group::reader_v2::reader_context::CONFIG_MERGE_TYPE;
use crate::storage::util::parse_uri;
fn context_with(merge_mode: &str) -> ReaderContext {
let mut ctx = ReaderContext::empty();
ctx.merge_mode = merge_mode.to_string();
ctx
}
fn split() -> InputSplit {
InputSplit::new(
Some("f1-0_0-1-1_20240101120000000.parquet".to_string()),
Some("20240101120000000".to_string()),
vec![".f1-0_20240101130000000.log.1_0-1-1".to_string()],
String::new(),
)
}
fn storage() -> Arc<Storage> {
Storage::new_with_base_url(parse_uri("file:///tmp").unwrap()).unwrap()
}
async fn load_err(ctx: ReaderContext, params: ReaderParameters, expected: &str) -> CoreError {
let mut stats = HoodieReadStats::default();
match DefaultFileGroupRecordBufferLoader::new()
.get_record_buffer(Arc::new(ctx), storage(), &split(), ¶ms, &mut stats)
.await
{
Ok(_) => panic!("{expected}"),
Err(e) => e,
}
}
#[tokio::test]
async fn test_get_record_buffer_unsupported_merge_mode_is_refused_by_name() {
let err = load_err(
context_with("CUSTOM"),
ReaderParameters::default(),
"CUSTOM must be refused",
)
.await;
assert!(
matches!(&err, CoreError::ReadFileSliceError(m) if m.contains("CUSTOM")),
"the error must name the mode it refused, got: {err}"
);
}
#[tokio::test]
async fn test_get_record_buffer_merge_mode_is_case_insensitive() {
let err = load_err(
context_with("event_time_ordering"),
ReaderParameters::default(),
"the fixture log file does not exist",
)
.await;
assert!(
!err.to_string().contains("Unsupported merge mode"),
"a lower-case mode must pass the gate, got: {err}"
);
}
#[tokio::test]
async fn test_get_record_buffer_empty_merge_mode_defaults_to_commit_time() {
let err = load_err(
context_with(""),
ReaderParameters::default(),
"the fixture log file does not exist",
)
.await;
assert!(
!err.to_string().contains("Unsupported merge mode"),
"an unset mode must default rather than be refused, got: {err}"
);
}
#[tokio::test]
async fn test_get_record_buffer_skip_merge_is_refused() {
let mut ctx = context_with("COMMIT_TIME_ORDERING");
ctx.hoodie_reader_config
.insert(CONFIG_MERGE_TYPE.to_string(), "skip_merge".to_string());
let err = load_err(
ctx,
ReaderParameters::default(),
"skip_merge must be refused",
)
.await;
assert!(
matches!(&err, CoreError::Unsupported(m) if m.contains("skip_merge")),
"got: {err}"
);
}
#[tokio::test]
async fn test_get_record_buffer_sort_output_is_refused() {
let params = ReaderParameters {
sort_output: true,
..Default::default()
};
let err = load_err(
context_with("COMMIT_TIME_ORDERING"),
params,
"sort_output must be refused",
)
.await;
assert!(
matches!(&err, CoreError::Unsupported(m) if m.contains("sort_output")),
"got: {err}"
);
}
#[tokio::test]
async fn test_get_record_buffer_emit_delete_is_refused() {
let params = ReaderParameters {
emit_delete: true,
..Default::default()
};
let err = load_err(
context_with("COMMIT_TIME_ORDERING"),
params,
"emit_delete must be refused",
)
.await;
assert!(
matches!(&err, CoreError::Unsupported(m) if m.contains("emit_delete")),
"got: {err}"
);
}
#[test]
fn test_base_file_is_parquet_accepts_unset_and_any_casing() {
assert!(base_file_is_parquet(""));
assert!(base_file_is_parquet("parquet"));
assert!(base_file_is_parquet("PARQUET"));
assert!(base_file_is_parquet("Parquet"));
assert!(!base_file_is_parquet("lance"));
assert!(!base_file_is_parquet("hfile"));
assert!(!base_file_is_parquet("orc"));
}
}