pub(crate) mod hudi_exec;
pub(crate) mod util;
use std::collections::HashMap;
use std::error::Error;
use std::fmt::Debug;
use std::sync::Arc;
use arrow_schema::{Schema, SchemaRef};
use async_trait::async_trait;
use datafusion::catalog::{Session, TableProviderFactory};
use datafusion::datasource::TableProvider;
use datafusion::datasource::listing::PartitionedFile;
use datafusion::datasource::object_store::ObjectStoreUrl;
use datafusion::datasource::physical_plan::FileGroup;
use datafusion::datasource::physical_plan::FileScanConfigBuilder;
use datafusion::datasource::physical_plan::parquet::source::ParquetSource;
use datafusion::datasource::source::DataSourceExec;
use datafusion::error::Result;
use datafusion::logical_expr::Operator;
use datafusion::physical_plan::ExecutionPlan;
use datafusion_common::DFSchema;
use datafusion_common::DataFusionError::Execution;
use datafusion_common::config::TableParquetOptions;
use datafusion_common::stats::Precision;
use datafusion_common::{DataFusionError, Statistics};
use datafusion_expr::utils::split_conjunction;
use datafusion_expr::{CreateExternalTable, Expr, TableProviderFilterPushDown, TableType};
use datafusion_physical_expr::create_physical_expr;
use log::warn;
use crate::hudi_exec::HudiScanExec;
use crate::util::expr::exprs_to_filters;
use hudi_core::config::read::HudiReadConfig::{
FileSliceReadConcurrency, InputPartitions, UseReadOptimizedMode,
};
use hudi_core::config::table::{BaseFileFormatValue, HudiTableConfig};
use hudi_core::config::util::empty_options;
use hudi_core::config::{ConfigParser, HudiConfigs};
use hudi_core::file_group::file_slice::FileSlice;
use hudi_core::storage::util::{get_scheme_authority, join_url_segments};
use hudi_core::table::{ReadOptions, Table as HudiTable};
fn default_file_slice_read_concurrency() -> usize {
match FileSliceReadConcurrency.default_value() {
Some(value) => value.into(),
None => unreachable!("FileSliceReadConcurrency has a default value defined in hudi-core"),
}
}
pub(crate) fn inexact_usize_from_u64(value: u64) -> Precision<usize> {
match usize::try_from(value) {
Ok(value) => Precision::Inexact(value),
Err(_) => Precision::Absent,
}
}
pub(crate) fn external_error<E>(context: impl Into<String>, error: E) -> DataFusionError
where
E: Error + Send + Sync + 'static,
{
DataFusionError::External(Box::new(error)).context(context)
}
fn filter_field_matches_partition_column(filter_field: &str, partition_column: &str) -> bool {
filter_field == partition_column
|| filter_field
.rsplit_once('.')
.is_some_and(|(_, name)| name == partition_column)
}
#[derive(Clone)]
pub struct HudiDataSource {
table: Arc<HudiTable>,
schema: SchemaRef,
partition_schema: Schema,
cached_stats: Option<Statistics>,
input_partitions: usize,
read_optimized_mode: bool,
file_slice_read_concurrency: usize,
base_file_format: Option<BaseFileFormatValue>,
}
impl std::fmt::Debug for HudiDataSource {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("HudiDataSource")
.field("table", &self.table)
.field(
"partition_columns",
&self
.partition_schema
.fields()
.iter()
.map(|field| field.name())
.collect::<Vec<_>>(),
)
.finish()
}
}
impl HudiDataSource {
pub async fn new(base_uri: &str) -> Result<Self> {
Self::new_with_options(base_uri, empty_options()).await
}
pub async fn new_with_options<I, K, V>(base_uri: &str, options: I) -> Result<Self>
where
I: IntoIterator<Item = (K, V)>,
K: AsRef<str>,
V: Into<String>,
{
let all_options: Vec<(String, String)> = options
.into_iter()
.map(|(k, v)| (k.as_ref().to_string(), v.into()))
.collect();
let input_partitions: usize = match all_options
.iter()
.find(|(k, _)| k == InputPartitions.as_ref())
{
Some((_, v)) => v.parse().map_err(|_| {
Execution(format!(
"Invalid value '{v}' for {}: expected a non-negative integer",
InputPartitions.as_ref()
))
})?,
None => 0,
};
let read_optimized_mode: bool = match all_options
.iter()
.find(|(k, _)| k == UseReadOptimizedMode.as_ref())
{
Some((_, v)) => v.parse().map_err(|e| {
Execution(format!(
"Invalid value '{v}' for {}: {e}",
UseReadOptimizedMode.as_ref()
))
})?,
None => false,
};
let file_slice_read_concurrency: usize = match all_options
.iter()
.find(|(k, _)| k == FileSliceReadConcurrency.as_ref())
{
Some((_, v)) => {
let parsed = v.parse().map_err(|_| {
Execution(format!(
"Invalid value '{v}' for {}: expected a positive integer",
FileSliceReadConcurrency.as_ref()
))
})?;
if parsed == 0 {
return Err(Execution(format!(
"Invalid value '0' for {}: expected a positive integer",
FileSliceReadConcurrency.as_ref()
)));
}
parsed
}
None => default_file_slice_read_concurrency(),
};
let table = HudiTable::new_with_options(base_uri, all_options)
.await
.map_err(|e| external_error("Failed to create Hudi table", e))?;
let base_file_format =
BaseFileFormatValue::from_configs(&table.hudi_configs).map_err(|e| {
external_error(
format!(
"Invalid {} config",
HudiTableConfig::BaseFileFormat.as_ref()
),
e,
)
})?;
if matches!(base_file_format, Some(BaseFileFormatValue::HFile)) {
return Err(Execution(
"HFile is only supported for Hudi metadata tables, not regular DataFusion scans"
.to_string(),
));
}
let schema = table
.get_schema_with_meta_fields()
.await
.map(SchemaRef::from)
.unwrap_or_else(|e| {
warn!("Failed to get table schema, using empty schema: {e}");
SchemaRef::from(Schema::empty())
});
let partition_schema = match table.get_partition_schema().await {
Ok(s) => s,
Err(e) => {
warn!("Failed to get partition schema, using empty schema: {e}");
Schema::empty()
}
};
let cached_stats = match table.compute_table_stats(None).await {
Some((num_rows, total_byte_size)) => {
let num_fields = schema.fields().len();
Some(Statistics {
num_rows: inexact_usize_from_u64(num_rows),
total_byte_size: inexact_usize_from_u64(total_byte_size),
column_statistics: vec![
datafusion_common::ColumnStatistics::new_unknown();
num_fields
],
})
}
None => None,
};
Ok(Self {
table: Arc::new(table),
schema,
partition_schema,
cached_stats,
input_partitions,
read_optimized_mode,
file_slice_read_concurrency,
base_file_format,
})
}
fn get_input_partitions(&self) -> usize {
self.input_partitions
}
#[cfg(test)]
fn get_file_slice_read_concurrency(&self) -> usize {
self.file_slice_read_concurrency
}
fn effective_read_optimized(&self) -> bool {
self.read_optimized_mode
}
fn scan_read_options(
&self,
pushdown_filters: Vec<(String, String, String)>,
read_optimized: bool,
) -> Result<ReadOptions> {
let mut read_options = ReadOptions::new()
.with_filters(pushdown_filters)
.map_err(|e| external_error("Invalid pushdown filter", e))?;
if read_optimized {
read_options = read_options.with_hudi_option(UseReadOptimizedMode.as_ref(), "true");
}
Ok(read_options)
}
fn read_options_for_hudi_exec(
hudi_configs: &HudiConfigs,
options: &ReadOptions,
) -> ReadOptions {
let drops_partition_columns: bool = hudi_configs
.get_or_default(HudiTableConfig::DropsPartitionFields)
.into();
if !drops_partition_columns || options.filters.is_empty() {
return options.clone();
}
let partition_columns: Vec<String> = hudi_configs
.get_or_default(HudiTableConfig::PartitionFields)
.into();
let mut applicable = options.clone();
applicable.filters = options
.filters
.iter()
.filter(|filter| {
!partition_columns
.iter()
.any(|p| filter_field_matches_partition_column(&filter.field, p))
})
.cloned()
.collect();
applicable
}
fn file_slices_are_parquet(file_slices: &[FileSlice]) -> Result<bool> {
if file_slices.is_empty() {
return Ok(false);
}
for file_slice in file_slices {
let relative_path = file_slice.base_file_relative_path().map_err(|e| {
external_error(
format!("Failed to get base file relative path for {file_slice:?}"),
e,
)
})?;
let Some(relative_path) = relative_path else {
return Ok(true);
};
if !BaseFileFormatValue::Parquet.matches_extension(&relative_path) {
return Ok(false);
}
}
Ok(true)
}
fn use_parquet_source(
&self,
read_options: &ReadOptions,
file_slices: &[FileSlice],
) -> Result<bool> {
let parquet_base_files = match &self.base_file_format {
Some(format) => matches!(format, BaseFileFormatValue::Parquet),
None => Self::file_slices_are_parquet(file_slices)?,
};
if !parquet_base_files {
return Ok(false);
}
if !self.table.is_mor() {
return Ok(true);
}
read_options
.is_read_optimized()
.map_err(|e| external_error("Invalid read-optimized option", e))
}
fn can_push_down_expr(schema: &Schema, expr: &Expr) -> bool {
match expr {
Expr::BinaryExpr(binary_expr) => {
let left = &binary_expr.left;
let op = &binary_expr.op;
let right = &binary_expr.right;
match op {
Operator::And => {
Self::can_push_down_expr(schema, left)
|| Self::can_push_down_expr(schema, right)
}
Operator::Or => {
false
}
_ => {
Self::is_supported_operator(op)
&& Self::is_supported_operand(schema, left)
&& Self::is_supported_operand(schema, right)
}
}
}
Expr::Not(inner_expr) => {
Self::can_push_down_expr(schema, inner_expr)
}
Expr::Between(between) => {
!between.negated
&& matches!(&*between.expr, Expr::Column(col) if schema.column_with_name(&col.name).is_some())
&& matches!(&*between.low, Expr::Literal(..))
&& matches!(&*between.high, Expr::Literal(..))
}
Expr::InList(in_list) => {
!in_list.list.is_empty()
&& matches!(in_list.expr.as_ref(), Expr::Column(col) if schema.column_with_name(&col.name).is_some())
&& in_list
.list
.iter()
.all(|expr| matches!(expr, Expr::Literal(..)))
}
_ => false,
}
}
fn is_supported_operator(op: &Operator) -> bool {
matches!(
op,
Operator::Eq
| Operator::NotEq
| Operator::Gt
| Operator::Lt
| Operator::GtEq
| Operator::LtEq
)
}
fn is_supported_operand(schema: &Schema, expr: &Expr) -> bool {
match expr {
Expr::Column(col) => schema.column_with_name(&col.name).is_some(),
Expr::Literal(..) => true,
_ => false,
}
}
fn get_partition_columns(&self) -> Vec<String> {
self.partition_schema
.fields()
.iter()
.map(|f| f.name().clone())
.collect()
}
fn is_partition_column_filter(expr: &Expr, partition_cols: &[String]) -> bool {
match expr {
Expr::BinaryExpr(binary_expr) => match binary_expr.op {
Operator::And => {
Self::is_partition_column_filter(&binary_expr.left, partition_cols)
&& Self::is_partition_column_filter(&binary_expr.right, partition_cols)
}
Operator::Or => false,
_ => match (&*binary_expr.left, &*binary_expr.right) {
(Expr::Column(col), Expr::Literal(..))
| (Expr::Literal(..), Expr::Column(col)) => partition_cols.contains(&col.name),
_ => false,
},
},
Expr::Not(inner) => Self::is_partition_column_filter(inner, partition_cols),
Expr::Between(between) => {
!between.negated
&& matches!(&*between.expr, Expr::Column(col) if partition_cols.contains(&col.name))
&& matches!(&*between.low, Expr::Literal(..))
&& matches!(&*between.high, Expr::Literal(..))
}
Expr::InList(in_list) => {
!in_list.list.is_empty()
&& matches!(in_list.expr.as_ref(), Expr::Column(col) if partition_cols.contains(&col.name))
&& in_list
.list
.iter()
.all(|expr| matches!(expr, Expr::Literal(..)))
}
_ => false,
}
}
fn is_exact_partition_equality_filter(expr: &Expr, partition_cols: &[String]) -> bool {
match expr {
Expr::BinaryExpr(binary_expr) if binary_expr.op == Operator::Eq => {
match (&*binary_expr.left, &*binary_expr.right) {
(Expr::Column(col), Expr::Literal(..))
| (Expr::Literal(..), Expr::Column(col)) => partition_cols.contains(&col.name),
_ => false,
}
}
_ => false,
}
}
fn filter_pushdown_support(
table_schema: &Schema,
partition_cols: &[String],
expr: &Expr,
) -> TableProviderFilterPushDown {
let conjuncts = split_conjunction(expr);
let has_pushable_conjunct = conjuncts
.iter()
.any(|conjunct| Self::can_push_down_expr(table_schema, conjunct));
if !has_pushable_conjunct {
return TableProviderFilterPushDown::Unsupported;
}
let all_conjuncts_are_exact_partition_eq = conjuncts.iter().all(|conjunct| {
Self::can_push_down_expr(table_schema, conjunct)
&& Self::is_exact_partition_equality_filter(conjunct, partition_cols)
});
if all_conjuncts_are_exact_partition_eq {
TableProviderFilterPushDown::Exact
} else {
TableProviderFilterPushDown::Inexact
}
}
fn split_scan_pushdown_exprs(&self, filters: &[Expr]) -> (Vec<Expr>, Vec<Expr>) {
let partition_cols = self.get_partition_columns();
Self::split_scan_pushdown_exprs_for_schema(self.schema.as_ref(), &partition_cols, filters)
}
fn split_scan_pushdown_exprs_for_schema(
table_schema: &Schema,
partition_cols: &[String],
filters: &[Expr],
) -> (Vec<Expr>, Vec<Expr>) {
let all_pushdown_exprs: Vec<Expr> = filters
.iter()
.flat_map(|expr| split_conjunction(expr).into_iter())
.filter(|expr| Self::can_push_down_expr(table_schema, expr))
.cloned()
.collect();
let partition_pushdown_exprs = all_pushdown_exprs
.iter()
.filter(|expr| Self::is_partition_column_filter(expr, partition_cols))
.cloned()
.collect();
(partition_pushdown_exprs, all_pushdown_exprs)
}
fn use_parquet_source_without_file_slices(
&self,
read_options: &ReadOptions,
) -> Result<Option<bool>> {
if self.table.is_mor()
&& !read_options
.is_read_optimized()
.map_err(|e| external_error("Invalid read-optimized option", e))?
{
return Ok(Some(false));
}
match &self.base_file_format {
Some(BaseFileFormatValue::Parquet) => Ok(Some(true)),
Some(_) => Ok(Some(false)),
None => Ok(None),
}
}
#[allow(clippy::too_many_arguments)]
async fn scan_parquet(
&self,
state: &dyn Session,
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
flat_slices: Vec<FileSlice>,
) -> Result<Arc<dyn ExecutionPlan>> {
let input_partitions = self.get_input_partitions_for_scan(state);
let file_slices =
hudi_core::util::collection::split_into_chunks(flat_slices, input_partitions);
let base_url = self.table.base_url();
let mut parquet_file_groups: Vec<Vec<PartitionedFile>> = Vec::new();
for file_slice_vec in file_slices {
let mut parquet_file_group_vec = Vec::new();
for f in file_slice_vec {
let relative_path = f.base_file_relative_path().map_err(|e| {
external_error(
format!("Failed to get base file relative path for {f:?}"),
e,
)
})?;
let Some(relative_path) = relative_path else {
continue;
};
let url = join_url_segments(&base_url, &[relative_path.as_str()])
.map_err(|e| external_error("Failed to join URL segments", e))?;
let size = f
.base_file
.as_ref()
.and_then(|b| b.file_metadata.as_ref())
.map_or(0, |m| m.size);
let partitioned_file = PartitionedFile::new(url.path(), size);
parquet_file_group_vec.push(partitioned_file);
}
parquet_file_groups.push(parquet_file_group_vec)
}
let url = ObjectStoreUrl::parse(get_scheme_authority(&base_url))?;
let parquet_opts = TableParquetOptions {
global: state.config_options().execution.parquet.clone(),
column_specific_options: Default::default(),
key_value_metadata: Default::default(),
crypto: Default::default(),
};
let table_schema = self.schema();
let mut parquet_source = ParquetSource::new(table_schema.clone())
.with_table_parquet_options(parquet_opts)
.with_pushdown_filters(true)
.with_reorder_filters(true)
.with_enable_page_index(true);
let filter = filters.iter().cloned().reduce(|acc, new| acc.and(new));
if let Some(expr) = filter {
let df_schema = DFSchema::try_from(table_schema.clone())?;
let predicate = create_physical_expr(&expr, &df_schema, state.execution_props())?;
parquet_source = parquet_source.with_predicate(predicate)
}
let file_groups: Vec<FileGroup> = parquet_file_groups
.into_iter()
.map(FileGroup::from)
.collect();
let mut fsc_builder = FileScanConfigBuilder::new(url, Arc::new(parquet_source))
.with_file_groups(file_groups)
.with_projection_indices(projection.cloned())?
.with_limit(limit);
if let Some(stats) = &self.cached_stats {
fsc_builder = fsc_builder.with_statistics(stats.clone());
}
let fsc = fsc_builder.build();
Ok(Arc::new(DataSourceExec::new(Arc::new(fsc))))
}
async fn scan_hudi(
&self,
projection: Option<&Vec<usize>>,
limit: Option<usize>,
input_partitions: usize,
flat_slices: Vec<FileSlice>,
read_options: ReadOptions,
) -> Result<Arc<dyn ExecutionPlan>> {
let slice_log_bytes: Vec<Option<u64>> =
flat_slices.iter().map(FileSlice::log_size_bytes).collect();
let file_slice_read_concurrency = hudi_core::file_group::admission::slices_in_flight(
self.scan_max_memory_size(),
input_partitions,
&slice_log_bytes,
self.file_slice_read_concurrency,
);
let file_slices =
hudi_core::util::collection::split_into_chunks(flat_slices, input_partitions);
let file_group_reader = Arc::new(
self.table
.create_file_group_reader_with_options(Some(&read_options), empty_options())
.await
.map_err(|e| external_error("Failed to create FileGroupReader", e))?,
);
let mut hudi_read_options =
Self::read_options_for_hudi_exec(&self.table.hudi_configs, &read_options);
if let Some(proj) = projection {
let col_names: Vec<String> = proj
.iter()
.map(|&i| self.schema.field(i).name().clone())
.collect();
hudi_read_options = hudi_read_options.with_projection(col_names);
}
Ok(Arc::new(HudiScanExec::new(
file_slices,
file_group_reader,
hudi_read_options,
input_partitions,
file_slice_read_concurrency,
self.schema.clone(),
projection.cloned(),
limit,
)))
}
fn scan_max_memory_size(&self) -> Option<u64> {
self.table
.hudi_configs
.try_get(hudi_core::config::read::HudiReadConfig::ScanMaxMemorySize)
.ok()
.flatten()
.map(|v| -> usize { v.into() })
.map(|v| v as u64)
}
fn get_input_partitions_for_scan(&self, state: &dyn Session) -> usize {
match self.get_input_partitions() {
0 => state.config_options().execution.target_partitions,
n => n,
}
}
}
#[async_trait]
impl TableProvider for HudiDataSource {
fn schema(&self) -> SchemaRef {
self.schema.clone()
}
fn table_type(&self) -> TableType {
TableType::Base
}
fn statistics(&self) -> Option<Statistics> {
self.cached_stats.clone()
}
async fn scan(
&self,
state: &dyn Session,
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>> {
self.table.register_storage(state.runtime_env().clone());
let input_partitions = self.get_input_partitions_for_scan(state);
let (partition_pushdown_exprs, all_pushdown_exprs) =
self.split_scan_pushdown_exprs(filters);
let partition_pushdown_filters = exprs_to_filters(&partition_pushdown_exprs);
let all_pushdown_filters = exprs_to_filters(&all_pushdown_exprs);
let all_filters_are_partition_filters = all_pushdown_filters == partition_pushdown_filters;
let read_optimized = self.effective_read_optimized();
let partition_read_options =
self.scan_read_options(partition_pushdown_filters.clone(), read_optimized)?;
let all_read_options = if all_filters_are_partition_filters {
partition_read_options.clone()
} else {
self.scan_read_options(all_pushdown_filters, read_optimized)?
};
match self.use_parquet_source_without_file_slices(&partition_read_options)? {
Some(true) => {
let flat_slices = self
.table
.get_file_slices(&partition_read_options)
.await
.map_err(|e| external_error("Failed to get file slices from Hudi table", e))?;
self.scan_parquet(state, projection, filters, limit, flat_slices)
.await
}
Some(false) => {
let flat_slices = self
.table
.get_file_slices(&all_read_options)
.await
.map_err(|e| external_error("Failed to get file slices from Hudi table", e))?;
self.scan_hudi(
projection,
limit,
input_partitions,
flat_slices,
all_read_options,
)
.await
}
None => {
let partition_flat_slices = self
.table
.get_file_slices(&partition_read_options)
.await
.map_err(|e| external_error("Failed to get file slices from Hudi table", e))?;
if self.use_parquet_source(&partition_read_options, &partition_flat_slices)? {
self.scan_parquet(state, projection, filters, limit, partition_flat_slices)
.await
} else {
let flat_slices = if all_filters_are_partition_filters {
partition_flat_slices
} else {
self.table
.get_file_slices(&all_read_options)
.await
.map_err(|e| {
external_error("Failed to get file slices from Hudi table", e)
})?
};
self.scan_hudi(
projection,
limit,
input_partitions,
flat_slices,
all_read_options,
)
.await
}
}
}
}
fn supports_filters_pushdown(
&self,
filters: &[&Expr],
) -> Result<Vec<TableProviderFilterPushDown>> {
let partition_cols = self.get_partition_columns();
filters
.iter()
.map(|expr| {
Ok(Self::filter_pushdown_support(
self.schema.as_ref(),
&partition_cols,
expr,
))
})
.collect()
}
}
#[derive(Debug)]
pub struct HudiTableFactory {}
impl HudiTableFactory {
pub fn new() -> Self {
Self {}
}
fn resolve_options(
state: &dyn Session,
cmd: &CreateExternalTable,
) -> Result<HashMap<String, String>> {
let mut options: HashMap<_, _> = state
.config_options()
.entries()
.iter()
.filter_map(|e| {
let value = e.value.as_ref().filter(|v| !v.is_empty())?;
Some((e.key.clone(), value.clone()))
})
.collect();
options.extend(cmd.options.iter().map(|(k, v)| (k.clone(), v.clone())));
Ok(options)
}
}
impl Default for HudiTableFactory {
fn default() -> Self {
Self::new()
}
}
#[async_trait]
impl TableProviderFactory for HudiTableFactory {
async fn create(
&self,
state: &dyn Session,
cmd: &CreateExternalTable,
) -> Result<Arc<dyn TableProvider>> {
let options = HudiTableFactory::resolve_options(state, cmd)?;
let base_uri = cmd.location.as_str();
let table_provider = HudiDataSource::new_with_options(base_uri, options).await?;
Ok(Arc::new(table_provider))
}
}
#[cfg(test)]
mod tests {
use super::*;
use arrow_schema::{DataType, Field};
use datafusion_common::{Column, ScalarValue};
use hudi_core::config::internal::HudiInternalConfig;
use hudi_core::config::table::{BaseFileFormatValue, HudiTableConfig};
use std::fs::canonicalize;
use std::path::Path;
use url::Url;
use datafusion::logical_expr::BinaryExpr;
use datafusion::prelude::SessionContext;
use hudi_test::SampleTable::{V6Nonpartitioned, V6SimplekeygenNonhivestyle, V9LanceTxnsSimple};
use crate::HudiDataSource;
#[tokio::test]
async fn get_default_input_partitions() {
let base_url =
Url::from_file_path(canonicalize(Path::new("tests/data/table_props_valid")).unwrap())
.unwrap();
let hudi = HudiDataSource::new(base_url.as_str()).await.unwrap();
assert_eq!(hudi.get_input_partitions(), 0);
assert_eq!(
hudi.get_file_slice_read_concurrency(),
default_file_slice_read_concurrency()
);
assert_eq!(hudi.table_type(), TableType::Base);
assert_eq!(hudi.statistics(), None);
}
#[tokio::test]
async fn test_new_with_options_sets_file_slice_read_concurrency() {
let hudi = HudiDataSource::new_with_options(
V6Nonpartitioned.path_to_cow().as_str(),
[(FileSliceReadConcurrency.as_ref(), "2")],
)
.await
.unwrap();
assert_eq!(hudi.get_file_slice_read_concurrency(), 2);
}
#[tokio::test]
async fn test_new_with_options_rejects_invalid_file_slice_read_concurrency() {
for invalid in ["0", "abc"] {
let result = HudiDataSource::new_with_options(
V6Nonpartitioned.path_to_cow().as_str(),
[(FileSliceReadConcurrency.as_ref(), invalid)],
)
.await;
assert!(result.is_err());
let error = result.unwrap_err().to_string();
assert!(error.contains(FileSliceReadConcurrency.as_ref()));
assert!(error.contains(invalid));
}
}
#[test]
fn test_file_slices_are_parquet_empty_is_false() {
assert!(!HudiDataSource::file_slices_are_parquet(&[]).unwrap());
}
#[tokio::test]
async fn test_new_with_options_rejects_hfile_format_for_regular_scan() {
let result = HudiDataSource::new_with_options(
V6Nonpartitioned.path_to_cow().as_str(),
[
(
HudiTableConfig::BaseFileFormat.as_ref(),
BaseFileFormatValue::HFile.as_ref(),
),
(HudiInternalConfig::SkipConfigValidation.as_ref(), "true"),
],
)
.await;
assert!(
result.is_err(),
"HFile format should be rejected for regular DataFusion scans"
);
assert!(result.unwrap_err().to_string().contains("HFile"));
}
#[tokio::test]
async fn test_new_with_options_rejects_invalid_base_file_format_config() {
let result = HudiDataSource::new_with_options(
V6Nonpartitioned.path_to_cow().as_str(),
[
(HudiTableConfig::BaseFileFormat.as_ref(), "orc"),
(HudiInternalConfig::SkipConfigValidation.as_ref(), "true"),
],
)
.await;
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("orc"));
}
fn pushdown_test_schema() -> Schema {
Schema::new(vec![
Field::new("byteField", DataType::Int8, false),
Field::new("name", DataType::Utf8, true),
Field::new("intField", DataType::Int32, true),
])
}
fn partition_cols() -> Vec<String> {
vec!["byteField".to_string()]
}
fn col_lit(name: &str, op: Operator, lit: ScalarValue) -> Expr {
Expr::BinaryExpr(BinaryExpr {
left: Box::new(Expr::Column(Column::from_name(name.to_string()))),
op,
right: Box::new(Expr::Literal(lit, None)),
})
}
fn pushdown_support(
schema: &Schema,
partition_cols: &[String],
expr: &Expr,
) -> TableProviderFilterPushDown {
HudiDataSource::filter_pushdown_support(schema, partition_cols, expr)
}
async fn scan_hudi_exec_filter_triplets(
hudi: &HudiDataSource,
filters: &[Expr],
) -> Vec<(String, String, Vec<String>)> {
let ctx = SessionContext::new();
let state = ctx.state();
let plan = hudi.scan(&state, None, filters, None).await.unwrap();
let exec = plan
.downcast_ref::<HudiScanExec>()
.expect("scan should route to HudiScanExec");
exec.read_options()
.filters
.iter()
.map(|filter| {
(
filter.field.clone(),
filter.operator.to_string(),
filter.values.clone(),
)
})
.collect()
}
fn assert_scan_filter(
filters: &[(String, String, Vec<String>)],
field: &str,
operator: &str,
values: &[&str],
) {
assert!(
filters
.iter()
.any(|(actual_field, actual_operator, actual_values)| {
actual_field == field
&& actual_operator == operator
&& actual_values
.iter()
.map(String::as_str)
.eq(values.iter().copied())
}),
"expected filter ({field}, {operator}, {values:?}) in {filters:?}"
);
}
#[test]
fn test_filter_pushdown_support_for_non_partitioned_schema() {
let schema = pushdown_test_schema();
let partition_cols = vec![];
let filters = [
col_lit(
"name",
Operator::Eq,
ScalarValue::Utf8(Some("Alice".to_string())),
),
col_lit("intField", Operator::Gt, ScalarValue::Int32(Some(20000))),
col_lit(
"nonexistent_column",
Operator::Eq,
ScalarValue::Int32(Some(1)),
),
col_lit(
"name",
Operator::NotEq,
ScalarValue::Utf8(Some("Diana".to_string())),
),
Expr::Literal(ScalarValue::Int32(Some(10)), None),
Expr::Not(Box::new(col_lit(
"intField",
Operator::Gt,
ScalarValue::Int32(Some(20000)),
))),
];
let result = filters
.iter()
.map(|expr| pushdown_support(&schema, &partition_cols, expr))
.collect::<Vec<_>>();
assert_eq!(
result,
vec![
TableProviderFilterPushDown::Inexact,
TableProviderFilterPushDown::Inexact,
TableProviderFilterPushDown::Unsupported,
TableProviderFilterPushDown::Inexact,
TableProviderFilterPushDown::Unsupported,
TableProviderFilterPushDown::Inexact,
]
);
}
#[test]
fn test_filter_pushdown_exact_only_for_partition_equality() {
let schema = pushdown_test_schema();
let partition_cols = partition_cols();
let partition_eq = col_lit("byteField", Operator::Eq, ScalarValue::Int8(Some(1)));
let partition_gt = col_lit("byteField", Operator::Gt, ScalarValue::Int8(Some(1)));
let non_partition_eq = col_lit(
"name",
Operator::Eq,
ScalarValue::Utf8(Some("Alice".to_string())),
);
assert_eq!(
pushdown_support(&schema, &partition_cols, &partition_eq),
TableProviderFilterPushDown::Exact
);
assert_eq!(
pushdown_support(&schema, &partition_cols, &partition_gt),
TableProviderFilterPushDown::Inexact
);
assert_eq!(
pushdown_support(&schema, &partition_cols, &non_partition_eq),
TableProviderFilterPushDown::Inexact
);
}
#[test]
fn test_filter_pushdown_splits_conjunctions_for_classification() {
let schema = pushdown_test_schema();
let partition_cols = partition_cols();
let partition_eq = col_lit("byteField", Operator::Eq, ScalarValue::Int8(Some(1)));
let second_partition_eq = col_lit("byteField", Operator::Eq, ScalarValue::Int8(Some(2)));
let non_partition_eq = col_lit(
"name",
Operator::Eq,
ScalarValue::Utf8(Some("Alice".to_string())),
);
let unsupported = Expr::Literal(ScalarValue::Boolean(Some(true)), None);
assert_eq!(
pushdown_support(
&schema,
&partition_cols,
&partition_eq.clone().and(second_partition_eq)
),
TableProviderFilterPushDown::Exact
);
assert_eq!(
pushdown_support(
&schema,
&partition_cols,
&partition_eq.clone().and(non_partition_eq)
),
TableProviderFilterPushDown::Inexact
);
assert_eq!(
pushdown_support(&schema, &partition_cols, &partition_eq.and(unsupported)),
TableProviderFilterPushDown::Inexact
);
}
#[test]
fn test_filter_pushdown_between_in_list_and_or() {
let schema = pushdown_test_schema();
let partition_cols = partition_cols();
let partition_between = Expr::Between(datafusion_expr::Between::new(
Box::new(Expr::Column(Column::from_name("byteField".to_string()))),
false,
Box::new(Expr::Literal(ScalarValue::Int8(Some(1)), None)),
Box::new(Expr::Literal(ScalarValue::Int8(Some(3)), None)),
));
let partition_in = Expr::InList(datafusion_expr::expr::InList::new(
Box::new(Expr::Column(Column::from_name("byteField".to_string()))),
vec![
Expr::Literal(ScalarValue::Int8(Some(1)), None),
Expr::Literal(ScalarValue::Int8(Some(2)), None),
],
false,
));
let or_expr = col_lit(
"name",
Operator::Eq,
ScalarValue::Utf8(Some("Alice".to_string())),
)
.or(col_lit(
"name",
Operator::Eq,
ScalarValue::Utf8(Some("Bob".to_string())),
));
assert_eq!(
pushdown_support(&schema, &partition_cols, &partition_between),
TableProviderFilterPushDown::Inexact
);
assert_eq!(
pushdown_support(&schema, &partition_cols, &partition_in),
TableProviderFilterPushDown::Inexact
);
assert_eq!(
pushdown_support(&schema, &partition_cols, &or_expr),
TableProviderFilterPushDown::Unsupported
);
}
#[test]
fn test_scan_pushdown_exprs_separates_partition_filters_for_listing() {
let schema = pushdown_test_schema();
let partition_cols = partition_cols();
let partition_eq = col_lit("byteField", Operator::Eq, ScalarValue::Int8(Some(1)));
let non_partition_eq = col_lit(
"name",
Operator::Eq,
ScalarValue::Utf8(Some("Alice".to_string())),
);
let partition_in = Expr::InList(datafusion_expr::expr::InList::new(
Box::new(Expr::Column(Column::from_name("byteField".to_string()))),
vec![
Expr::Literal(ScalarValue::Int8(Some(1)), None),
Expr::Literal(ScalarValue::Int8(Some(2)), None),
],
false,
));
let (partition_exprs, all_exprs) = HudiDataSource::split_scan_pushdown_exprs_for_schema(
&schema,
&partition_cols,
&[partition_eq.and(non_partition_eq), partition_in],
);
let partition_filters = exprs_to_filters(&partition_exprs);
let all_filters = exprs_to_filters(&all_exprs);
let partition_fields = partition_filters
.iter()
.map(|(field, _, _)| field.as_str())
.collect::<Vec<_>>();
let all_fields = all_filters
.iter()
.map(|(field, _, _)| field.as_str())
.collect::<Vec<_>>();
assert_eq!(partition_fields, ["byteField", "byteField"]);
assert_eq!(all_fields, ["byteField", "name", "byteField"]);
}
#[test]
fn test_read_options_for_hudi_exec_strips_dropped_partition_filters() {
let hudi_configs = HudiConfigs::new([
(HudiTableConfig::DropsPartitionFields, "true"),
(HudiTableConfig::PartitionFields, "region,country"),
]);
let read_options = ReadOptions::new()
.with_filters([
("region", "=", "us"),
("amount", ">", "10"),
("txns.country", "=", "ca"),
])
.unwrap()
.with_projection(["txn_id", "amount"]);
let actual = HudiDataSource::read_options_for_hudi_exec(&hudi_configs, &read_options);
assert_eq!(actual.filters.len(), 1);
assert_eq!(actual.filters[0].field, "amount");
assert_eq!(actual.projection, read_options.projection);
assert_eq!(actual.hudi_options, read_options.hudi_options);
}
#[tokio::test]
async fn test_scan_hudi_keeps_inexact_non_partition_filters() {
let hudi = HudiDataSource::new(V6SimplekeygenNonhivestyle.url_to_mor_parquet().as_str())
.await
.unwrap();
let read_options = hudi
.scan_read_options(
vec![("id".to_string(), ">".to_string(), "1".to_string())],
false,
)
.unwrap();
let flat_slices = hudi.table.get_file_slices(&read_options).await.unwrap();
let plan = hudi
.scan_hudi(None, None, 1, flat_slices, read_options)
.await
.unwrap();
let exec = plan
.downcast_ref::<HudiScanExec>()
.expect("MOR snapshot scan should use HudiScanExec");
assert_eq!(exec.read_options().filters.len(), 1);
let filter = &exec.read_options().filters[0];
assert_eq!(filter.field, "id");
assert_eq!(filter.values, vec!["1".to_string()]);
}
#[tokio::test]
async fn a_scan_memory_budget_lowers_the_planned_slice_concurrency() {
use datafusion::physical_plan::displayable;
async fn planned_concurrency(options: Vec<(&str, &str)>) -> String {
let hudi = HudiDataSource::new_with_options(
V6SimplekeygenNonhivestyle.url_to_mor_parquet().as_str(),
options,
)
.await
.unwrap();
let ctx = SessionContext::new();
let state = ctx.state();
let plan = hudi.scan(&state, None, &[], None).await.unwrap();
displayable(plan.as_ref())
.set_show_schema(false)
.indent(true)
.to_string()
}
let unbounded = planned_concurrency(vec![]).await;
assert!(
unbounded.contains("file_slice_read_concurrency=4"),
"the default ceiling should stand with no budget: {unbounded}"
);
let bounded =
planned_concurrency(vec![("hoodie.read.scan.max.memory.size", "1048576")]).await;
assert!(
bounded.contains("file_slice_read_concurrency=1"),
"a 1 MiB budget must lower the fan-out to 1: {bounded}"
);
}
#[tokio::test]
async fn test_scan_mor_snapshot_keeps_partition_and_non_partition_filters_for_hudi_exec() {
let hudi = HudiDataSource::new(V6SimplekeygenNonhivestyle.url_to_mor_parquet().as_str())
.await
.unwrap();
let partition_filter = col_lit("byteField", Operator::Eq, ScalarValue::Int32(Some(10)));
let non_partition_filter = col_lit("id", Operator::Gt, ScalarValue::Int32(Some(1)));
let filters =
scan_hudi_exec_filter_triplets(&hudi, &[partition_filter, non_partition_filter]).await;
assert_scan_filter(&filters, "byteField", "=", &["10"]);
assert_scan_filter(&filters, "id", ">", &["1"]);
}
#[tokio::test]
async fn test_scan_lance_keeps_partition_and_non_partition_filters_for_hudi_exec() {
let hudi = HudiDataSource::new(V9LanceTxnsSimple.url_to_cow().as_str())
.await
.unwrap();
let partition_filter = col_lit(
"region",
Operator::Eq,
ScalarValue::Utf8(Some("us".to_string())),
);
let non_partition_filter = col_lit(
"txn_id",
Operator::Eq,
ScalarValue::Utf8(Some("TXN-001".to_string())),
);
let filters =
scan_hudi_exec_filter_triplets(&hudi, &[partition_filter, non_partition_filter]).await;
assert_scan_filter(&filters, "region", "=", &["us"]);
assert_scan_filter(&filters, "txn_id", "=", &["TXN-001"]);
}
}