use crate::config::read::HudiReadConfig::ListingParallelism;
use crate::config::table::HudiTableConfig::BaseFileFormat;
use crate::config::HudiConfigs;
use crate::error::CoreError;
use crate::file_group::base_file::BaseFile;
use crate::file_group::log_file::LogFile;
use crate::file_group::FileGroup;
use crate::metadata::LAKE_FORMAT_METADATA_DIRS;
use crate::storage::{get_leaf_dirs, Storage};
use crate::table::partition::{
is_table_partitioned, PartitionPruner, EMPTY_PARTITION_PATH, PARTITION_METAFIELD_PREFIX,
};
use crate::Result;
use dashmap::DashMap;
use futures::{stream, StreamExt, TryStreamExt};
use std::collections::HashMap;
use std::string::ToString;
use std::sync::Arc;
#[derive(Clone, Debug)]
#[allow(dead_code)]
pub struct FileLister {
hudi_configs: Arc<HudiConfigs>,
storage: Arc<Storage>,
partition_pruner: PartitionPruner,
}
impl FileLister {
pub fn new(
hudi_configs: Arc<HudiConfigs>,
storage: Arc<Storage>,
partition_pruner: PartitionPruner,
) -> Self {
Self {
hudi_configs,
storage,
partition_pruner,
}
}
fn should_exclude_for_listing(file_name: &str) -> bool {
file_name.starts_with(PARTITION_METAFIELD_PREFIX) || file_name.ends_with(".crc")
}
async fn list_file_groups_for_partition(&self, partition_path: &str) -> Result<Vec<FileGroup>> {
let base_file_format = self
.hudi_configs
.get_or_default(BaseFileFormat)
.to::<String>();
let listed_file_metadata = self.storage.list_files(Some(partition_path)).await?;
let mut file_id_to_base_files: HashMap<String, Vec<BaseFile>> = HashMap::new();
let mut file_id_to_log_files: HashMap<String, Vec<LogFile>> = HashMap::new();
for file_metadata in listed_file_metadata {
if FileLister::should_exclude_for_listing(&file_metadata.name) {
continue;
}
let base_file_extension = format!(".{}", base_file_format);
if file_metadata.name.ends_with(&base_file_extension) {
let base_file = BaseFile::try_from(file_metadata)?;
let file_id = &base_file.file_id;
file_id_to_base_files
.entry(file_id.to_owned())
.or_default()
.push(base_file);
} else {
match LogFile::try_from(file_metadata) {
Ok(log_file) => {
let file_id = &log_file.file_id;
file_id_to_log_files
.entry(file_id.to_owned())
.or_default()
.push(log_file);
}
Err(e) => {
log::warn!("Failed to create a log file: {}", e);
continue;
}
}
}
}
let mut file_groups: Vec<FileGroup> = Vec::new();
for (file_id, base_files) in file_id_to_base_files.into_iter() {
let mut file_group = FileGroup::new(file_id.to_owned(), partition_path.to_string());
file_group.add_base_files(base_files)?;
let log_files = file_id_to_log_files.remove(&file_id).unwrap_or_default();
file_group.add_log_files(log_files)?;
file_groups.push(file_group);
}
Ok(file_groups)
}
async fn list_relevant_partition_paths(&self) -> Result<Vec<String>> {
if !is_table_partitioned(&self.hudi_configs) {
return Ok(vec![EMPTY_PARTITION_PATH.to_string()]);
}
let top_level_dirs: Vec<String> = self
.storage
.list_dirs(None)
.await?
.into_iter()
.filter(|dir| !LAKE_FORMAT_METADATA_DIRS.contains(&dir.as_str()))
.collect();
let mut partition_paths = Vec::new();
for dir in top_level_dirs {
partition_paths.extend(get_leaf_dirs(&self.storage, Some(&dir)).await?);
}
if partition_paths.is_empty() || self.partition_pruner.is_empty() {
return Ok(partition_paths);
}
Ok(partition_paths
.into_iter()
.filter(|path_str| self.partition_pruner.should_include(path_str))
.collect())
}
pub async fn list_file_groups_for_relevant_partitions(
&self,
) -> Result<DashMap<String, Vec<FileGroup>>> {
if !is_table_partitioned(&self.hudi_configs) {
let file_groups = self
.list_file_groups_for_partition(EMPTY_PARTITION_PATH)
.await?;
let file_groups_map = DashMap::with_capacity(1);
file_groups_map.insert(EMPTY_PARTITION_PATH.to_string(), file_groups);
return Ok(file_groups_map);
}
let pruned_partition_paths = self.list_relevant_partition_paths().await?;
let file_groups_map = Arc::new(DashMap::with_capacity(pruned_partition_paths.len()));
let parallelism = self
.hudi_configs
.get_or_default(ListingParallelism)
.to::<usize>();
stream::iter(pruned_partition_paths)
.map(|p| async move {
let file_groups = self.list_file_groups_for_partition(&p).await?;
Ok::<_, CoreError>((p, file_groups))
})
.buffer_unordered(parallelism)
.try_for_each(|(p, file_groups)| {
let file_groups_map = file_groups_map.clone();
async move {
file_groups_map.insert(p, file_groups);
Ok(())
}
})
.await?;
Ok(file_groups_map.as_ref().to_owned())
}
}
#[cfg(test)]
mod test {
use super::*;
use crate::table::Table;
use hudi_test::SampleTable;
use std::collections::HashSet;
#[tokio::test]
async fn list_partition_paths_for_nonpartitioned_table() {
let base_url = SampleTable::V6Nonpartitioned.url_to_cow();
let hudi_table = Table::new(base_url.path()).await.unwrap();
let lister = FileLister::new(
hudi_table.hudi_configs.clone(),
hudi_table.file_system_view.storage.clone(),
PartitionPruner::empty(),
);
let partition_paths = lister.list_relevant_partition_paths().await.unwrap();
let partition_path_set: HashSet<&str> =
HashSet::from_iter(partition_paths.iter().map(|p| p.as_str()));
assert_eq!(partition_path_set, HashSet::from([""]))
}
#[tokio::test]
async fn list_partition_paths_for_complexkeygen_table() {
let base_url = SampleTable::V6ComplexkeygenHivestyle.url_to_cow();
let hudi_table = Table::new(base_url.path()).await.unwrap();
let fs_view = &hudi_table.file_system_view;
let lister = FileLister::new(
fs_view.hudi_configs.clone(),
fs_view.storage.clone(),
PartitionPruner::empty(),
);
let partition_paths = lister.list_relevant_partition_paths().await.unwrap();
let partition_path_set: HashSet<&str> =
HashSet::from_iter(partition_paths.iter().map(|p| p.as_str()));
assert_eq!(
partition_path_set,
HashSet::from_iter(vec![
"byteField=10/shortField=300",
"byteField=20/shortField=100",
"byteField=30/shortField=100"
])
)
}
}