use std::collections::HashMap;
use std::collections::HashSet;
use std::sync::Arc;
use arrow_array::cast::AsArray;
use arrow_array::{Array, RecordBatch};
use crate::Result;
use crate::config::HudiConfigs;
use crate::error::CoreError;
use crate::file_group::base_file::reader::KeyPredicate;
use crate::file_group::file_slice::FileSlice;
use crate::file_group::reader_v2::engine::HoodieFileGroupReader;
use crate::file_group::reader_v2::input_split::InputSplit;
use crate::file_group::reader_v2::reader_parameters::ReaderParameters;
use crate::file_group::reader_v2::resolver::resolve_reader_context;
use crate::metadata::table::records::{
FilesPartitionRecord, HoodieMetadataFileInfo, MetadataRecordType,
};
use crate::storage::Storage;
const KEY_COLUMN: &str = "key";
const TYPE_COLUMN: &str = "type";
const FILESYSTEM_METADATA_COLUMN: &str = "filesystemMetadata";
const FILE_SIZE_FIELD: &str = "size";
const FILE_IS_DELETED_FIELD: &str = "isDeleted";
pub(crate) struct MetadataTableV2Reader {
hudi_configs: Arc<HudiConfigs>,
storage: Arc<Storage>,
valid_instants: Option<HashSet<String>>,
}
impl MetadataTableV2Reader {
pub(crate) fn new(hudi_configs: Arc<HudiConfigs>, storage: Arc<Storage>) -> Self {
Self {
hudi_configs,
storage,
valid_instants: None,
}
}
pub(crate) fn with_valid_instants(mut self, instants: HashSet<String>) -> Self {
self.valid_instants = Some(instants);
self
}
pub(crate) async fn read_files_partition(
&self,
file_slice: &FileSlice,
keys: &[&str],
) -> Result<HashMap<String, FilesPartitionRecord>> {
let batch = self.read_files_partition_batch(file_slice, keys).await?;
Self::records_from_batch(&batch, keys)
}
pub(crate) async fn read_files_partition_batch(
&self,
file_slice: &FileSlice,
keys: &[&str],
) -> Result<RecordBatch> {
let base_file_path = file_slice.base_file_relative_path()?;
let log_file_paths = if file_slice.has_log_file() {
file_slice
.log_files
.iter()
.map(|log_file| file_slice.log_file_relative_path(log_file))
.collect::<Result<Vec<String>>>()?
} else {
vec![]
};
let mut reader_context = resolve_reader_context(
&self.hudi_configs,
!log_file_paths.is_empty(),
base_file_path.as_deref(),
)?;
reader_context.rebuild_record_context(FilesPartitionRecord::PARTITION_NAME.to_string());
let key_predicate = (!keys.is_empty())
.then(|| KeyPredicate::Keys(keys.iter().map(|k| (*k).to_string()).collect()));
reader_context.key_predicate = key_predicate;
if let Some(valid) = &self.valid_instants {
reader_context.instant_range = Some(
crate::timeline::selector::InstantRange::exact_match(valid.iter().cloned(), "UTC"),
);
}
let base_file_commit_time = file_slice
.base_file
.as_ref()
.map(|base_file| base_file.commit_timestamp.clone());
let mut reader = HoodieFileGroupReader::new(
Arc::new(reader_context),
self.storage.clone(),
InputSplit::new(
base_file_path,
base_file_commit_time,
log_file_paths,
FilesPartitionRecord::PARTITION_NAME.to_string(),
),
ReaderParameters::default(),
None,
None,
)?;
let batch = reader.read().await?;
log::debug!(
"metadata read of '{}' with {} named key(s): merge map peaked at {} entries",
FilesPartitionRecord::PARTITION_NAME,
keys.len(),
reader.read_stats().merge_map_peak_entries
);
Ok(batch)
}
fn records_from_batch(
batch: &RecordBatch,
keys: &[&str],
) -> Result<HashMap<String, FilesPartitionRecord>> {
let missing = |name: &str| {
CoreError::MetadataTable(format!(
"A metadata record must carry a '{name}' column; the metadata table's \
schema has changed or the wrong partition was read"
))
};
let key_column = batch
.column_by_name(KEY_COLUMN)
.ok_or_else(|| missing(KEY_COLUMN))?
.as_string_opt::<i32>()
.ok_or_else(|| missing(KEY_COLUMN))?;
let types = batch
.column_by_name(TYPE_COLUMN)
.ok_or_else(|| missing(TYPE_COLUMN))?
.as_primitive_opt::<arrow_array::types::Int32Type>()
.ok_or_else(|| missing(TYPE_COLUMN))?;
let maps = batch
.column_by_name(FILESYSTEM_METADATA_COLUMN)
.ok_or_else(|| missing(FILESYSTEM_METADATA_COLUMN))?
.as_map_opt()
.ok_or_else(|| missing(FILESYSTEM_METADATA_COLUMN))?;
let wanted: Option<std::collections::HashSet<&str>> =
(!keys.is_empty()).then(|| keys.iter().copied().collect());
let mut out = HashMap::with_capacity(batch.num_rows());
for row in 0..batch.num_rows() {
if key_column.is_null(row) {
continue;
}
if let Some(wanted) = wanted.as_ref()
&& !wanted.contains(key_column.value(row))
{
continue;
}
let record_type = MetadataRecordType::from(types.value(row));
let key = normalise_partition(key_column.value(row));
let mut files = Self::file_map(maps, row)?;
if record_type == MetadataRecordType::AllPartitions
&& let Some(mut info) = files.remove(FilesPartitionRecord::NON_PARTITIONED_NAME)
{
info.name = String::new();
files.insert(String::new(), info);
}
out.insert(
key.clone(),
FilesPartitionRecord {
key,
record_type,
files,
},
);
}
Ok(out)
}
fn file_map(
maps: &arrow_array::MapArray,
row: usize,
) -> Result<HashMap<String, HoodieMetadataFileInfo>> {
let mut files = HashMap::new();
if maps.is_null(row) {
return Ok(files);
}
let entries = maps.value(row);
let names = entries.column(0).as_string_opt::<i32>().ok_or_else(|| {
CoreError::MetadataTable(
"'filesystemMetadata' keys must be Utf8 file names".to_string(),
)
})?;
let values = entries.column(1).as_struct_opt().ok_or_else(|| {
CoreError::MetadataTable("'filesystemMetadata' values must be structs".to_string())
})?;
let sizes = values
.column_by_name(FILE_SIZE_FIELD)
.and_then(|c| c.as_primitive_opt::<arrow_array::types::Int64Type>())
.ok_or_else(|| {
CoreError::MetadataTable(format!("'{FILE_SIZE_FIELD}' must be Int64"))
})?;
let deleted = values
.column_by_name(FILE_IS_DELETED_FIELD)
.and_then(|c| c.as_boolean_opt())
.ok_or_else(|| {
CoreError::MetadataTable(format!("'{FILE_IS_DELETED_FIELD}' must be Boolean"))
})?;
for i in 0..names.len() {
if names.is_null(i) || values.is_null(i) {
continue;
}
let name = names.value(i).to_string();
files.insert(
name.clone(),
HoodieMetadataFileInfo {
name,
size: if sizes.is_null(i) { 0 } else { sizes.value(i) },
is_deleted: !deleted.is_null(i) && deleted.value(i),
},
);
}
Ok(files)
}
}
fn normalise_partition(raw: &str) -> String {
if raw == FilesPartitionRecord::NON_PARTITIONED_NAME {
String::new()
} else {
raw.to_string()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::file_group::FileGroup;
use crate::metadata::table::reader::MetadataTableFileGroupReader;
use hudi_test::QuickstartTripsTable;
use std::fs::canonicalize;
use std::path::PathBuf;
use url::Url;
const FILE_GROUP: &str = "files-0000-0";
const PRECOMPACT_BASE: &str = "files-0000-0_0-955-2690_00000000000000000.hfile";
const PRECOMPACT_LOGS: &[&str] = &[
".files-0000-0_00000000000000000.log.1_0-0-0",
".files-0000-0_20251220210108078.log.1_10-999-2838",
".files-0000-0_20251220210123755.log.1_3-1032-2950",
".files-0000-0_20251220210125441.log.1_5-1057-3024",
".files-0000-0_20251220210127080.log.1_3-1082-3100",
".files-0000-0_20251220210128625.log.1_5-1107-3174",
".files-0000-0_20251220210129235.log.1_3-1118-3220",
".files-0000-0_20251220210130911.log.1_3-1149-3338",
];
fn metadata_configs() -> Arc<HudiConfigs> {
let table_path = QuickstartTripsTable::V8Trips8I3U1D.path_to_mor_avro();
let mdt = PathBuf::from(table_path).join(".hoodie").join("metadata");
let uri = Url::from_file_path(canonicalize(&mdt).unwrap())
.unwrap()
.as_ref()
.to_string();
crate::metadata::table::test_support::metadata_table_configs(&uri)
}
fn precompact_slice() -> crate::Result<FileSlice> {
let mut fg = FileGroup::new(
FILE_GROUP.to_string(),
FilesPartitionRecord::PARTITION_NAME.to_string(),
);
fg.add_base_file_from_name(PRECOMPACT_BASE)?;
fg.add_log_files_from_names(PRECOMPACT_LOGS.iter().copied())?;
Ok(fg
.get_file_slice_as_of(crate::file_group::reader_v2::MAX_INSTANT_TIME)
.expect("the file group has a slice")
.clone())
}
use crate::metadata::table::test_support::comparable;
fn base_only_slice() -> crate::Result<FileSlice> {
let mut fg = FileGroup::new(
FILE_GROUP.to_string(),
FilesPartitionRecord::PARTITION_NAME.to_string(),
);
fg.add_base_file_from_name(PRECOMPACT_BASE)?;
Ok(fg
.get_file_slice_as_of(crate::file_group::reader_v2::MAX_INSTANT_TIME)
.expect("the file group has a slice")
.clone())
}
#[tokio::test]
#[ignore]
async fn reader_cost_on_a_metadata_slice() -> crate::Result<()> {
use std::time::Instant;
let configs = metadata_configs();
let storage = Storage::new(Arc::new(HashMap::new()), configs.clone())?;
let v1 = MetadataTableFileGroupReader::new(configs.clone(), storage.clone());
let v2 = MetadataTableV2Reader::new(configs.clone(), storage.clone());
let full = precompact_slice()?;
let base_only = base_only_slice()?;
const WARM: usize = 20;
const ITERS: usize = 60;
println!("{:28} {:>10} {:>10} {:>8}", "slice", "v1", "v2", "ratio");
for (label, slice, keys) in [
("base + 8 logs, all keys", &full, &[][..]),
("base + 8 logs, one key", &full, &["city=chennai"][..]),
("base only, all keys", &base_only, &[][..]),
] {
let a = v1.read_files_partition(slice, keys).await?;
let b = v2.read_files_partition(slice, keys).await?;
assert_eq!(
comparable(&a),
comparable(&b),
"{label}: the readers must agree before a timing means anything"
);
for _ in 0..WARM {
v1.read_files_partition(slice, keys).await?;
v2.read_files_partition(slice, keys).await?;
}
let t = Instant::now();
for _ in 0..ITERS {
v1.read_files_partition(slice, keys).await?;
}
let d1 = t.elapsed() / ITERS as u32;
let t = Instant::now();
for _ in 0..ITERS {
v2.read_files_partition(slice, keys).await?;
}
let d2 = t.elapsed() / ITERS as u32;
println!(
"{label:28} {:>8}us {:>8}us {:>7.2}x",
d1.as_micros(),
d2.as_micros(),
d2.as_secs_f64() / d1.as_secs_f64()
);
}
Ok(())
}
#[tokio::test]
async fn it_matches_the_existing_reader_value_for_value() -> crate::Result<()> {
let configs = metadata_configs();
let storage = Storage::new(Arc::new(HashMap::new()), configs.clone())?;
let slice = precompact_slice()?;
let existing = MetadataTableFileGroupReader::new(configs.clone(), storage.clone())
.read_files_partition(&slice, &[])
.await?;
assert!(
comparable(&existing)
.iter()
.any(|(_, _, files)| files.len() > 1),
"the oracle must fold several file entries into at least one record, or \
the comparison says nothing about the merge"
);
let through_v2 = MetadataTableV2Reader::new(configs, storage)
.read_files_partition(&slice, &[])
.await?;
assert_eq!(comparable(&through_v2), comparable(&existing));
Ok(())
}
#[tokio::test]
async fn a_named_keys_read_matches_the_existing_reader() -> crate::Result<()> {
let configs = metadata_configs();
let storage = Storage::new(Arc::new(HashMap::new()), configs.clone())?;
let slice = precompact_slice()?;
let all = MetadataTableFileGroupReader::new(configs.clone(), storage.clone())
.read_files_partition(&slice, &[])
.await?;
let mut names: Vec<&str> = all.keys().map(String::as_str).collect();
names.sort();
assert!(names.len() > 1, "the fixture must hold several keys");
let wanted = vec![names[names.len() - 1]];
let existing = MetadataTableFileGroupReader::new(configs.clone(), storage.clone())
.read_files_partition(&slice, &wanted)
.await?;
let through_v2 = MetadataTableV2Reader::new(configs, storage)
.read_files_partition(&slice, &wanted)
.await?;
assert_eq!(
existing.len(),
1,
"the oracle must return just the named key"
);
assert_eq!(comparable(&through_v2), comparable(&existing));
Ok(())
}
fn one_record(key: &str, record_type: i32, entries: &[(&str, i64, bool)]) -> RecordBatch {
use arrow_array::builder::{
BooleanBuilder, Int64Builder, MapBuilder, StringBuilder, StructBuilder,
};
use arrow_schema::{DataType, Field, Fields, Schema};
let value_fields = Fields::from(vec![
Field::new(FILE_SIZE_FIELD, DataType::Int64, false),
Field::new(FILE_IS_DELETED_FIELD, DataType::Boolean, false),
]);
let mut builder = MapBuilder::new(
None,
StringBuilder::new(),
StructBuilder::new(
value_fields,
vec![
Box::new(Int64Builder::new()),
Box::new(BooleanBuilder::new()),
],
),
);
for (name, size, is_deleted) in entries {
builder.keys().append_value(name);
let values = builder.values();
values
.field_builder::<Int64Builder>(0)
.unwrap()
.append_value(*size);
values
.field_builder::<BooleanBuilder>(1)
.unwrap()
.append_value(*is_deleted);
values.append(true);
}
builder.append(true).unwrap();
let map = builder.finish();
let schema = Arc::new(Schema::new(vec![
Field::new(KEY_COLUMN, DataType::Utf8, false),
Field::new(TYPE_COLUMN, DataType::Int32, false),
Field::new(FILESYSTEM_METADATA_COLUMN, map.data_type().clone(), false),
]));
RecordBatch::try_new(
schema,
vec![
Arc::new(arrow_array::StringArray::from(vec![key])),
Arc::new(arrow_array::Int32Array::from(vec![record_type])),
Arc::new(map),
],
)
.unwrap()
}
#[test]
fn a_non_partitioned_key_normalises_to_the_empty_string() {
let batch = one_record(
".",
MetadataRecordType::Files as i32,
&[("a.parquet", 10, false)],
);
let records = MetadataTableV2Reader::records_from_batch(&batch, &[]).unwrap();
assert_eq!(records.keys().collect::<Vec<_>>(), vec![""]);
assert_eq!(records[""].key, "");
}
#[test]
fn a_non_partitioned_entry_inside_all_partitions_normalises_too() {
let batch = one_record(
"__all_partitions__",
MetadataRecordType::AllPartitions as i32,
&[(".", 0, false)],
);
let records = MetadataTableV2Reader::records_from_batch(&batch, &[]).unwrap();
let record = &records["__all_partitions__"];
assert_eq!(record.files.keys().collect::<Vec<_>>(), vec![""]);
assert_eq!(
record.files[""].name, "",
"the entry's own name field has to be cleared as well as its map key"
);
}
#[test]
fn a_null_map_value_is_omitted_rather_than_defaulted() {
use arrow_array::builder::{
BooleanBuilder, Int64Builder, MapBuilder, StringBuilder, StructBuilder,
};
use arrow_schema::{DataType, Field, Fields, Schema};
let value_fields = Fields::from(vec![
Field::new(FILE_SIZE_FIELD, DataType::Int64, true),
Field::new(FILE_IS_DELETED_FIELD, DataType::Boolean, true),
]);
let mut builder = MapBuilder::new(
None,
StringBuilder::new(),
StructBuilder::new(
value_fields,
vec![
Box::new(Int64Builder::new()),
Box::new(BooleanBuilder::new()),
],
),
);
for (name, present) in [("live.parquet", true), ("null-valued.parquet", false)] {
builder.keys().append_value(name);
let values = builder.values();
let size = values.field_builder::<Int64Builder>(0).unwrap();
if present {
size.append_value(10);
} else {
size.append_null();
}
let deleted = values.field_builder::<BooleanBuilder>(1).unwrap();
if present {
deleted.append_value(false);
} else {
deleted.append_null();
}
values.append(present);
}
builder.append(true).unwrap();
let map = builder.finish();
let schema = Arc::new(Schema::new(vec![
Field::new(KEY_COLUMN, DataType::Utf8, false),
Field::new(TYPE_COLUMN, DataType::Int32, false),
Field::new(FILESYSTEM_METADATA_COLUMN, map.data_type().clone(), false),
]));
let batch = RecordBatch::try_new(
schema,
vec![
Arc::new(arrow_array::StringArray::from(vec!["city=x"])),
Arc::new(arrow_array::Int32Array::from(vec![
MetadataRecordType::Files as i32,
])),
Arc::new(map),
],
)
.unwrap();
let records = MetadataTableV2Reader::records_from_batch(&batch, &[]).unwrap();
let files = &records["city=x"].files;
assert_eq!(
files.keys().collect::<Vec<_>>(),
vec!["live.parquet"],
"a null-valued entry must be dropped, not turned into a phantom active file"
);
}
#[test]
fn a_tombstone_is_reported_as_deleted() {
let batch = one_record(
"city=x",
MetadataRecordType::Files as i32,
&[("live.parquet", 10, false), ("gone.parquet", 0, true)],
);
let records = MetadataTableV2Reader::records_from_batch(&batch, &[]).unwrap();
let files = &records["city=x"].files;
assert!(!files["live.parquet"].is_deleted);
assert!(files["gone.parquet"].is_deleted);
assert_eq!(files["live.parquet"].size, 10);
}
#[test]
fn an_unknown_record_type_decodes_to_unknown_as_the_existing_reader_does() {
let batch = one_record("k", 99, &[]);
let records = MetadataTableV2Reader::records_from_batch(&batch, &[]).unwrap();
assert_eq!(records["k"].record_type, MetadataRecordType::Unknown);
}
async fn merge_peak(
configs: Arc<HudiConfigs>,
storage: Arc<Storage>,
slice: &FileSlice,
keys: &[&str],
) -> crate::Result<u64> {
let base_file_path = slice.base_file_relative_path()?;
let mut context =
resolve_reader_context(&configs, slice.has_log_file(), base_file_path.as_deref())?;
context.rebuild_record_context(FilesPartitionRecord::PARTITION_NAME.to_string());
context.key_predicate = (!keys.is_empty())
.then(|| KeyPredicate::Keys(keys.iter().map(|k| (*k).to_string()).collect()));
let logs = slice
.log_files
.iter()
.map(|f| slice.log_file_relative_path(f))
.collect::<crate::Result<Vec<String>>>()?;
let mut reader = HoodieFileGroupReader::new(
Arc::new(context),
storage,
InputSplit::new(
slice.base_file_relative_path()?,
slice.base_file.as_ref().map(|b| b.commit_timestamp.clone()),
logs,
FilesPartitionRecord::PARTITION_NAME.to_string(),
),
ReaderParameters::default(),
None,
None,
)?;
reader.read().await?;
Ok(reader.read_stats().merge_map_peak_entries as u64)
}
#[tokio::test]
async fn a_named_keys_read_merges_less_than_a_full_read() -> crate::Result<()> {
let configs = metadata_configs();
let storage = Storage::new(Arc::new(HashMap::new()), configs.clone())?;
let slice = precompact_slice()?;
let full = merge_peak(configs.clone(), storage.clone(), &slice, &[]).await?;
assert!(
full > 1,
"the slice must put several records through the merge, or there is no work \
to avoid; peaked at {full}"
);
let all = MetadataTableFileGroupReader::new(configs.clone(), storage.clone())
.read_files_partition(&slice, &[])
.await?;
let mut names: Vec<&str> = all.keys().map(String::as_str).collect();
names.sort();
let one = merge_peak(configs, storage, &slice, &[names[names.len() - 1]]).await?;
assert!(
one < full,
"a one-key read must merge fewer records than a full read; both peaked at \
{one} against {full}, so the predicate is not reaching the log read"
);
Ok(())
}
}