pub(crate) mod backend;
#[cfg(feature = "datafusion")]
pub mod datafusion;
pub(crate) mod deletion_vector;
pub(crate) mod metrics;
mod options;
pub(crate) mod partition_target;
pub(crate) mod planning;
pub(crate) mod predicate;
#[allow(dead_code)]
pub(crate) mod scheduling;
pub(crate) mod transform;
#[doc(hidden)]
pub use metrics::ParquetRangePlanningDiagnosticSnapshot;
pub use metrics::{DeltaScanMetrics, DeltaScanMetricsSnapshot};
#[doc(hidden)]
pub use options::ParquetRangeReadPolicy;
pub use options::{
DeltaScanExecutionOptions, DeltaSnapshotSelection, DeltaStorageOptions, ParquetReaderBackend,
};
pub use predicate::{DeltaComparison, DeltaPredicate, DeltaScalar};
use std::{
collections::VecDeque,
fmt,
pin::Pin,
sync::Arc,
task::{Context, Poll},
};
#[cfg(feature = "experimental-parquet-metadata-preparation")]
use std::time::{Duration, Instant};
use arrow::{datatypes::SchemaRef, record_batch::RecordBatch};
use futures_util::Stream;
use snafu::ResultExt;
#[cfg(any(
feature = "datafusion",
feature = "experimental-parquet-metadata-preparation"
))]
use self::backend::direct_parquet::ParquetMetadataCache;
#[cfg(feature = "experimental-parquet-metadata-preparation")]
use self::backend::direct_parquet::prepare_parquet_metadata;
use self::{
planning::{DeltaScanPartitionTargetOptions, DeltaScanPlan, plan_scan},
predicate::{evaluate_predicate, referenced_columns, validate_predicate},
scheduling::{
DeltaScanScheduler, FileAdmissionDecision, FileAdmissionPolicy, FileBatchStream,
FileExecutor, PartitionStream,
},
};
use crate::{
DeltaProtocol, DeltaReaderError,
delta::{
kernel::{kernel_pruning_is_exact, kernel_pruning_predicate},
protocol::validate_protocol,
snapshot::{
ArrowTableSnapshot, KernelTableSnapshot, load_delta_table_snapshot,
load_kernel_table_snapshot,
},
},
error::{DataFileReadSnafu, InvalidConfigurationSnafu, ScanPlanningSnafu},
};
const TRACING_TARGET: &str = "delta_arrow_reader";
const DELTA_LOG_SCAN_METADATA_SOURCE: &str = "delta_log";
const EAGER_CACHE_SCAN_METADATA_SOURCE: &str = "eager_cache";
#[cfg(feature = "experimental-parquet-metadata-preparation")]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ParquetMetadataPreparationLimits {
pub max_files: usize,
pub max_retained_metadata_bytes: usize,
}
#[cfg(feature = "experimental-parquet-metadata-preparation")]
#[non_exhaustive]
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ParquetMetadataPreparationReport {
pub files_prepared: usize,
pub estimated_retained_metadata_bytes: usize,
pub preparation_duration: Duration,
pub read_metrics: DeltaScanMetricsSnapshot,
}
#[cfg(feature = "experimental-parquet-metadata-preparation")]
struct PreparedParquetMetadata {
cache: Arc<ParquetMetadataCache>,
report: ParquetMetadataPreparationReport,
}
#[must_use = "table builder settings do nothing unless the builder is loaded"]
pub struct DeltaTableBuilder {
table_location: String,
storage_options: DeltaStorageOptions,
snapshot_selection: DeltaSnapshotSelection,
execution_options: DeltaScanExecutionOptions,
}
impl DeltaTableBuilder {
pub fn new(table_location: impl Into<String>) -> Self {
Self {
table_location: table_location.into(),
storage_options: DeltaStorageOptions::new(),
snapshot_selection: DeltaSnapshotSelection::Latest,
execution_options: DeltaScanExecutionOptions::new(),
}
}
pub fn with_storage_options(mut self, storage_options: DeltaStorageOptions) -> Self {
self.storage_options = storage_options;
self
}
pub const fn with_snapshot_selection(
mut self,
snapshot_selection: DeltaSnapshotSelection,
) -> Self {
self.snapshot_selection = snapshot_selection;
self
}
pub const fn with_execution_options(
mut self,
execution_options: DeltaScanExecutionOptions,
) -> Self {
self.execution_options = execution_options;
self
}
pub async fn load_table(self) -> Result<DeltaTable, DeltaReaderError> {
let snapshot = load_delta_table_snapshot(
self.table_location,
self.storage_options,
self.snapshot_selection,
)
.await?;
Ok(DeltaTable::new(snapshot, self.execution_options))
}
pub async fn load_table_with_eager_scan_metadata(self) -> Result<DeltaTable, DeltaReaderError> {
let snapshot = load_delta_table_snapshot(
self.table_location,
self.storage_options,
self.snapshot_selection,
)
.await?;
validate_protocol(snapshot.protocol())?;
let snapshot = materialize_eager_scan_metadata(snapshot).await?;
Ok(DeltaTable::new(snapshot, self.execution_options))
}
#[cfg(feature = "experimental-parquet-metadata-preparation")]
pub async fn load_table_with_prepared_parquet_metadata(
self,
limits: ParquetMetadataPreparationLimits,
) -> Result<DeltaTable, DeltaReaderError> {
limits.validate()?;
if self.execution_options.parquet_backend() != ParquetReaderBackend::Direct {
return InvalidConfigurationSnafu {
reason: "parquet_metadata_preparation_requires_direct_backend",
}
.fail();
}
self.load_table_with_eager_scan_metadata()
.await?
.prepare_parquet_metadata(limits)
.await
}
pub async fn load_snapshot(self) -> Result<DeltaTableSnapshot, DeltaReaderError> {
let snapshot = load_kernel_table_snapshot(
self.table_location,
self.storage_options,
self.snapshot_selection,
)
.await?;
Ok(DeltaTableSnapshot::new(snapshot, self.execution_options))
}
}
async fn materialize_eager_scan_metadata(
snapshot: ArrowTableSnapshot,
) -> Result<ArrowTableSnapshot, DeltaReaderError> {
tokio::task::spawn_blocking(move || snapshot.materialize_eager_scan_metadata())
.await
.boxed()
.context(ScanPlanningSnafu {
reason: "eager_scan_metadata_task_failed",
})?
}
#[cfg(feature = "experimental-parquet-metadata-preparation")]
impl ParquetMetadataPreparationLimits {
fn validate(self) -> Result<(), DeltaReaderError> {
if self.max_files == 0 {
return InvalidConfigurationSnafu {
reason: "parquet_metadata_preparation_max_files_must_be_positive",
}
.fail();
}
if self.max_retained_metadata_bytes == 0 {
return InvalidConfigurationSnafu {
reason: "parquet_metadata_preparation_memory_limit_must_be_positive",
}
.fail();
}
Ok(())
}
}
impl fmt::Debug for DeltaTableBuilder {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("DeltaTableBuilder")
.field("table_location", &"<redacted>")
.field("storage_options", &"<redacted>")
.field("snapshot_selection", &self.snapshot_selection)
.field("execution_options", &self.execution_options)
.finish()
}
}
pub struct DeltaTableSnapshot {
snapshot: KernelTableSnapshot,
execution_options: DeltaScanExecutionOptions,
}
impl DeltaTableSnapshot {
fn new(snapshot: KernelTableSnapshot, execution_options: DeltaScanExecutionOptions) -> Self {
Self {
snapshot,
execution_options,
}
}
pub fn version(&self) -> u64 {
self.snapshot.version()
}
pub fn protocol(&self) -> &DeltaProtocol {
self.snapshot.protocol()
}
pub fn table_url(&self) -> &str {
self.snapshot.table_url()
}
pub fn validate_protocol(&self) -> Result<(), DeltaReaderError> {
validate_protocol(self.protocol())
}
pub fn into_table(self) -> Result<DeltaTable, DeltaReaderError> {
Ok(DeltaTable::new(
self.snapshot.into_arrow_snapshot()?,
self.execution_options,
))
}
}
impl fmt::Debug for DeltaTableSnapshot {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("DeltaTableSnapshot")
.field("version", &self.version())
.finish_non_exhaustive()
}
}
#[derive(Clone)]
pub struct DeltaTable {
snapshot: Arc<ArrowTableSnapshot>,
execution_options: DeltaScanExecutionOptions,
#[cfg(feature = "experimental-parquet-metadata-preparation")]
prepared_parquet_metadata: Option<Arc<PreparedParquetMetadata>>,
}
impl DeltaTable {
fn new(snapshot: ArrowTableSnapshot, execution_options: DeltaScanExecutionOptions) -> Self {
Self {
snapshot: Arc::new(snapshot),
execution_options,
#[cfg(feature = "experimental-parquet-metadata-preparation")]
prepared_parquet_metadata: None,
}
}
#[cfg(feature = "experimental-parquet-metadata-preparation")]
async fn prepare_parquet_metadata(
mut self,
limits: ParquetMetadataPreparationLimits,
) -> Result<Self, DeltaReaderError> {
let snapshot = Arc::clone(&self.snapshot);
let execution_options = self.execution_options;
let plan = tokio::task::spawn_blocking(move || {
plan_scan(
snapshot.as_ref(),
None,
&[],
None,
false,
execution_options,
DeltaScanPartitionTargetOptions {
explicit_target_partitions: None,
datafusion_target_partitions: None,
},
)
})
.await
.boxed()
.context(ScanPlanningSnafu {
reason: "parquet_metadata_preparation_planning_task_failed",
})??;
let metadata_cache = Arc::new(ParquetMetadataCache::default());
let started = Instant::now();
let prepared =
prepare_parquet_metadata(&plan, Arc::clone(&metadata_cache), Arc::default(), limits)
.await?;
self.prepared_parquet_metadata = Some(Arc::new(PreparedParquetMetadata {
cache: metadata_cache,
report: ParquetMetadataPreparationReport {
files_prepared: prepared.file_count,
estimated_retained_metadata_bytes: prepared.memory_bytes,
preparation_duration: started.elapsed(),
read_metrics: plan.metrics.snapshot(),
},
}));
Ok(self)
}
pub fn version(&self) -> u64 {
self.snapshot.version()
}
pub fn schema(&self) -> SchemaRef {
self.snapshot.schema()
}
pub fn protocol(&self) -> &DeltaProtocol {
self.snapshot.protocol()
}
pub fn table_url(&self) -> &str {
self.snapshot.table_url()
}
#[allow(dead_code)]
pub(crate) fn partition_columns(&self) -> &[String] {
self.snapshot.partition_columns()
}
#[allow(dead_code)]
pub(crate) fn snapshot(&self) -> &ArrowTableSnapshot {
self.snapshot.as_ref()
}
pub fn validate_protocol(&self) -> Result<(), DeltaReaderError> {
validate_protocol(self.protocol())
}
#[cfg(feature = "experimental-parquet-metadata-preparation")]
pub fn parquet_metadata_preparation_report(&self) -> Option<&ParquetMetadataPreparationReport> {
self.prepared_parquet_metadata
.as_deref()
.map(|prepared| &prepared.report)
}
#[cfg(feature = "datafusion")]
pub(crate) fn prepared_parquet_metadata_cache(&self) -> Option<Arc<ParquetMetadataCache>> {
#[cfg(feature = "experimental-parquet-metadata-preparation")]
{
self.prepared_parquet_metadata
.as_ref()
.map(|prepared| Arc::clone(&prepared.cache))
}
#[cfg(not(feature = "experimental-parquet-metadata-preparation"))]
{
None
}
}
pub fn scan(&self) -> DeltaScanBuilder<'_> {
DeltaScanBuilder {
table: self,
projection: None,
predicate: None,
limit: None,
target_partitions: None,
execution_options: self.execution_options,
}
}
}
impl fmt::Debug for DeltaTable {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("DeltaTable")
.field("version", &self.version())
.finish_non_exhaustive()
}
}
#[must_use = "scan builder settings do nothing unless the scan is built"]
pub struct DeltaScanBuilder<'table> {
table: &'table DeltaTable,
projection: Option<Vec<String>>,
predicate: Option<DeltaPredicate>,
limit: Option<usize>,
target_partitions: Option<usize>,
execution_options: DeltaScanExecutionOptions,
}
impl<'table> DeltaScanBuilder<'table> {
pub fn with_projection(
mut self,
logical_columns: impl IntoIterator<Item = impl Into<String>>,
) -> Self {
self.projection = Some(logical_columns.into_iter().map(Into::into).collect());
self
}
pub fn with_predicate(mut self, predicate: DeltaPredicate) -> Self {
self.predicate = Some(predicate);
self
}
pub const fn with_limit(mut self, limit: usize) -> Self {
self.limit = Some(limit);
self
}
pub fn with_target_partitions(
mut self,
target_partitions: usize,
) -> Result<Self, DeltaReaderError> {
if target_partitions == 0 {
return InvalidConfigurationSnafu {
reason: "scan_partition_target_must_be_positive",
}
.fail();
}
self.target_partitions = Some(target_partitions);
Ok(self)
}
pub const fn with_execution_options(
mut self,
execution_options: DeltaScanExecutionOptions,
) -> Self {
self.execution_options = execution_options;
self
}
pub async fn build(self) -> Result<DeltaScan, DeltaReaderError> {
self.table.validate_protocol()?;
if let Some(predicate) = self.predicate.as_ref() {
validate_predicate(predicate, self.table.schema().as_ref())?;
}
let snapshot_version = self.table.version();
let backend = self.execution_options.parquet_backend();
let scan_metadata_source = if self.table.snapshot.eager_scan_metadata().is_some() {
EAGER_CACHE_SCAN_METADATA_SOURCE
} else {
DELTA_LOG_SCAN_METADATA_SOURCE
};
trace_planning_started(snapshot_version, backend, scan_metadata_source);
let snapshot = Arc::clone(&self.table.snapshot);
let projection = self.projection;
let predicate = self.predicate;
let hidden_columns = predicate
.as_ref()
.map(referenced_columns)
.unwrap_or_default();
let enforce_physical_predicate_rows =
predicate.as_ref().is_some_and(kernel_pruning_is_exact);
let kernel_predicate = predicate.as_ref().and_then(kernel_pruning_predicate);
let include_stats = kernel_predicate.is_some();
let execution_options = self.execution_options;
let target_partitions = self.target_partitions;
let result = tokio::task::spawn_blocking(move || {
plan_scan(
snapshot.as_ref(),
projection.as_deref(),
&hidden_columns,
kernel_predicate,
include_stats,
execution_options,
DeltaScanPartitionTargetOptions {
explicit_target_partitions: target_partitions,
datafusion_target_partitions: None,
},
)
})
.await
.boxed()
.context(ScanPlanningSnafu {
reason: "scan_planning_task_failed",
})
.and_then(|result| result);
match result {
Ok(plan) => {
trace_planning_completed(
snapshot_version,
backend,
plan.partitions.len(),
scan_metadata_source,
);
Ok(DeltaScan {
plan: Arc::new(plan),
predicate,
limit: self.limit,
enforce_physical_predicate_rows,
#[cfg(feature = "experimental-parquet-metadata-preparation")]
parquet_metadata_cache: self
.table
.prepared_parquet_metadata
.as_ref()
.map(|prepared| Arc::clone(&prepared.cache)),
})
}
Err(error) => {
trace_planning_failed(snapshot_version, backend, scan_metadata_source, &error);
Err(error)
}
}
}
}
#[must_use = "scans do nothing unless converted into a stream"]
pub struct DeltaScan {
plan: Arc<DeltaScanPlan>,
predicate: Option<DeltaPredicate>,
limit: Option<usize>,
enforce_physical_predicate_rows: bool,
#[cfg(feature = "experimental-parquet-metadata-preparation")]
parquet_metadata_cache: Option<Arc<ParquetMetadataCache>>,
}
impl DeltaScan {
pub fn schema(&self) -> SchemaRef {
Arc::clone(&self.plan.projected_schema)
}
pub fn partition_count(&self) -> usize {
self.plan.partitions.len()
}
pub fn into_stream(self) -> DeltaBatchStream {
let metrics = self.plan.metrics.clone();
let schema = Arc::clone(&self.plan.projected_schema);
let partition_count = self.plan.partitions.len();
let snapshot_version = self.plan.snapshot_version;
let backend = self.plan.execution_options.parquet_backend();
let projection = (self.plan.logical_schema.as_ref() != schema.as_ref())
.then(|| (0..schema.fields().len()).collect::<Vec<_>>());
let parquet_metadata_cache = {
#[cfg(feature = "experimental-parquet-metadata-preparation")]
{
self.parquet_metadata_cache
}
#[cfg(not(feature = "experimental-parquet-metadata-preparation"))]
{
None
}
};
let partitions = if self.limit == Some(0) {
VecDeque::new()
} else {
let scheduler = DeltaScanScheduler::new(Arc::clone(&self.plan));
let admission: FileAdmissionPolicy<_> = Arc::new(|_| Ok(FileAdmissionDecision::Admit));
let executor = match backend {
ParquetReaderBackend::Direct => direct_parquet_executor(
&self.plan,
None,
self.enforce_physical_predicate_rows
.then(|| self.plan.physical_predicate.clone())
.flatten(),
parquet_metadata_cache,
),
ParquetReaderBackend::DeltaKernel => delta_kernel_executor(&self.plan),
};
scheduler.partition_streams(admission, executor)
};
DeltaBatchStream {
schema,
metrics,
partitions,
predicate: self.predicate,
projection,
remaining: self.limit,
snapshot_version,
backend,
partition_count,
started: false,
done: false,
}
}
}
#[must_use = "streams do nothing unless polled"]
pub struct DeltaBatchStream {
schema: SchemaRef,
metrics: DeltaScanMetrics,
partitions: VecDeque<PartitionStream>,
predicate: Option<DeltaPredicate>,
projection: Option<Vec<usize>>,
remaining: Option<usize>,
snapshot_version: u64,
backend: ParquetReaderBackend,
partition_count: usize,
started: bool,
done: bool,
}
impl DeltaBatchStream {
pub fn schema(&self) -> SchemaRef {
Arc::clone(&self.schema)
}
pub fn metrics(&self) -> DeltaScanMetrics {
self.metrics.clone()
}
fn start(&mut self) {
if self.started {
return;
}
self.started = true;
trace_execution_started(self.snapshot_version, self.backend, self.partition_count);
for partition in &mut self.partitions {
partition.start();
}
}
fn complete(&mut self) {
if self.done {
return;
}
self.partitions.clear();
self.done = true;
trace_execution_completed(self.snapshot_version, self.backend, self.partition_count);
}
fn fail(&mut self, error: &DeltaReaderError) {
self.partitions.clear();
self.done = true;
trace_execution_failed(
self.snapshot_version,
self.backend,
self.partition_count,
error,
);
}
fn finalize_batch(&self, mut batch: RecordBatch) -> Result<RecordBatch, DeltaReaderError> {
if let Some(predicate) = self.predicate.as_ref() {
batch = evaluate_predicate(&batch, predicate)?;
}
if let Some(projection) = self.projection.as_ref() {
batch = batch
.project(projection)
.boxed()
.context(DataFileReadSnafu {
reason: "direct_projection_failed",
})?;
}
Ok(batch)
}
}
impl Stream for DeltaBatchStream {
type Item = Result<RecordBatch, DeltaReaderError>;
fn poll_next(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let this = self.get_mut();
if this.done {
return Poll::Ready(None);
}
this.start();
loop {
let Some(partition) = this.partitions.front_mut() else {
this.complete();
return Poll::Ready(None);
};
match Pin::new(partition).poll_next(context) {
Poll::Ready(Some(Ok(batch))) => {
let mut batch = match this.finalize_batch(batch) {
Ok(batch) => batch,
Err(error) => {
this.fail(&error);
return Poll::Ready(Some(Err(error)));
}
};
if let Some(remaining) = this.remaining.as_mut() {
if batch.num_rows() >= *remaining {
batch = batch.slice(0, *remaining);
*remaining = 0;
this.complete();
} else {
*remaining -= batch.num_rows();
}
}
return Poll::Ready(Some(Ok(batch)));
}
Poll::Ready(Some(Err(error))) => {
this.fail(&error);
return Poll::Ready(Some(Err(error)));
}
Poll::Ready(None) => {
this.partitions.pop_front();
}
Poll::Pending => return Poll::Pending,
}
}
}
}
impl Drop for DeltaBatchStream {
fn drop(&mut self) {
if self.done {
return;
}
self.partitions.clear();
self.done = true;
trace_execution_dropped(self.snapshot_version, self.backend, self.partition_count);
}
}
pub(crate) fn direct_parquet_executor(
plan: &Arc<DeltaScanPlan>,
output_batch_size_rows: Option<usize>,
row_predicate: Option<crate::delta::kernel::DeltaKernelPredicate>,
metadata_cache: Option<Arc<backend::direct_parquet::ParquetMetadataCache>>,
) -> FileExecutor<planning::DeltaScanFileTask, FileBatchStream> {
backend::direct_parquet::direct_parquet_file_executor(
plan,
output_batch_size_rows,
row_predicate,
Arc::default(),
metadata_cache,
)
}
pub(crate) fn delta_kernel_executor(
plan: &Arc<DeltaScanPlan>,
) -> FileExecutor<planning::DeltaScanFileTask, FileBatchStream> {
backend::kernel_reader::delta_kernel_file_executor(plan)
}
fn trace_planning_started(
snapshot_version: u64,
backend: ParquetReaderBackend,
scan_metadata_source: &'static str,
) {
tracing::debug!(
target: TRACING_TARGET,
event = "scan_planning.started",
snapshot_version,
backend = ?backend,
partition_count = tracing::field::Empty,
scan_metadata_source,
outcome = "started"
);
}
fn trace_planning_completed(
snapshot_version: u64,
backend: ParquetReaderBackend,
partition_count: usize,
scan_metadata_source: &'static str,
) {
tracing::debug!(
target: TRACING_TARGET,
event = "scan_planning.completed",
snapshot_version,
backend = ?backend,
partition_count,
scan_metadata_source,
outcome = "completed"
);
}
fn trace_planning_failed(
snapshot_version: u64,
backend: ParquetReaderBackend,
scan_metadata_source: &'static str,
error: &DeltaReaderError,
) {
tracing::debug!(
target: TRACING_TARGET,
event = "scan_planning.failed",
snapshot_version,
backend = ?backend,
partition_count = tracing::field::Empty,
scan_metadata_source,
outcome = "failed",
error_code = error.code(),
error_phase = error.phase().as_str()
);
}
fn trace_execution_started(
snapshot_version: u64,
backend: ParquetReaderBackend,
partition_count: usize,
) {
tracing::debug!(
target: TRACING_TARGET,
event = "scan_execution.started",
snapshot_version,
backend = ?backend,
partition_count,
outcome = "started"
);
}
fn trace_execution_completed(
snapshot_version: u64,
backend: ParquetReaderBackend,
partition_count: usize,
) {
tracing::debug!(
target: TRACING_TARGET,
event = "scan_execution.completed",
snapshot_version,
backend = ?backend,
partition_count,
outcome = "completed"
);
}
fn trace_execution_failed(
snapshot_version: u64,
backend: ParquetReaderBackend,
partition_count: usize,
error: &DeltaReaderError,
) {
tracing::debug!(
target: TRACING_TARGET,
event = "scan_execution.failed",
snapshot_version,
backend = ?backend,
partition_count,
outcome = "failed",
error_code = error.code(),
error_phase = error.phase().as_str()
);
}
fn trace_execution_dropped(
snapshot_version: u64,
backend: ParquetReaderBackend,
partition_count: usize,
) {
tracing::debug!(
target: TRACING_TARGET,
event = "scan_execution.dropped",
snapshot_version,
backend = ?backend,
partition_count,
outcome = "dropped"
);
}
#[cfg(test)]
mod tests {
use std::{
collections::{BTreeMap, VecDeque},
fmt, fs,
future::pending,
path::{Path, PathBuf},
sync::{Arc, Mutex, Once},
time::{Duration, SystemTime, UNIX_EPOCH},
};
use arrow::{
array::Int32Array,
datatypes::{DataType, Field, Schema, SchemaRef},
record_batch::RecordBatch,
};
use futures_util::{FutureExt, StreamExt, stream};
use tokio::{sync::Notify, time::timeout};
use tracing::{
Event, Level, Metadata, Subscriber,
field::{Field as TracingField, Visit},
span::{Attributes, Id, Record},
subscriber::{Interest, with_default},
};
use super::{
DeltaBatchStream, DeltaTable, trace_execution_completed, trace_execution_dropped,
trace_execution_failed, trace_execution_started, trace_planning_completed,
trace_planning_failed, trace_planning_started,
};
use crate::{
DeltaScanExecutionOptions, DeltaScanMetrics, DeltaSnapshotSelection, DeltaStorageOptions,
ParquetReaderBackend,
delta::snapshot::load_delta_table_snapshot_blocking,
error::InvalidConfigurationSnafu,
reader::{
metrics::DeltaScanMetricsConfig,
scheduling::{
FileAdmissionDecision, FileAdmissionPolicy, FileBatchStream, FileExecutor,
FileReadPermit, PartitionStream, ScanCancellation, ScanReadLimiter,
},
},
};
static TRACING_TEST_LOCK: Mutex<()> = Mutex::new(());
static TRACING_TEST_GLOBAL_SUBSCRIBER: Once = Once::new();
#[derive(Clone, Default)]
struct EventFields(Arc<Mutex<Vec<BTreeMap<String, String>>>>);
impl Subscriber for EventFields {
fn register_callsite(&self, metadata: &'static Metadata<'static>) -> Interest {
if metadata.target() == "delta_arrow_reader" && *metadata.level() == Level::DEBUG {
Interest::always()
} else {
Interest::sometimes()
}
}
fn enabled(&self, metadata: &Metadata<'_>) -> bool {
metadata.target() == "delta_arrow_reader" && *metadata.level() == Level::DEBUG
}
fn new_span(&self, _attributes: &Attributes<'_>) -> Id {
Id::from_u64(1)
}
fn record(&self, _span: &Id, _values: &Record<'_>) {}
fn record_follows_from(&self, _span: &Id, _follows: &Id) {}
fn event(&self, event: &Event<'_>) {
let metadata = event.metadata();
assert_eq!(metadata.target(), "delta_arrow_reader");
let mut fields = metadata
.fields()
.iter()
.map(|field| (field.name().to_owned(), "<empty>".to_owned()))
.collect();
event.record(&mut FieldVisitor(&mut fields));
self.0.lock().expect("event lock").push(fields);
}
fn enter(&self, _span: &Id) {}
fn exit(&self, _span: &Id) {}
}
struct FieldVisitor<'fields>(&'fields mut BTreeMap<String, String>);
impl Visit for FieldVisitor<'_> {
fn record_debug(&mut self, field: &TracingField, value: &dyn fmt::Debug) {
self.0.insert(field.name().to_owned(), format!("{value:?}"));
}
fn record_str(&mut self, field: &TracingField, value: &str) {
self.0.insert(field.name().to_owned(), value.to_owned());
}
fn record_u64(&mut self, field: &TracingField, value: u64) {
self.0.insert(field.name().to_owned(), value.to_string());
}
fn record_i64(&mut self, field: &TracingField, value: i64) {
self.0.insert(field.name().to_owned(), value.to_string());
}
fn record_u128(&mut self, field: &TracingField, value: u128) {
self.0.insert(field.name().to_owned(), value.to_string());
}
}
struct DeltaLogTable(PathBuf);
impl DeltaLogTable {
fn new(name: &str) -> Result<Self, Box<dyn std::error::Error>> {
let nanos = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos();
let path = Path::new("target")
.join("delta-arrow-reader-tracing-tests")
.join(format!("{}-{name}-{nanos}", std::process::id()));
fs::create_dir_all(path.join("_delta_log"))?;
fs::write(
path.join("_delta_log/00000000000000000000.json"),
r#"{"protocol":{"minReaderVersion":1,"minWriterVersion":2}}
{"metaData":{"id":"tracing-test","format":{"provider":"parquet","options":{}},"schemaString":"{\"type\":\"struct\",\"fields\":[{\"name\":\"id\",\"type\":\"integer\",\"nullable\":true,\"metadata\":{}}]}","partitionColumns":[],"configuration":{},"createdTime":1587968585495}}
{"add":{"path":"secret-planning-object.parquet","partitionValues":{},"size":10,"modificationTime":1587968586000,"dataChange":true}}
"#,
)?;
Ok(Self(path))
}
}
impl Drop for DeltaLogTable {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.0);
}
}
fn capture_events<T>(run: impl FnOnce() -> T) -> (T, Vec<BTreeMap<String, String>>) {
let _lock = TRACING_TEST_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let events = Arc::new(Mutex::new(Vec::new()));
let subscriber = EventFields(Arc::clone(&events));
TRACING_TEST_GLOBAL_SUBSCRIBER.call_once(|| {
let _ = tracing::subscriber::set_global_default(EventFields::default());
});
let result = with_default(subscriber, || {
tracing::callsite::rebuild_interest_cache();
run()
});
tracing::callsite::rebuild_interest_cache();
let captured = events
.lock()
.map(|events| events.clone())
.unwrap_or_default();
(result, captured)
}
struct ControlledMerge {
stream: DeltaBatchStream,
limiter: Arc<ScanReadLimiter>,
cancellation: ScanCancellation,
metrics: DeltaScanMetrics,
first_partition_gate: Arc<Notify>,
}
fn schema() -> SchemaRef {
Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]))
}
fn batch(id: i32) -> RecordBatch {
RecordBatch::try_new(schema(), vec![Arc::new(Int32Array::from(vec![id]))])
.expect("valid test batch")
}
fn batch_id(batch: &RecordBatch) -> i32 {
batch
.column(0)
.as_any()
.downcast_ref::<Int32Array>()
.expect("Int32 id")
.value(0)
}
fn execution_options() -> Result<DeltaScanExecutionOptions, crate::DeltaReaderError> {
DeltaScanExecutionOptions::new()
.with_prefetch_files_per_partition(0)
.with_max_concurrent_file_reads_per_partition(1)?
.with_max_concurrent_file_reads_per_scan(Some(2))?
.with_output_buffer_batches_per_partition(1)
}
fn metrics() -> DeltaScanMetrics {
DeltaScanMetrics::new(DeltaScanMetricsConfig {
snapshot_version: 7,
parquet_backend: ParquetReaderBackend::Direct,
scan_partitions_planned: 2,
files_planned: 2,
add_actions_excluded_during_planning: Some(0),
estimated_input_rows: Some(4),
estimated_input_bytes: Some(4),
})
}
fn file_stream(permit: FileReadPermit, batches: Vec<RecordBatch>) -> FileBatchStream {
Box::pin(stream::unfold(
(VecDeque::from(batches), permit),
|(mut batches, permit)| async move {
batches
.pop_front()
.map(|batch| (Ok(batch), (batches, permit)))
},
))
}
fn gated_file_stream(
permit: FileReadPermit,
batches: Vec<RecordBatch>,
gate: Arc<Notify>,
) -> FileBatchStream {
Box::pin(stream::unfold(
(false, VecDeque::from(batches), permit, gate),
|(wait, mut batches, permit, gate)| async move {
let batch = batches.pop_front()?;
if wait {
gate.notified().await;
}
Some((Ok(batch), (true, batches, permit, gate)))
},
))
}
fn direct_stream(
partitions: VecDeque<PartitionStream>,
metrics: DeltaScanMetrics,
) -> DeltaBatchStream {
DeltaBatchStream {
schema: schema(),
metrics,
partitions,
predicate: None,
projection: None,
remaining: None,
snapshot_version: 7,
backend: ParquetReaderBackend::Direct,
partition_count: 2,
started: false,
done: false,
}
}
fn controlled_merge() -> Result<ControlledMerge, Box<dyn std::error::Error>> {
let options = execution_options()?;
let limiter = ScanReadLimiter::new(options, 2, 2);
let cancellation = ScanCancellation::new();
let metrics = metrics();
let first_partition_gate = Arc::new(Notify::new());
let executor: FileExecutor<i32, FileBatchStream> = {
let gate = Arc::clone(&first_partition_gate);
Arc::new(move |task, permit, _| {
let gate = Arc::clone(&gate);
async move {
let batches = vec![batch(task), batch(task * 2)];
Ok(if task == 1 {
gated_file_stream(permit, batches, gate)
} else {
file_stream(permit, batches)
})
}
.boxed()
})
};
let admission: FileAdmissionPolicy<i32> =
Arc::new(|_: &i32| Ok(FileAdmissionDecision::Admit));
let first = PartitionStream::new(
vec![1],
limiter.partition(0)?,
options,
admission.clone(),
Arc::clone(&executor),
metrics.clone(),
cancellation.clone(),
);
let second = PartitionStream::new(
vec![10],
limiter.partition(1)?,
options,
admission,
executor,
metrics.clone(),
cancellation.clone(),
);
Ok(ControlledMerge {
stream: direct_stream(VecDeque::from([first, second]), metrics.clone()),
limiter,
cancellation,
metrics,
first_partition_gate,
})
}
async fn wait_for_batches(metrics: &DeltaScanMetrics, expected: u64) {
timeout(Duration::from_secs(5), async {
while metrics.snapshot().scheduler_batches_emitted < expected {
tokio::task::yield_now().await;
}
})
.await
.expect("batch production reached expected bound");
}
#[test]
fn lifecycle_tracing_has_only_bounded_fields() {
let error = InvalidConfigurationSnafu { reason: "test" }.build();
let (_, events) = capture_events(|| {
trace_planning_started(7, ParquetReaderBackend::Direct, "eager_cache");
trace_planning_completed(7, ParquetReaderBackend::Direct, 2, "eager_cache");
trace_planning_failed(7, ParquetReaderBackend::Direct, "eager_cache", &error);
trace_execution_started(7, ParquetReaderBackend::Direct, 2);
trace_execution_completed(7, ParquetReaderBackend::Direct, 2);
trace_execution_failed(7, ParquetReaderBackend::Direct, 2, &error);
trace_execution_dropped(7, ParquetReaderBackend::Direct, 2);
});
assert_eq!(events.len(), 7);
let allowed = [
"backend",
"error_phase",
"error_code",
"event",
"outcome",
"partition_count",
"scan_metadata_source",
"snapshot_version",
];
for fields in events.iter() {
assert!(fields.keys().all(|field| allowed.contains(&field.as_str())));
assert!(fields.contains_key("event"));
assert!(fields.contains_key("snapshot_version"));
assert!(fields.contains_key("backend"));
assert!(fields.contains_key("partition_count"));
assert!(fields.contains_key("outcome"));
if fields
.get("event")
.is_some_and(|event| event.starts_with("scan_planning."))
{
assert_eq!(
fields.get("scan_metadata_source").map(String::as_str),
Some("eager_cache")
);
}
}
}
#[test]
fn planning_tracing_reports_the_table_metadata_source_on_success_and_failure()
-> Result<(), Box<dyn std::error::Error>> {
const OBJECT_KEY: &str = "secret-planning-object.parquet";
const STORAGE_VALUE: &str = "secret-planning-storage-value";
let fixture = DeltaLogTable::new("planning-source")?;
let mut storage_options = DeltaStorageOptions::new();
storage_options.insert("secret-option".to_owned(), STORAGE_VALUE.to_owned());
let snapshot = load_delta_table_snapshot_blocking(
&fixture.0.to_string_lossy(),
&storage_options,
DeltaSnapshotSelection::Latest,
)?;
let eager_snapshot = snapshot.clone().materialize_eager_scan_metadata()?;
let lazy = DeltaTable::new(snapshot, DeltaScanExecutionOptions::new());
let eager = DeltaTable::new(eager_snapshot, DeltaScanExecutionOptions::new());
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
let (eager_result, eager_events) =
capture_events(|| runtime.block_on(eager.scan().build()));
let _ = eager_result?;
assert_eq!(eager_events.len(), 2);
assert_eq!(
eager_events[0].get("event").map(String::as_str),
Some("scan_planning.started")
);
assert_eq!(
eager_events[1].get("event").map(String::as_str),
Some("scan_planning.completed")
);
assert!(eager_events.iter().all(|event| {
event.get("scan_metadata_source").map(String::as_str) == Some("eager_cache")
}));
let (failed_result, failed_events) =
capture_events(|| runtime.block_on(eager.scan().with_projection(["missing"]).build()));
assert!(failed_result.is_err());
assert_eq!(failed_events.len(), 2);
assert_eq!(
failed_events[0].get("event").map(String::as_str),
Some("scan_planning.started")
);
assert_eq!(
failed_events[1].get("event").map(String::as_str),
Some("scan_planning.failed")
);
assert!(failed_events.iter().all(|event| {
event.get("scan_metadata_source").map(String::as_str) == Some("eager_cache")
}));
let (lazy_result, lazy_events) = capture_events(|| runtime.block_on(lazy.scan().build()));
let _ = lazy_result?;
assert_eq!(lazy_events.len(), 2);
assert!(lazy_events.iter().all(|event| {
event.get("scan_metadata_source").map(String::as_str) == Some("delta_log")
}));
let captured = format!("{eager_events:?}{failed_events:?}{lazy_events:?}");
assert!(!captured.contains(&fixture.0.to_string_lossy().into_owned()));
assert!(!captured.contains(OBJECT_KEY));
assert!(!captured.contains(STORAGE_VALUE));
Ok(())
}
#[tokio::test]
async fn merged_stream_is_ordered_and_bounds_later_partition_queues()
-> Result<(), Box<dyn std::error::Error>> {
let ControlledMerge {
mut stream,
limiter,
metrics,
first_partition_gate,
..
} = controlled_merge()?;
let first = stream.next().await.ok_or("first batch missing")??;
assert_eq!(batch_id(&first), 1);
wait_for_batches(&metrics, 2).await;
for _ in 0..32 {
tokio::task::yield_now().await;
}
assert_eq!(metrics.snapshot().scheduler_batches_emitted, 2);
assert_eq!(limiter.active_file_reads(), 2);
first_partition_gate.notify_one();
let mut ids = vec![batch_id(
&stream.next().await.ok_or("second batch missing")??,
)];
while let Some(batch) = stream.next().await {
ids.push(batch_id(&batch?));
}
assert_eq!(ids, [2, 10, 20]);
assert_eq!(metrics.snapshot().scheduler_batches_emitted, 4);
assert_eq!(metrics.snapshot().scan_partitions_completed, 2);
assert_eq!(limiter.active_file_reads(), 0);
Ok(())
}
#[tokio::test]
async fn merged_stream_drop_cancels_blocked_partitions_and_releases_permits()
-> Result<(), Box<dyn std::error::Error>> {
let ControlledMerge {
mut stream,
limiter,
cancellation,
metrics,
..
} = controlled_merge()?;
let first = stream.next().await.ok_or("first batch missing")??;
assert_eq!(batch_id(&first), 1);
wait_for_batches(&metrics, 2).await;
assert_eq!(limiter.active_file_reads(), 2);
drop(stream);
assert!(cancellation.is_cancelled());
timeout(Duration::from_secs(5), async {
while limiter.active_file_reads() != 0 {
tokio::task::yield_now().await;
}
})
.await?;
assert_eq!(metrics.snapshot().scheduler_batches_emitted, 2);
assert_eq!(metrics.snapshot().scan_partitions_completed, 0);
Ok(())
}
#[tokio::test]
async fn merged_stream_forwards_one_concurrent_error_and_releases_permits()
-> Result<(), Box<dyn std::error::Error>> {
let options = execution_options()?;
let limiter = ScanReadLimiter::new(options, 2, 2);
let cancellation = ScanCancellation::new();
let metrics = metrics();
let executor: FileExecutor<i32, FileBatchStream> = Arc::new(|task, permit, _| {
async move {
Ok(if task == 1 {
Box::pin(stream::once(async move {
let _permit = permit;
pending::<Result<RecordBatch, crate::DeltaReaderError>>().await
})) as FileBatchStream
} else {
Box::pin(stream::once(async move {
let _permit = permit;
Err(InvalidConfigurationSnafu {
reason: "controlled_partition_failure",
}
.build())
})) as FileBatchStream
})
}
.boxed()
});
let admission = Arc::new(|_: &i32| Ok(FileAdmissionDecision::Admit));
let first = PartitionStream::new(
vec![1],
limiter.partition(0)?,
options,
admission.clone(),
Arc::clone(&executor),
metrics.clone(),
cancellation.clone(),
);
let second = PartitionStream::new(
vec![2],
limiter.partition(1)?,
options,
admission,
executor,
metrics.clone(),
cancellation.clone(),
);
let mut stream = direct_stream(VecDeque::from([first, second]), metrics);
let error = timeout(Duration::from_secs(5), stream.next())
.await?
.ok_or("error item missing")?
.expect_err("controlled partition must fail");
assert_eq!(error.code(), "invalid_configuration");
assert!(stream.next().await.is_none());
assert!(cancellation.is_cancelled());
timeout(Duration::from_secs(5), async {
while limiter.active_file_reads() != 0 {
tokio::task::yield_now().await;
}
})
.await?;
Ok(())
}
}