use crate::Result;
use crate::config::HudiConfigs;
use crate::config::error::ConfigError;
use crate::config::read::HudiReadConfig;
use crate::config::table::{BaseFileFormatValue, HudiTableConfig};
use crate::error::CoreError;
use crate::file_group::reader_v2::buffer::spillable_map;
use crate::file_group::reader_v2::metadata_merger::resolve_custom_merger;
use crate::file_group::reader_v2::reader_context::{CONFIG_MERGE_TYPE, MergeMode, ReaderContext};
use crate::file_group::reader_v2::record_context::RecordContext;
use crate::file_group::reader_v2::schema_handler::FileGroupReaderSchemaHandler;
use crate::timeline::selector::InstantRange;
use std::collections::HashMap;
#[allow(dead_code)]
pub(crate) fn resolve_reader_context(
hudi_configs: &HudiConfigs,
has_log_files: bool,
base_file_path: Option<&str>,
) -> Result<ReaderContext> {
let table_path: String = hudi_configs.get(HudiTableConfig::BasePath)?.into();
let latest_commit_time: String = hudi_configs
.try_get(HudiReadConfig::EndTimestamp)?
.ok_or_else(|| ConfigError::NotFound(HudiReadConfig::EndTimestamp.as_ref().to_string()))?
.into();
let merge_mode = match resolve_merge_mode(hudi_configs) {
Ok(mode) => mode,
Err(CoreError::Unsupported(_)) if !has_log_files => MergeMode::CommitTimeOrdering,
Err(e) => return Err(e),
};
let base_file_format = BaseFileFormatValue::resolve_from_configs(hudi_configs, base_file_path)?;
let instant_range = resolve_instant_range(hudi_configs)?;
let (table_config, hoodie_reader_config) = partition_configs(hudi_configs);
let merge_strategy_id = RECORD_MERGE_STRATEGY_ID_KEYS
.iter()
.find_map(|k| table_config.get(*k))
.cloned()
.unwrap_or_default();
Ok(ReaderContext {
table_path,
latest_commit_time,
merge_mode: merge_mode.as_ref().to_string(),
base_file_format: base_file_format.as_ref().to_string(),
has_log_files,
instant_range: Some(instant_range),
table_config,
hoodie_reader_config,
should_merge_use_record_position: resolve_use_record_position(hudi_configs)?,
iterator_mode: "ENGINE_RECORD".to_string(),
merge_strategy_id,
has_bootstrap_base_file: false,
needs_bootstrap_merge: false,
enable_logical_timestamp_field_repair: false,
row_filter_builder: None,
row_group_selector: None,
key_predicate: None,
mor_pk_safe: false,
repair_risk_columns: Vec::new(),
completion_gate_inputs: None,
record_context: RecordContext::default(),
schema_handler: FileGroupReaderSchemaHandler::new(),
})
}
#[allow(dead_code)]
const RECORD_MERGE_MODE: &str = "hoodie.record.merge.mode";
#[allow(dead_code)]
const READ_CONFIG_PREFIX: &str = "hoodie.read.";
#[allow(dead_code)]
const CRATE_CONFIG_PREFIXES: [&str; 2] = ["hoodie.internal.", "hoodie.plan."];
#[allow(dead_code)]
const READER_CONFIG_KEYS: [&str; 5] = [
spillable_map::CONFIG_MERGE_MAX_SIZE,
spillable_map::CONFIG_MAX_PEAK_MEMORY,
spillable_map::CONFIG_SPILLABLE_MAP_PATH,
spillable_map::CONFIG_DISKMAP_TYPE,
CONFIG_MERGE_TYPE,
];
#[allow(dead_code)]
fn partition_configs(
hudi_configs: &HudiConfigs,
) -> (HashMap<String, String>, HashMap<String, String>) {
let mut table_config = HashMap::new();
let mut reader_config = HashMap::new();
for (key, value) in hudi_configs.as_options() {
if key.starts_with(READ_CONFIG_PREFIX) || READER_CONFIG_KEYS.contains(&key.as_str()) {
reader_config.insert(key, value);
} else if !CRATE_CONFIG_PREFIXES
.iter()
.any(|prefix| key.starts_with(prefix))
{
table_config.insert(key, value);
}
}
(table_config, reader_config)
}
#[allow(dead_code)]
fn resolve_instant_range(hudi_configs: &HudiConfigs) -> Result<InstantRange> {
let timezone: String = hudi_configs
.get_or_default(HudiTableConfig::TimelineTimezone)
.into();
let start_timestamp = hudi_configs
.try_get(HudiReadConfig::StartTimestamp)?
.map(|v| -> String { v.into() });
let end_timestamp = hudi_configs
.try_get(HudiReadConfig::EndTimestamp)?
.map(|v| -> String { v.into() });
Ok(InstantRange::new(
timezone,
start_timestamp,
end_timestamp,
false,
true,
))
}
#[allow(dead_code)]
fn resolve_use_record_position(hudi_configs: &HudiConfigs) -> Result<bool> {
Ok(hudi_configs
.try_get(HudiReadConfig::MergeUseRecordPositions)?
.map(|v| -> bool { v.into() })
.unwrap_or(false))
}
pub(crate) const RECORD_MERGE_STRATEGY_ID_KEYS: [&str; 2] = [
"hoodie.record.merge.strategy.id",
"hoodie.compaction.record.merger.strategy",
];
pub(crate) const PAYLOAD_CLASS_KEYS: [&str; 3] = [
"hoodie.compaction.payload.class",
"hoodie.datasource.write.payload.class",
"hoodie.table.legacy.payload.class",
];
const EVENT_TIME_STRATEGY_ID: &str = "eeb8d96f-b1e4-49fd-bbf8-28ac514178e5";
const COMMIT_TIME_STRATEGY_ID: &str = "ce9acb64-bde0-424c-9b91-f6ebba25356d";
const EVENT_TIME_PAYLOADS: [&str; 2] = [
"org.apache.hudi.common.model.DefaultHoodieRecordPayload",
"org.apache.hudi.common.model.EventTimeAvroPayload",
];
const COMMIT_TIME_PAYLOAD: &str = "org.apache.hudi.common.model.OverwriteWithLatestAvroPayload";
#[derive(Debug, PartialEq, Eq)]
enum InferredMode {
CommitTime,
EventTime,
Custom,
}
impl InferredMode {
fn as_hudi_name(&self) -> &'static str {
match self {
Self::CommitTime => "COMMIT_TIME_ORDERING",
Self::EventTime => "EVENT_TIME_ORDERING",
Self::Custom => "CUSTOM",
}
}
}
fn mode_from_payload_class(payload_class: &str) -> Option<InferredMode> {
if payload_class.is_empty() {
return None;
}
if EVENT_TIME_PAYLOADS.contains(&payload_class) {
Some(InferredMode::EventTime)
} else if payload_class == COMMIT_TIME_PAYLOAD {
Some(InferredMode::CommitTime)
} else {
Some(InferredMode::Custom)
}
}
fn mode_from_strategy_id(strategy_id: &str) -> Option<InferredMode> {
if strategy_id.is_empty() {
return None;
}
match strategy_id {
EVENT_TIME_STRATEGY_ID => Some(InferredMode::EventTime),
COMMIT_TIME_STRATEGY_ID => Some(InferredMode::CommitTime),
_ => Some(InferredMode::Custom),
}
}
fn infer_merge_mode(hudi_configs: &HudiConfigs) -> Result<InferredMode> {
let options = hudi_configs.as_options();
let payload_class = PAYLOAD_CLASS_KEYS
.iter()
.find_map(|key| options.get(*key))
.map(String::as_str)
.unwrap_or_default();
let strategy_id = RECORD_MERGE_STRATEGY_ID_KEYS
.iter()
.find_map(|key| options.get(*key))
.map(String::as_str)
.unwrap_or_default();
if payload_class.is_empty() && strategy_id.is_empty() {
let has_ordering_field = hudi_configs
.try_get(HudiTableConfig::OrderingFields)?
.map(|v| -> Vec<String> { v.into() })
.is_some_and(|fields| fields.iter().any(|f| !f.trim().is_empty()));
let mode = if has_ordering_field {
InferredMode::EventTime
} else {
InferredMode::CommitTime
};
let ordering_desc = if has_ordering_field {
"having"
} else {
"not having"
};
log::debug!(
"merge mode {}: the table names no payload class or merge strategy, so it was \
inferred from {ordering_desc} an ordering field",
mode.as_hudi_name()
);
return Ok(mode);
}
let from_payload = mode_from_payload_class(payload_class);
let from_strategy = mode_from_strategy_id(strategy_id);
let table_version: isize = hudi_configs
.try_get(HudiTableConfig::TableVersion)?
.map(|v| v.into())
.unwrap_or(6);
let (inferred, source) = if table_version >= 8 {
match from_strategy {
Some(mode) => (Some(mode), "merge strategy"),
None => (from_payload, "payload class"),
}
} else {
match from_payload {
Some(mode) => (Some(mode), "payload class"),
None => (from_strategy, "merge strategy"),
}
};
if let Some(ref mode) = inferred {
log::debug!(
"merge mode {}: inferred from the {source} of a version {table_version} table \
(payload class {payload_class:?}, merge strategy {strategy_id:?})",
mode.as_hudi_name()
);
}
inferred.ok_or_else(|| {
CoreError::Unsupported(format!(
"Cannot infer a merge mode from payload class '{payload_class}' \
or merge strategy '{strategy_id}'."
))
})
}
fn resolve_merge_mode(hudi_configs: &HudiConfigs) -> Result<MergeMode> {
if let Some(mode) = hudi_configs.as_options().get(RECORD_MERGE_MODE) {
log::debug!("merge mode {mode}: stated by the table");
return match mode.to_ascii_uppercase().as_str() {
"COMMIT_TIME_ORDERING" => Ok(MergeMode::CommitTimeOrdering),
"EVENT_TIME_ORDERING" => Ok(MergeMode::EventTimeOrdering),
"CUSTOM" if resolve_custom_merger(&hudi_configs.as_options()).is_some() => {
Ok(MergeMode::Custom)
}
other => Err(CoreError::Unsupported(format!(
"Record merge mode '{other}' is not supported. Set \
hoodie.read.file.group.reader.version=1 to read with the \
reader that served this table before"
))),
};
}
match infer_merge_mode(hudi_configs)? {
InferredMode::CommitTime => Ok(MergeMode::CommitTimeOrdering),
InferredMode::EventTime => Ok(MergeMode::EventTimeOrdering),
InferredMode::Custom if resolve_custom_merger(&hudi_configs.as_options()).is_some() => {
Ok(MergeMode::Custom)
}
InferredMode::Custom => Err(CoreError::Unsupported(
"This table merges with a merger of its own, which the merge-on-read \
reader cannot reproduce. Set hoodie.read.file.group.reader.version=1 \
to read with the reader that served it before, which merges without \
that merger"
.to_string(),
)),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::internal::HudiInternalConfig;
use crate::config::plan::HudiPlanConfig;
use crate::file_group::reader::FileGroupReader;
use crate::file_group::reader_v2::reader_parameters::ReaderParameters;
use std::sync::Arc;
fn minimal_configs() -> Vec<(String, String)> {
vec![
(
HudiTableConfig::BasePath.as_ref().to_string(),
"file:///tmp/t".to_string(),
),
(
HudiReadConfig::EndTimestamp.as_ref().to_string(),
"20240101000000000".to_string(),
),
(
HudiTableConfig::OrderingFields.as_ref().to_string(),
"ts".to_string(),
),
]
}
fn configs_without(key: &str) -> HudiConfigs {
HudiConfigs::new(
minimal_configs()
.into_iter()
.filter(|(k, _)| k != key)
.collect::<Vec<_>>(),
)
}
#[test]
fn instant_range_matches_file_group_reader_version_one() {
let mut options = minimal_configs();
options.push((
HudiReadConfig::StartTimestamp.as_ref().to_string(),
"20230101000000000".to_string(),
));
let configs = HudiConfigs::new(options);
let version_one = FileGroupReader::new_with_overrides(
Arc::new(configs.clone()),
HashMap::new(),
HashMap::new(),
)
.unwrap()
.create_instant_range_for_log_file_scan()
.unwrap();
let resolved = resolve_reader_context(&configs, true, None)
.unwrap()
.instant_range
.expect("the resolver always derives a range");
assert_eq!(format!("{version_one:?}"), format!("{resolved:?}"));
}
#[test]
fn resolves_table_path_from_base_path_config() {
let configs = HudiConfigs::new(minimal_configs());
let ctx = resolve_reader_context(&configs, false, None).unwrap();
assert_eq!(ctx.table_path, "file:///tmp/t");
}
#[test]
fn resolves_a_base_only_slice_of_a_custom_merge_table() {
let mut options = minimal_configs();
options.push(("hoodie.record.merge.mode".to_string(), "CUSTOM".to_string()));
let ctx =
resolve_reader_context(&HudiConfigs::new(options), false, None).unwrap();
assert_eq!(ctx.merge_mode, "COMMIT_TIME_ORDERING");
}
#[test]
fn refuses_a_merging_slice_of_a_custom_merge_table_naming_the_way_back() {
let mut options = minimal_configs();
options.push(("hoodie.record.merge.mode".to_string(), "CUSTOM".to_string()));
let err =
resolve_reader_context(&HudiConfigs::new(options), true, None).unwrap_err();
assert!(
err.to_string()
.contains("hoodie.read.file.group.reader.version=1"),
"the error must name the way back, got: {err}"
);
}
#[test]
fn custom_payload_class_follows_the_same_no_merge_rule() {
let mut options = minimal_configs();
options.push((
"hoodie.compaction.payload.class".to_string(),
"com.example.MyPayload".to_string(),
));
let configs = HudiConfigs::new(options);
let ctx = resolve_reader_context(&configs, false, None).unwrap();
assert_eq!(ctx.merge_mode, "COMMIT_TIME_ORDERING");
let err = resolve_reader_context(&configs, true, None).unwrap_err();
assert!(
err.to_string()
.contains("hoodie.read.file.group.reader.version=1"),
"the error must name the way back, got: {err}"
);
}
#[test]
fn resolves_latest_commit_time_from_end_timestamp() {
let configs = HudiConfigs::new(minimal_configs());
let ctx = resolve_reader_context(&configs, false, None).unwrap();
assert_eq!(ctx.latest_commit_time, "20240101000000000");
}
#[test]
fn errors_when_end_timestamp_is_absent() {
let configs = configs_without(HudiReadConfig::EndTimestamp.as_ref());
let err = resolve_reader_context(&configs, false, None).unwrap_err();
assert!(
err.to_string()
.contains(HudiReadConfig::EndTimestamp.as_ref()),
"error should name the missing key, got: {err}"
);
}
#[test]
fn resolves_has_log_files_from_the_caller() {
let configs = HudiConfigs::new(minimal_configs());
assert!(
resolve_reader_context(&configs, true, None)
.unwrap()
.has_log_files
);
assert!(
!resolve_reader_context(&configs, false, None)
.unwrap()
.has_log_files
);
}
#[test]
fn defaults_base_file_format_to_parquet() {
let configs = HudiConfigs::new(minimal_configs());
let ctx = resolve_reader_context(&configs, true, None).unwrap();
assert_eq!(ctx.base_file_format, BaseFileFormatValue::Parquet.as_ref());
}
#[test]
fn resolves_base_file_format_from_the_path_when_the_table_names_none() {
let configs = HudiConfigs::new(minimal_configs());
assert!(
BaseFileFormatValue::from_configs(&configs)
.unwrap()
.is_none(),
"the fixture must name no format, or the path is never consulted"
);
for (path, expected) in [
("part/f.lance", BaseFileFormatValue::Lance),
("part/f.hfile", BaseFileFormatValue::HFile),
("part/f.parquet", BaseFileFormatValue::Parquet),
("part/f.unknown", BaseFileFormatValue::Parquet),
] {
let ctx = resolve_reader_context(&configs, true, Some(path)).unwrap();
assert_eq!(
ctx.base_file_format,
expected.as_ref(),
"'{path}' must resolve to {}",
expected.as_ref()
);
}
}
#[test]
fn an_explicit_base_file_format_outranks_the_path() {
let mut options = minimal_configs();
options.push((
HudiTableConfig::BaseFileFormat.as_ref().to_string(),
"hfile".to_string(),
));
let configs = HudiConfigs::new(options);
let ctx = resolve_reader_context(&configs, true, Some("part/f.lance")).unwrap();
assert_eq!(ctx.base_file_format, BaseFileFormatValue::HFile.as_ref());
}
#[test]
fn bounds_instant_range_by_the_read_window() {
let mut options = minimal_configs();
options.push((
HudiReadConfig::StartTimestamp.as_ref().to_string(),
"20230101000000000".to_string(),
));
let configs = HudiConfigs::new(options);
let ctx = resolve_reader_context(&configs, true, None).unwrap();
let rendered = format!("{:?}", ctx.instant_range.expect("always derived"));
assert!(rendered.contains("20230101000000000"), "got: {rendered}");
assert!(rendered.contains("20240101000000000"), "got: {rendered}");
assert!(
rendered.contains("start_inclusive: false"),
"incremental reads are exclusive at the start, got: {rendered}"
);
assert!(
rendered.contains("end_inclusive: true"),
"the pinned commit is included, got: {rendered}"
);
}
#[test]
fn routes_the_non_read_prefixed_per_read_keys_to_the_reader_config() {
let mut options = minimal_configs();
let expected = [
(spillable_map::CONFIG_MERGE_MAX_SIZE, "104857600"),
(spillable_map::CONFIG_MAX_PEAK_MEMORY, "209715200"),
(spillable_map::CONFIG_SPILLABLE_MAP_PATH, "/scratch/spill"),
(spillable_map::CONFIG_DISKMAP_TYPE, "ROCKS_DB"),
(CONFIG_MERGE_TYPE, "skip_merge"),
];
for (key, value) in expected {
options.push((key.to_string(), value.to_string()));
}
let configs = HudiConfigs::new(options);
let ctx = resolve_reader_context(&configs, true, None).unwrap();
for (key, value) in expected {
assert_eq!(
ctx.hoodie_reader_config.get(key),
Some(&value.to_string()),
"{key} is a per-read override and must reach the reader config"
);
assert!(
!ctx.table_config.contains_key(key),
"{key} is not something the table declares about itself"
);
}
let spill = spillable_map::SpillConfig::from_config(&ctx.hoodie_reader_config).unwrap();
assert_eq!(
spill.max_in_memory_size,
(104857600.0 * spillable_map::SPILL_TRIGGER_FRACTION) as u64,
"the configured merge budget must be honored"
);
assert_eq!(spill.spill_path, std::path::PathBuf::from("/scratch/spill"));
assert_eq!(spill.max_peak_in_memory_size, Some(209715200));
}
#[test]
fn separates_table_configs_from_read_configs() {
let configs = HudiConfigs::new(minimal_configs());
let ctx = resolve_reader_context(&configs, true, None).unwrap();
assert_eq!(
ctx.table_config
.get(HudiTableConfig::OrderingFields.as_ref()),
Some(&"ts".to_string())
);
assert!(
!ctx.table_config
.contains_key(HudiReadConfig::EndTimestamp.as_ref()),
"read configs must not leak into table config"
);
assert_eq!(
ctx.hoodie_reader_config
.get(HudiReadConfig::EndTimestamp.as_ref()),
Some(&"20240101000000000".to_string())
);
assert!(
!ctx.hoodie_reader_config
.contains_key(HudiTableConfig::OrderingFields.as_ref()),
"table configs must not leak into reader config"
);
}
#[test]
fn drops_configs_that_are_neither_table_nor_read() {
let crate_configs = [
HudiInternalConfig::SkipConfigValidation.as_ref(),
HudiInternalConfig::TimelineArchivedReadEnabled.as_ref(),
HudiPlanConfig::ListingParallelism.as_ref(),
];
let mut options = minimal_configs();
for key in crate_configs {
options.push((key.to_string(), "1".to_string()));
}
let configs = HudiConfigs::new(options);
let ctx = resolve_reader_context(&configs, true, None).unwrap();
for key in crate_configs {
assert!(
!ctx.table_config.contains_key(key),
"{key} is not a table property"
);
assert!(
!ctx.hoodie_reader_config.contains_key(key),
"{key} is not a per-read override"
);
}
}
#[test]
fn leaves_position_based_merge_off() {
let configs = HudiConfigs::new(minimal_configs());
let ctx = resolve_reader_context(&configs, true, None).unwrap();
assert!(!ctx.should_merge_use_record_position);
}
#[test]
fn turns_position_based_merge_on_from_hudis_own_config_key() {
let mut options = minimal_configs();
options.push((
HudiReadConfig::MergeUseRecordPositions.as_ref().to_string(),
"true".to_string(),
));
let ctx =
resolve_reader_context(&HudiConfigs::new(options), true, None).unwrap();
assert!(ctx.should_merge_use_record_position);
}
#[test]
fn rejects_a_non_boolean_position_merge_setting() {
let mut options = minimal_configs();
options.push((
HudiReadConfig::MergeUseRecordPositions.as_ref().to_string(),
"yes".to_string(),
));
let err =
resolve_reader_context(&HudiConfigs::new(options), true, None).unwrap_err();
assert!(
err.to_string()
.contains("hoodie.merge.use.record.positions"),
"the error should name the offending key, got: {err}"
);
}
#[test]
fn defaults_reader_parameters_to_version_one_behavior() {
let params = ReaderParameters::default();
assert!(!params.emit_delete);
assert!(!params.sort_output);
assert!(!params.allow_inflight_instants);
}
#[test]
fn resolves_event_time_ordering_from_an_ordering_field() {
let configs = HudiConfigs::new(minimal_configs());
let ctx = resolve_reader_context(&configs, true, None).unwrap();
assert_eq!(ctx.merge_mode, MergeMode::EventTimeOrdering.as_ref());
}
#[test]
fn infers_commit_time_ordering_without_an_ordering_field() {
let configs = configs_without(HudiTableConfig::OrderingFields.as_ref());
let ctx = resolve_reader_context(&configs, true, None).unwrap();
assert_eq!(ctx.merge_mode, MergeMode::CommitTimeOrdering.as_ref());
}
#[test]
fn infers_event_time_ordering_for_a_virtual_key_table() {
let mut options = minimal_configs();
options.push((
HudiTableConfig::PopulatesMetaFields.as_ref().to_string(),
"false".to_string(),
));
let configs = HudiConfigs::new(options);
let ctx = resolve_reader_context(&configs, true, None).unwrap();
assert_eq!(ctx.merge_mode, MergeMode::EventTimeOrdering.as_ref());
}
#[test]
fn infers_the_mode_java_infers() {
let cases: Vec<(&str, &str, Option<&str>, &str, InferredMode)> = vec![
("", "", None, "6", InferredMode::CommitTime),
("", "", Some("ts"), "6", InferredMode::EventTime),
(
"org.apache.hudi.common.model.DefaultHoodieRecordPayload",
"",
None,
"6",
InferredMode::EventTime,
),
(
"org.apache.hudi.common.model.EventTimeAvroPayload",
"",
None,
"6",
InferredMode::EventTime,
),
(
"org.apache.hudi.common.model.OverwriteWithLatestAvroPayload",
"",
None,
"6",
InferredMode::CommitTime,
),
("com.example.MyPayload", "", None, "6", InferredMode::Custom),
(
"",
EVENT_TIME_STRATEGY_ID,
None,
"8",
InferredMode::EventTime,
),
(
"",
COMMIT_TIME_STRATEGY_ID,
None,
"8",
InferredMode::CommitTime,
),
(
"",
"00000000-0000-0000-0000-000000000000",
None,
"8",
InferredMode::Custom,
),
(
"org.apache.hudi.common.model.OverwriteWithLatestAvroPayload",
EVENT_TIME_STRATEGY_ID,
None,
"8",
InferredMode::EventTime,
),
(
"org.apache.hudi.common.model.OverwriteWithLatestAvroPayload",
EVENT_TIME_STRATEGY_ID,
None,
"6",
InferredMode::CommitTime,
),
];
for (payload, strategy, ordering, version, expected) in cases {
let mut options = vec![
(
HudiTableConfig::BasePath.as_ref().to_string(),
"file:///tmp/t".to_string(),
),
(
HudiTableConfig::TableVersion.as_ref().to_string(),
version.to_string(),
),
];
if !payload.is_empty() {
options.push((PAYLOAD_CLASS_KEYS[0].to_string(), payload.to_string()));
}
if !strategy.is_empty() {
options.push((
RECORD_MERGE_STRATEGY_ID_KEYS[0].to_string(),
strategy.to_string(),
));
}
if let Some(field) = ordering {
options.push((
HudiTableConfig::OrderingFields.as_ref().to_string(),
field.to_string(),
));
}
let inferred = infer_merge_mode(&HudiConfigs::new(options)).unwrap();
assert_eq!(
inferred, expected,
"payload={payload:?} strategy={strategy:?} ordering={ordering:?} version={version}"
);
}
}
#[tokio::test]
async fn resolves_a_v6_table_that_was_refused_before() {
use hudi_test::SampleTable;
let url = SampleTable::V6SimplekeygenHivestyleNoMetafields.url_to_mor_parquet();
let table = crate::table::Table::new(url.as_ref()).await.unwrap();
let mut options = table.hudi_configs.as_options();
options.insert(
"hoodie.read.end.timestamp".to_string(),
"99991231235959999".to_string(),
);
let ctx =
resolve_reader_context(&HudiConfigs::new(options), true, None).unwrap();
assert_eq!(ctx.merge_mode, MergeMode::CommitTimeOrdering.as_ref());
}
#[test]
fn refuses_a_table_with_its_own_merger() {
let mut options = minimal_configs();
options.push((
PAYLOAD_CLASS_KEYS[0].to_string(),
"com.example.MyPayload".to_string(),
));
let configs = HudiConfigs::new(options);
let err = resolve_reader_context(&configs, true, None).unwrap_err();
assert!(
err.to_string().contains("merger of its own"),
"error should say why, got: {err}"
);
}
}