Skip to main content

deltalake_core/operations/
optimize.rs

1//! Optimize a Delta Table
2//!
3//! Perform bin-packing on a Delta Table which merges small files into a large
4//! file. Bin-packing reduces the number of API calls required for read
5//! operations.
6//!
7//! Optimize will fail if a concurrent write operation removes files from the
8//! table (such as in an overwrite). It will always succeed if concurrent writers
9//! are only appending.
10//!
11//! Optimize increments the table's version and creates remove actions for
12//! optimized files. Optimize does not delete files from storage. To delete
13//! files that were removed, call `vacuum` on [`DeltaTable`].
14//!
15//! See [`OptimizeBuilder`] for configuration.
16//!
17//! # Example
18//! ```rust ignore
19//! let table = open_table(Url::from_directory_path("/abs/path/to/table").unwrap())?;
20//! let (table, metrics) = OptimizeBuilder::new(table.object_store(), table.state).await?;
21//! ````
22
23use std::collections::HashMap;
24use std::fmt;
25use std::num::NonZeroU64;
26use std::sync::Arc;
27use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
28
29use arrow::array::RecordBatch;
30use arrow::datatypes::SchemaRef;
31use datafusion::catalog::Session;
32use datafusion::execution::context::{SessionContext, SessionState};
33use delta_kernel::engine::arrow_conversion::TryIntoArrow as _;
34use delta_kernel::expressions::Scalar;
35use delta_kernel::table_features::ColumnMappingMode;
36use delta_kernel::table_properties::DataSkippingNumIndexedCols;
37use futures::future::BoxFuture;
38use futures::stream::BoxStream;
39use futures::{Future, StreamExt, TryStreamExt};
40use indexmap::IndexMap;
41use itertools::Itertools;
42use num_cpus;
43use parquet::basic::{Compression, ZstdLevel};
44use parquet::errors::ParquetError;
45use parquet::file::properties::WriterProperties;
46use serde::{Deserialize, Deserializer, Serialize, Serializer, de::Error as DeError};
47use tracing::*;
48use uuid::Uuid;
49
50use super::write::writer::{PartitionWriter, PartitionWriterConfig};
51use super::{CustomExecuteHandler, Operation};
52use crate::delta_datafusion::{
53    DataFusionMixins, DeltaScanConfig, DeltaScanNext, SessionFallbackPolicy, SessionResolveContext,
54    create_session_state_with_spill_config, resolve_session_state, update_datafusion_session,
55};
56use crate::errors::{DeltaResult, DeltaTableError, unsupported_column_mapping_write};
57use crate::kernel::transaction::{CommitBuilder, CommitProperties, DEFAULT_RETRIES, PROTOCOL};
58use crate::kernel::{Action, Add, DataType, PartitionsExt, Remove, StructType, Version};
59use crate::kernel::{EagerSnapshot, resolve_snapshot};
60use crate::logstore::{LogStore, LogStoreRef, ObjectStoreRef};
61use crate::parquet_utils::default_writer_properties;
62use crate::protocol::DeltaOperation;
63use crate::table::config::TablePropertiesExt as _;
64use crate::table::state::DeltaTableState;
65use crate::writer::utils::arrow_schema_without_partitions;
66use crate::{DeltaTable, ObjectMeta, PartitionFilter, to_kernel_predicate};
67
68/// Planner used by optimize.
69#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
70#[serde(rename_all = "camelCase")]
71pub enum PlannerStrategy {
72    /// Older metrics with no planner field.
73    #[default]
74    UnknownLegacy,
75    /// Compact planner.
76    PreserveLocality,
77    /// Z order planner.
78    ZOrder,
79}
80
81/// Metrics from Optimize
82#[derive(Default, Debug, PartialEq, Clone, Serialize, Deserialize)]
83#[serde(rename_all = "camelCase", from = "MetricsSerde")]
84pub struct Metrics {
85    /// Number of optimized files added
86    pub num_files_added: u64,
87    /// Number of unoptimized files removed
88    pub num_files_removed: u64,
89    /// Detailed metrics for the add operation
90    #[serde(
91        serialize_with = "serialize_metric_details",
92        deserialize_with = "deserialize_metric_details"
93    )]
94    pub files_added: MetricDetails,
95    /// Detailed metrics for the remove operation
96    #[serde(
97        serialize_with = "serialize_metric_details",
98        deserialize_with = "deserialize_metric_details"
99    )]
100    pub files_removed: MetricDetails,
101    /// Number of partitions that had at least one file optimized
102    pub partitions_optimized: u64,
103    /// The number of batches written
104    pub num_batches: u64,
105    /// How many files were considered during optimization. Not every file considered is optimized
106    pub total_considered_files: usize,
107    /// How many files were considered for optimization but were skipped
108    pub total_files_skipped: usize,
109    /// Compatibility field for `preserved_stable_order`
110    pub preserve_insertion_order: bool,
111    /// Planner used for this run
112    pub planner_strategy: PlannerStrategy,
113    /// True when file order is kept within a partition
114    pub preserved_stable_order: bool,
115    /// Largest count of adjacent input files in one bin
116    pub max_bin_span_files: usize,
117}
118
119#[derive(Debug, Deserialize)]
120#[serde(rename_all = "camelCase")]
121struct MetricsSerde {
122    num_files_added: u64,
123    num_files_removed: u64,
124    #[serde(deserialize_with = "deserialize_metric_details")]
125    files_added: MetricDetails,
126    #[serde(deserialize_with = "deserialize_metric_details")]
127    files_removed: MetricDetails,
128    partitions_optimized: u64,
129    num_batches: u64,
130    total_considered_files: usize,
131    total_files_skipped: usize,
132    #[serde(default)]
133    preserve_insertion_order: Option<bool>,
134    #[serde(default)]
135    planner_strategy: PlannerStrategy,
136    #[serde(default)]
137    preserved_stable_order: Option<bool>,
138    #[serde(default)]
139    max_bin_span_files: usize,
140}
141
142impl From<MetricsSerde> for Metrics {
143    fn from(value: MetricsSerde) -> Self {
144        let preserve_insertion_order = value
145            .preserve_insertion_order
146            .or(value.preserved_stable_order)
147            .unwrap_or(false);
148        let preserved_stable_order = value.preserved_stable_order.unwrap_or(false);
149
150        Self {
151            num_files_added: value.num_files_added,
152            num_files_removed: value.num_files_removed,
153            files_added: value.files_added,
154            files_removed: value.files_removed,
155            partitions_optimized: value.partitions_optimized,
156            num_batches: value.num_batches,
157            total_considered_files: value.total_considered_files,
158            total_files_skipped: value.total_files_skipped,
159            preserve_insertion_order,
160            planner_strategy: value.planner_strategy,
161            preserved_stable_order,
162            max_bin_span_files: value.max_bin_span_files,
163        }
164    }
165}
166
167// Custom serialization function that serializes metric details as a string
168fn serialize_metric_details<S>(value: &MetricDetails, serializer: S) -> Result<S::Ok, S::Error>
169where
170    S: Serializer,
171{
172    serializer.serialize_str(&value.to_string())
173}
174
175// Custom deserialization that parses a JSON string into MetricDetails
176fn deserialize_metric_details<'de, D>(deserializer: D) -> Result<MetricDetails, D::Error>
177where
178    D: Deserializer<'de>,
179{
180    let s: String = Deserialize::deserialize(deserializer)?;
181    serde_json::from_str(&s).map_err(DeError::custom)
182}
183
184/// Statistics on files for a particular operation
185/// Operation can be remove or add
186#[derive(Debug, PartialEq, Clone, Serialize, Deserialize)]
187#[serde(rename_all = "camelCase")]
188pub struct MetricDetails {
189    /// Average file size of a operation
190    pub avg: f64,
191    /// Maximum file size of a operation
192    pub max: i64,
193    /// Minimum file size of a operation
194    pub min: i64,
195    /// Number of files encountered during operation
196    pub total_files: usize,
197    /// Sum of file sizes of a operation
198    pub total_size: i64,
199}
200
201impl MetricDetails {
202    /// Add a partial metric to the metrics
203    pub fn add(&mut self, partial: &MetricDetails) {
204        self.min = std::cmp::min(self.min, partial.min);
205        self.max = std::cmp::max(self.max, partial.max);
206        self.total_files += partial.total_files;
207        self.total_size += partial.total_size;
208        self.avg = self.total_size as f64 / self.total_files as f64;
209    }
210}
211
212impl fmt::Display for MetricDetails {
213    /// Display the metric details using serde serialization
214    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
215        serde_json::to_string(self).map_err(|_| fmt::Error)?.fmt(f)
216    }
217}
218
219#[derive(Debug)]
220/// Metrics for a single partition
221pub struct PartialMetrics {
222    /// Number of optimized files added
223    pub num_files_added: u64,
224    /// Number of unoptimized files removed
225    pub num_files_removed: u64,
226    /// Detailed metrics for the add operation
227    pub files_added: MetricDetails,
228    /// Detailed metrics for the remove operation
229    pub files_removed: MetricDetails,
230    /// The number of batches written
231    pub num_batches: u64,
232}
233
234impl Metrics {
235    /// Add a partial metric to the metrics
236    pub fn add(&mut self, partial: &PartialMetrics) {
237        self.num_files_added += partial.num_files_added;
238        self.num_files_removed += partial.num_files_removed;
239        self.files_added.add(&partial.files_added);
240        self.files_removed.add(&partial.files_removed);
241        self.num_batches += partial.num_batches;
242    }
243
244    fn apply_planner_stats(&mut self, planner_stats: &PlannerStats) {
245        self.planner_strategy = planner_stats.planner_strategy;
246        self.preserved_stable_order = planner_stats.preserved_stable_order;
247        self.preserve_insertion_order = planner_stats.preserved_stable_order;
248        self.max_bin_span_files = planner_stats.max_bin_span_files;
249    }
250}
251
252impl Default for MetricDetails {
253    fn default() -> Self {
254        MetricDetails {
255            min: i64::MAX,
256            max: 0,
257            avg: 0.0,
258            total_files: 0,
259            total_size: 0,
260        }
261    }
262}
263
264/// Type of optimization to perform.
265#[derive(Debug)]
266pub enum OptimizeType {
267    /// Compact files into pre-determined bins
268    Compact,
269    /// Z-order files based on provided columns
270    ZOrder(Vec<String>),
271}
272
273/// Optimize a Delta table with given options
274///
275/// If a target file size is not provided then `delta.targetFileSize` from the
276/// table's configuration is read. Otherwise a default value is used.
277pub struct OptimizeBuilder<'a> {
278    /// A snapshot of the to-be-optimized table's state
279    snapshot: Option<EagerSnapshot>,
280    /// Delta object store for handling data files
281    log_store: LogStoreRef,
282    /// Filters to select specific table partitions to be optimized
283    filters: &'a [PartitionFilter],
284    /// Desired file size after bin-packing files
285    target_size: Option<NonZeroU64>,
286    /// Properties passed to underlying parquet writer
287    writer_properties: Option<WriterProperties>,
288    /// Commit properties and configuration
289    commit_properties: CommitProperties,
290    /// Maximum number of concurrent tasks (default is number of cpus)
291    max_concurrent_tasks: usize,
292    /// Optimize type
293    optimize_type: OptimizeType,
294    /// Datafusion session state relevant for executing the input plan
295    session: Option<Arc<dyn Session>>,
296    session_fallback_policy: SessionFallbackPolicy,
297    min_commit_interval: Option<Duration>,
298    custom_execute_handler: Option<Arc<dyn CustomExecuteHandler>>,
299}
300
301impl super::Operation for OptimizeBuilder<'_> {
302    fn log_store(&self) -> &LogStoreRef {
303        &self.log_store
304    }
305    fn get_custom_execute_handler(&self) -> Option<Arc<dyn CustomExecuteHandler>> {
306        self.custom_execute_handler.clone()
307    }
308}
309
310impl<'a> OptimizeBuilder<'a> {
311    /// Create a new [`OptimizeBuilder`]
312    pub(crate) fn new(log_store: LogStoreRef, snapshot: Option<EagerSnapshot>) -> Self {
313        Self {
314            snapshot,
315            log_store,
316            filters: &[],
317            target_size: None,
318            writer_properties: None,
319            commit_properties: CommitProperties::default(),
320            max_concurrent_tasks: num_cpus::get(),
321            optimize_type: OptimizeType::Compact,
322            min_commit_interval: None,
323            session: None,
324            session_fallback_policy: SessionFallbackPolicy::default(),
325            custom_execute_handler: None,
326        }
327    }
328
329    /// Choose the type of optimization to perform. Defaults to [OptimizeType::Compact].
330    pub fn with_type(mut self, optimize_type: OptimizeType) -> Self {
331        self.optimize_type = optimize_type;
332        self
333    }
334
335    /// Only optimize files that return true for the specified partition filter
336    pub fn with_filters(mut self, filters: &'a [PartitionFilter]) -> Self {
337        self.filters = filters;
338        self
339    }
340
341    /// Set the target file size
342    pub fn with_target_size(mut self, target: NonZeroU64) -> Self {
343        self.target_size = Some(target);
344        self
345    }
346
347    /// Writer properties passed to parquet writer
348    pub fn with_writer_properties(mut self, writer_properties: WriterProperties) -> Self {
349        self.writer_properties = Some(writer_properties);
350        self
351    }
352
353    /// Additional information to write to the commit
354    pub fn with_commit_properties(mut self, commit_properties: CommitProperties) -> Self {
355        self.commit_properties = commit_properties;
356        self
357    }
358
359    /// Deprecated. This setting has no effect.
360    #[deprecated(
361        since = "0.32.0",
362        note = "compact always keeps partition file order, and z order does not; this setting has no effect"
363    )]
364    pub fn with_preserve_insertion_order(self, _preserve_insertion_order: bool) -> Self {
365        self
366    }
367
368    /// Max number of concurrent tasks
369    pub fn with_max_concurrent_tasks(mut self, max_concurrent_tasks: usize) -> Self {
370        self.max_concurrent_tasks = max_concurrent_tasks;
371        self
372    }
373
374    /// Min commit interval
375    pub fn with_min_commit_interval(mut self, min_commit_interval: Duration) -> Self {
376        self.min_commit_interval = Some(min_commit_interval);
377        self
378    }
379
380    /// Set a custom execute handler, for pre and post execution
381    pub fn with_custom_execute_handler(mut self, handler: Arc<dyn CustomExecuteHandler>) -> Self {
382        self.custom_execute_handler = Some(handler);
383        self
384    }
385
386    /// Set the DataFusion session used for planning and execution.
387    ///
388    /// The provided `session` should wrap a concrete `datafusion::execution::context::SessionState`.
389    ///
390    /// If `session` is not a `SessionState`, the default policy is to log a warning and fall back to
391    /// internal defaults. To make this strict (error instead), set
392    /// `with_session_fallback_policy(SessionFallbackPolicy::RequireSessionState)`.
393    ///
394    /// Example: `Arc::new(create_session().state())`.
395    pub fn with_session_state(mut self, session: Arc<dyn Session>) -> Self {
396        self.session = Some(session);
397        self
398    }
399
400    /// Control how delta-rs resolves the provided session when it is not a concrete `SessionState`.
401    ///
402    /// Defaults to `SessionFallbackPolicy::InternalDefaults` to preserve existing behavior.
403    pub fn with_session_fallback_policy(mut self, policy: SessionFallbackPolicy) -> Self {
404        self.session_fallback_policy = policy;
405        self
406    }
407}
408
409impl<'a> std::future::IntoFuture for OptimizeBuilder<'a> {
410    type Output = DeltaResult<(DeltaTable, Metrics)>;
411    type IntoFuture = BoxFuture<'a, Self::Output>;
412
413    fn into_future(self) -> Self::IntoFuture {
414        let this = self;
415
416        Box::pin(async move {
417            let snapshot =
418                resolve_snapshot(&this.log_store, this.snapshot.clone(), true, None).await?;
419            if snapshot.table_configuration().column_mapping_mode() != ColumnMappingMode::None {
420                return Err(unsupported_column_mapping_write("OPTIMIZE"));
421            }
422            PROTOCOL.can_write_to(&snapshot)?;
423
424            let operation_id = this.get_operation_id();
425            this.pre_execute(operation_id).await?;
426
427            let writer_properties = this.writer_properties.unwrap_or_else(|| {
428                default_writer_properties(Compression::ZSTD(ZstdLevel::try_new(4).unwrap()))
429            });
430            let (session, _) = resolve_session_state(
431                this.session.as_deref(),
432                this.session_fallback_policy,
433                || create_session_state_with_spill_config(None, None),
434                SessionResolveContext {
435                    operation: "optimize",
436                    table_uri: Some(this.log_store.root_url()),
437                    cdc: false,
438                },
439            )?;
440            let plan = create_merge_plan(
441                &this.log_store,
442                this.optimize_type,
443                &snapshot,
444                this.filters,
445                this.target_size.to_owned(),
446                writer_properties,
447                session,
448            )
449            .await?;
450
451            let metrics = plan
452                .execute(
453                    this.log_store.clone(),
454                    &snapshot,
455                    this.max_concurrent_tasks,
456                    this.min_commit_interval,
457                    this.commit_properties.clone(),
458                    operation_id,
459                    this.custom_execute_handler.as_ref(),
460                )
461                .await?;
462
463            if let Some(handler) = this.custom_execute_handler {
464                handler.post_execute(&this.log_store, operation_id).await?;
465            }
466            let mut table =
467                DeltaTable::new_with_state(this.log_store, DeltaTableState::new(snapshot));
468            table.update_state().await?;
469            Ok((table, metrics))
470        })
471    }
472}
473
474#[derive(Debug, Clone)]
475struct OptimizeInput {
476    target_size: NonZeroU64,
477    predicate: Option<String>,
478}
479
480const MAX_OPTIMIZE_TARGET_SIZE: u64 = i64::MAX as u64;
481
482fn optimize_target_size_to_i64(target_size: NonZeroU64) -> Result<i64, DeltaTableError> {
483    i64::try_from(target_size.get()).map_err(|_| {
484        DeltaTableError::Generic(format!(
485            "optimize target_size {} exceeds i64::MAX ({MAX_OPTIMIZE_TARGET_SIZE})",
486            target_size.get()
487        ))
488    })
489}
490
491impl TryFrom<OptimizeInput> for DeltaOperation {
492    type Error = DeltaTableError;
493
494    fn try_from(opt_input: OptimizeInput) -> Result<Self, Self::Error> {
495        Ok(DeltaOperation::Optimize {
496            target_size: optimize_target_size_to_i64(opt_input.target_size)?,
497            predicate: opt_input.predicate,
498        })
499    }
500}
501
502/// Generate an appropriate remove action for the optimization task
503fn create_remove(add: &Add) -> Action {
504    // NOTE unwrap is safe since UNIX_EPOCH will always be earlier then now.
505    let deletion_time = SystemTime::now().duration_since(UNIX_EPOCH).unwrap();
506    let deletion_time = deletion_time.as_millis() as i64;
507
508    Action::Remove(Remove {
509        path: add.path.clone(),
510        deletion_timestamp: Some(deletion_time),
511        data_change: false,
512        extended_file_metadata: Some(true),
513        partition_values: Some(add.partition_values.clone()),
514        size: Some(add.size),
515        deletion_vector: add.deletion_vector.clone(),
516        tags: add.tags.clone(),
517        base_row_id: add.base_row_id,
518        default_row_commit_version: add.default_row_commit_version,
519    })
520}
521
522/// Layout for optimizing a plan
523///
524/// Within each partition, we identify a set of files that need to be merged
525/// together and/or sorted together.
526#[derive(Debug)]
527enum OptimizeOperations {
528    /// Plan to compact files into bins
529    ///
530    /// Bins keep partition file order, stop at ordinal gaps, and stop at
531    /// skipped large files. Bins of size 1 are dropped.
532    Compact(HashMap<String, (IndexMap<String, Scalar>, Vec<MergeBin>)>),
533    /// Plan to Z-order each partition
534    ZOrder(
535        Vec<String>,
536        HashMap<String, (IndexMap<String, Scalar>, MergeBin)>,
537    ),
538    // TODO: Sort
539}
540
541impl Default for OptimizeOperations {
542    fn default() -> Self {
543        OptimizeOperations::Compact(HashMap::new())
544    }
545}
546
547#[derive(Debug)]
548/// Encapsulates the operations required to optimize a Delta Table
549pub struct MergePlan {
550    operations: OptimizeOperations,
551    /// Metrics collected during operation
552    metrics: Metrics,
553    /// Planner metadata copied into buffered and total metrics
554    planner_stats: PlannerStats,
555    /// Parameters passed down to merge tasks
556    task_parameters: Arc<MergeTaskParameters>,
557    /// Version of the table at beginning of optimization. Used for conflict resolution.
558    read_table_version: Version,
559    /// Session state used for provider owned rewrite scans.
560    read_session: Arc<SessionState>,
561}
562
563#[derive(Debug, Clone, Default)]
564struct PlannerStats {
565    planner_strategy: PlannerStrategy,
566    preserved_stable_order: bool,
567    max_bin_span_files: usize,
568}
569
570impl PlannerStats {
571    fn preserve_locality() -> Self {
572        Self {
573            planner_strategy: PlannerStrategy::PreserveLocality,
574            preserved_stable_order: true,
575            max_bin_span_files: 0,
576        }
577    }
578
579    fn z_order(max_bin_span_files: usize) -> Self {
580        Self {
581            planner_strategy: PlannerStrategy::ZOrder,
582            preserved_stable_order: false,
583            max_bin_span_files,
584        }
585    }
586
587    fn absorb(&mut self, other: &PlannerStats) {
588        self.max_bin_span_files = self.max_bin_span_files.max(other.max_bin_span_files);
589    }
590}
591
592/// Parameters passed to individual merge tasks
593#[derive(Debug)]
594pub struct MergeTaskParameters {
595    /// Schema of written files
596    file_schema: SchemaRef,
597    /// Properties passed to parquet writer
598    writer_properties: WriterProperties,
599    /// Input parameters for the optimize operation
600    input_parameters: OptimizeInput,
601    /// Num index cols to collect stats for
602    num_indexed_cols: DataSkippingNumIndexedCols,
603    /// Stats columns, specific columns to collect stats from, takes precedence over num_indexed_cols
604    stats_columns: Option<Vec<String>>,
605}
606
607/// A stream of record batches, with a ParquetError on failure.
608type ParquetReadStream = BoxStream<'static, Result<RecordBatch, ParquetError>>;
609
610#[derive(Clone)]
611struct SelectedFileScanFactory {
612    snapshot: EagerSnapshot,
613    log_store: LogStoreRef,
614    scan_config: DeltaScanConfig,
615    read_operation_id: Option<Uuid>,
616}
617
618impl SelectedFileScanFactory {
619    fn try_new(
620        snapshot: &EagerSnapshot,
621        log_store: LogStoreRef,
622        session: &dyn Session,
623        read_operation_id: Option<Uuid>,
624    ) -> Result<Self, DeltaTableError> {
625        Ok(Self {
626            snapshot: snapshot.clone(),
627            log_store,
628            // Mirror the caller's DataFusion session flags so rewrite scans keep
629            // the same parquet/view type behavior as the rest of optimize.
630            scan_config: DeltaScanConfig::new_from_session(session)
631                .with_schema(snapshot.input_schema()),
632            read_operation_id,
633        })
634    }
635
636    fn provider_for(
637        &self,
638        adds: impl IntoIterator<Item = Add>,
639    ) -> Result<DeltaScanNext, DeltaTableError> {
640        let provider = DeltaScanNext::new(self.snapshot.clone(), self.scan_config.clone())?
641            .with_log_store(self.log_store.clone());
642        let provider = if let Some(operation_id) = self.read_operation_id {
643            provider.with_operation_id(operation_id)
644        } else {
645            provider
646        };
647        provider.with_selected_adds(adds)
648    }
649}
650
651impl MergePlan {
652    /// Rewrites files in a single partition.
653    ///
654    /// Returns a vector of add and remove actions, as well as the partial metrics
655    /// collected during the operation.
656    async fn rewrite_files<F>(
657        task_parameters: Arc<MergeTaskParameters>,
658        partition_values: IndexMap<String, Scalar>,
659        files: MergeBin,
660        object_store: ObjectStoreRef,
661        read_stream: F,
662        ignore_target_size: bool,
663    ) -> Result<(Vec<Action>, PartialMetrics), DeltaTableError>
664    where
665        F: Future<Output = Result<ParquetReadStream, DeltaTableError>> + Send + 'static,
666    {
667        debug!("Rewriting files in partition: {partition_values:?}");
668        // First, initialize metrics
669        let mut partial_actions = files.iter().map(create_remove).collect::<Vec<_>>();
670
671        let files_removed = files
672            .iter()
673            .fold(MetricDetails::default(), |mut curr, file| {
674                curr.total_files += 1;
675                curr.total_size += file.size;
676                curr.max = std::cmp::max(curr.max, file.size);
677                curr.min = std::cmp::min(curr.min, file.size);
678                curr
679            });
680
681        let mut partial_metrics = PartialMetrics {
682            num_files_added: 0,
683            num_files_removed: files.len() as u64,
684            files_added: MetricDetails::default(),
685            files_removed,
686            num_batches: 0,
687        };
688
689        // Next, initialize the writer
690        let writer_config = PartitionWriterConfig::try_new(
691            task_parameters.file_schema.clone(),
692            partition_values.clone(),
693            Some(task_parameters.writer_properties.clone()),
694            // Since we know the total size of the bin, we can set the target file size to None.
695            if ignore_target_size {
696                None
697            } else {
698                Some(task_parameters.input_parameters.target_size)
699            },
700            None,
701            None,
702        )?;
703        let mut writer = PartitionWriter::try_with_config(
704            object_store,
705            writer_config,
706            task_parameters.num_indexed_cols,
707            task_parameters.stats_columns.clone(),
708        )?;
709
710        let mut read_stream = read_stream.await?;
711
712        while let Some(maybe_batch) = read_stream.next().await {
713            let mut batch = maybe_batch?;
714
715            batch = crate::kernel::schema::cast::cast_record_batch(
716                &batch,
717                task_parameters.file_schema.clone(),
718                false,
719                true,
720            )?;
721            partial_metrics.num_batches += 1;
722            writer.write(&batch).await?;
723        }
724
725        let add_actions = writer.close().await?.into_iter().map(|mut add| {
726            add.data_change = false;
727
728            let size = add.size;
729
730            partial_metrics.num_files_added += 1;
731            partial_metrics.files_added.total_files += 1;
732            partial_metrics.files_added.total_size += size;
733            partial_metrics.files_added.max = std::cmp::max(partial_metrics.files_added.max, size);
734            partial_metrics.files_added.min = std::cmp::min(partial_metrics.files_added.min, size);
735
736            Action::Add(add)
737        });
738        partial_actions.extend(add_actions);
739
740        debug!("Finished rewriting files in partition: {partition_values:?}");
741
742        Ok((partial_actions, partial_metrics))
743    }
744
745    async fn read_selected_files(
746        files: MergeBin,
747        context: Arc<SessionContext>,
748        scan_factory: SelectedFileScanFactory,
749    ) -> Result<ParquetReadStream, DeltaTableError> {
750        let provider = scan_factory.provider_for(files.iter().cloned())?;
751        let df = context.read_table(Arc::new(provider))?;
752        let stream = df
753            .execute_stream()
754            .await?
755            .map_err(|err| {
756                ParquetError::General(format!(
757                    "Optimize selected-file scan failed while scanning data: {err}"
758                ))
759            })
760            .boxed();
761        Ok(stream)
762    }
763
764    /// Datafusion-based z-order read.
765    async fn read_zorder(
766        files: MergeBin,
767        context: Arc<zorder::ZOrderExecContext>,
768        scan_factory: SelectedFileScanFactory,
769    ) -> Result<BoxStream<'static, Result<RecordBatch, ParquetError>>, DeltaTableError> {
770        use datafusion::functions::core::expr_ext::FieldAccessor;
771        use datafusion::logical_expr::expr::ScalarFunction;
772        use datafusion::logical_expr::{Expr, ScalarUDF, ident};
773
774        let provider = scan_factory.provider_for(files.iter().cloned())?;
775        let df = context.ctx.read_table(Arc::new(provider))?;
776
777        let cols = context
778            .columns
779            .iter()
780            .map(|col_name| {
781                let mut segments = col_name.split('.');
782                let first = segments.next().expect("column name cannot be empty");
783                let mut expr = ident(first);
784                for segment in segments {
785                    expr = expr.field(segment);
786                }
787                expr
788            })
789            .collect_vec();
790        let expr = Expr::ScalarFunction(ScalarFunction::new_udf(
791            Arc::new(ScalarUDF::from(zorder::datafusion::ZOrderUDF)),
792            cols,
793        ));
794        let df = df.sort(vec![expr.sort(true, true)])?;
795
796        let stream = df
797            .execute_stream()
798            .await?
799            .map_err(|err| {
800                ParquetError::General(format!("Z-order failed while scanning data: {err}"))
801            })
802            .boxed();
803
804        Ok(stream)
805    }
806
807    /// Perform the operations outlined in the plan.
808    #[allow(clippy::too_many_arguments)]
809    #[instrument(skip_all, fields(operation = "optimize", version = snapshot.version()))]
810    pub async fn execute(
811        mut self,
812        log_store: LogStoreRef,
813        snapshot: &EagerSnapshot,
814        max_concurrent_tasks: usize,
815        min_commit_interval: Option<Duration>,
816        commit_properties: CommitProperties,
817        operation_id: Uuid,
818        handle: Option<&Arc<dyn CustomExecuteHandler>>,
819    ) -> Result<Metrics, DeltaTableError> {
820        let operations = std::mem::take(&mut self.operations);
821        let read_session = self.read_session.clone();
822        info!("starting optimize execution");
823        let object_store = log_store.object_store(Some(operation_id));
824        update_datafusion_session(
825            read_session.as_ref(),
826            log_store.as_ref(),
827            Some(operation_id),
828        )?;
829
830        let mut stream = match operations {
831            OptimizeOperations::Compact(bins) => {
832                let read_context = Arc::new(SessionContext::new_with_state(
833                    read_session.as_ref().clone(),
834                ));
835                let scan_factory = SelectedFileScanFactory::try_new(
836                    snapshot,
837                    log_store.clone(),
838                    read_session.as_ref(),
839                    Some(operation_id),
840                )?;
841                let task_parameters = self.task_parameters.clone();
842
843                futures::stream::iter(bins)
844                    .flat_map(|(_, (partition, bins))| {
845                        futures::stream::iter(bins).map(move |bin| (partition.clone(), bin))
846                    })
847                    .map(move |(partition, files)| {
848                        debug!(
849                            "merging a group of {} files in partition {partition:?}",
850                            files.len(),
851                        );
852                        for file in files.iter() {
853                            debug!("  file {}", file.path);
854                        }
855
856                        let batch_stream = Self::read_selected_files(
857                            files.clone(),
858                            read_context.clone(),
859                            scan_factory.clone(),
860                        );
861
862                        let rewrite_result = tokio::task::spawn(Self::rewrite_files(
863                            task_parameters.clone(),
864                            partition,
865                            files,
866                            object_store.clone(),
867                            batch_stream,
868                            true,
869                        ));
870                        util::flatten_join_error(rewrite_result)
871                    })
872                    .buffered(max_concurrent_tasks)
873                    .boxed()
874            }
875            OptimizeOperations::ZOrder(zorder_columns, bins) => {
876                debug!("Starting zorder with the columns: {zorder_columns:?} {bins:?}");
877
878                let exec_context = Arc::new(zorder::ZOrderExecContext::new(
879                    zorder_columns,
880                    read_session.as_ref().clone(),
881                    object_store,
882                )?);
883                let task_parameters = self.task_parameters.clone();
884                let scan_factory = SelectedFileScanFactory::try_new(
885                    snapshot,
886                    log_store.clone(),
887                    read_session.as_ref(),
888                    Some(operation_id),
889                )?;
890
891                // For each rewrite evaluate the predicate and then modify each expression
892                // to either compute the new value or obtain the old one then write these batches
893                let log_store = log_store.clone();
894                futures::stream::iter(bins)
895                    .map(move |(_, (partition, files))| {
896                        let batch_stream = Self::read_zorder(
897                            files.clone(),
898                            exec_context.clone(),
899                            scan_factory.clone(),
900                        );
901                        let rewrite_result = tokio::task::spawn(Self::rewrite_files(
902                            task_parameters.clone(),
903                            partition,
904                            files,
905                            log_store.object_store(Some(operation_id)),
906                            batch_stream,
907                            false,
908                        ));
909                        util::flatten_join_error(rewrite_result)
910                    })
911                    .buffer_unordered(max_concurrent_tasks)
912                    .boxed()
913            }
914        };
915
916        let mut table =
917            DeltaTable::new_with_state(log_store.clone(), DeltaTableState::new(snapshot.clone()));
918
919        // Actions buffered so far. These will be flushed either at the end
920        // or when we reach the commit interval.
921        let mut actions = vec![];
922
923        // Each time we commit, we'll reset buffered_metrics to orig_metrics.
924        let mut orig_metrics = std::mem::take(&mut self.metrics);
925        orig_metrics.apply_planner_stats(&self.planner_stats);
926        let mut buffered_metrics = orig_metrics.clone();
927        let mut total_metrics = orig_metrics.clone();
928
929        let mut last_commit = Instant::now();
930        let mut commits_made = 0;
931        let mut snapshot = snapshot.clone();
932        loop {
933            let next = stream.next().await.transpose()?;
934
935            let end = next.is_none();
936
937            if let Some((partial_actions, partial_metrics)) = next {
938                debug!("Recording metrics for a completed partition");
939                actions.extend(partial_actions);
940                buffered_metrics.add(&partial_metrics);
941                total_metrics.add(&partial_metrics);
942            }
943
944            let now = Instant::now();
945            let mature = match min_commit_interval {
946                None => false,
947                Some(i) => now.duration_since(last_commit) > i,
948            };
949            if !actions.is_empty() && (mature || end) {
950                let actions = std::mem::take(&mut actions);
951                last_commit = now;
952
953                let mut properties = CommitProperties::default();
954                properties.app_metadata = commit_properties.app_metadata.clone();
955                properties
956                    .app_metadata
957                    .insert("readVersion".to_owned(), self.read_table_version.into());
958                let maybe_map_metrics = serde_json::to_value(std::mem::replace(
959                    &mut buffered_metrics,
960                    orig_metrics.clone(),
961                ));
962                if let Ok(map) = maybe_map_metrics {
963                    properties
964                        .app_metadata
965                        .insert("operationMetrics".to_owned(), map);
966                }
967
968                debug!("committing {} actions", actions.len());
969
970                let commit = CommitBuilder::from(properties)
971                    .with_actions(actions)
972                    .with_operation_id(operation_id)
973                    .with_post_commit_hook_handler(handle.cloned())
974                    .with_max_retries(DEFAULT_RETRIES + commits_made)
975                    .build(
976                        Some(&snapshot),
977                        log_store.clone(),
978                        self.task_parameters.input_parameters.clone().try_into()?,
979                    )
980                    .await?;
981                snapshot = commit.snapshot().snapshot;
982                commits_made += 1;
983            }
984
985            if end {
986                break;
987            }
988        }
989
990        if total_metrics.num_files_added == 0 {
991            total_metrics.files_added.min = 0;
992        }
993        if total_metrics.num_files_removed == 0 {
994            total_metrics.files_removed.min = 0;
995        }
996
997        table.state = Some(DeltaTableState::new(snapshot));
998
999        Ok(total_metrics)
1000    }
1001}
1002
1003/// Build a Plan on which files to merge together. See [OptimizeBuilder]
1004#[instrument(skip_all, fields(operation = "create_merge_plan", version = snapshot.version()))]
1005pub async fn create_merge_plan(
1006    log_store: &dyn LogStore,
1007    optimize_type: OptimizeType,
1008    snapshot: &EagerSnapshot,
1009    filters: &[PartitionFilter],
1010    target_size: Option<NonZeroU64>,
1011    writer_properties: WriterProperties,
1012    session: SessionState,
1013) -> Result<MergePlan, DeltaTableError> {
1014    let target_size = target_size.unwrap_or_else(|| snapshot.table_properties().target_file_size());
1015    let _ = optimize_target_size_to_i64(target_size)?;
1016    let partitions_keys = snapshot.metadata().partition_columns();
1017
1018    let (operations, metrics, planner_stats) = match optimize_type {
1019        OptimizeType::Compact => {
1020            info!("building compaction plan");
1021            build_compaction_plan(log_store, snapshot, filters, target_size).await?
1022        }
1023        OptimizeType::ZOrder(zorder_columns) => {
1024            info!("building z-order plan");
1025            build_zorder_plan(
1026                log_store,
1027                zorder_columns,
1028                snapshot,
1029                partitions_keys,
1030                filters,
1031            )
1032            .await?
1033        }
1034    };
1035
1036    info!(
1037        partitions_optimized = metrics.partitions_optimized,
1038        total_considered_files = metrics.total_considered_files,
1039        "merge plan created"
1040    );
1041
1042    let input_parameters = OptimizeInput {
1043        target_size,
1044        predicate: serde_json::to_string(filters).ok(),
1045    };
1046    let file_schema = arrow_schema_without_partitions(
1047        &Arc::new(snapshot.schema().as_ref().try_into_arrow()?),
1048        partitions_keys,
1049    );
1050
1051    Ok(MergePlan {
1052        operations,
1053        metrics,
1054        planner_stats,
1055        task_parameters: Arc::new(MergeTaskParameters {
1056            file_schema,
1057            writer_properties,
1058            input_parameters,
1059            num_indexed_cols: snapshot.table_properties().num_indexed_cols(),
1060            stats_columns: snapshot
1061                .table_properties()
1062                .data_skipping_stats_columns
1063                .as_ref()
1064                .map(|v| v.iter().map(|v| v.to_string()).collect::<Vec<String>>()),
1065        }),
1066        read_table_version: snapshot.version(),
1067        read_session: Arc::new(session),
1068    })
1069}
1070
1071/// A collection of bins for a particular partition
1072#[derive(Debug, Clone)]
1073struct MergeBin {
1074    files: Vec<Add>,
1075    size_bytes: u64,
1076}
1077
1078impl MergeBin {
1079    pub fn new() -> Self {
1080        MergeBin {
1081            files: Vec::new(),
1082            size_bytes: 0,
1083        }
1084    }
1085
1086    fn total_file_size(&self) -> u64 {
1087        self.size_bytes
1088    }
1089
1090    fn is_empty(&self) -> bool {
1091        self.files.is_empty()
1092    }
1093
1094    fn len(&self) -> usize {
1095        self.files.len()
1096    }
1097
1098    fn from_file(add: Add) -> Self {
1099        let mut bin = Self::new();
1100        bin.add(add);
1101        bin
1102    }
1103
1104    fn add(&mut self, add: Add) {
1105        self.size_bytes += add.size as u64;
1106        self.files.push(add);
1107    }
1108
1109    fn iter(&self) -> impl Iterator<Item = &Add> {
1110        self.files.iter()
1111    }
1112}
1113
1114impl IntoIterator for MergeBin {
1115    type Item = Add;
1116    type IntoIter = std::vec::IntoIter<Self::Item>;
1117
1118    fn into_iter(self) -> Self::IntoIter {
1119        self.files.into_iter()
1120    }
1121}
1122
1123#[derive(Debug, Clone)]
1124struct OrderedFileCandidate {
1125    add: Add,
1126    stable_ordinal: usize,
1127    size_bytes: u64,
1128}
1129
1130fn plan_compaction_bins_in_stable_order(
1131    files: Vec<OrderedFileCandidate>,
1132    target_size: u64,
1133) -> (Vec<MergeBin>, PlannerStats) {
1134    let mut bins = Vec::new();
1135    let mut current = MergeBin::new();
1136    let mut current_first_ordinal = None;
1137    let mut current_last_ordinal = None;
1138    let mut planner_stats = PlannerStats::preserve_locality();
1139
1140    for file in files {
1141        if current.is_empty() {
1142            current = MergeBin::from_file(file.add);
1143            current_first_ordinal = Some(file.stable_ordinal);
1144            current_last_ordinal = Some(file.stable_ordinal);
1145            continue;
1146        }
1147
1148        let extends_contiguous_span = current_last_ordinal
1149            .map(|last| file.stable_ordinal == last + 1)
1150            .unwrap_or(false);
1151        if !extends_contiguous_span {
1152            if let (Some(first), Some(last)) = (current_first_ordinal, current_last_ordinal) {
1153                planner_stats.max_bin_span_files =
1154                    planner_stats.max_bin_span_files.max(last - first + 1);
1155            }
1156
1157            bins.push(current);
1158            current = MergeBin::from_file(file.add);
1159            current_first_ordinal = Some(file.stable_ordinal);
1160            current_last_ordinal = Some(file.stable_ordinal);
1161            continue;
1162        }
1163
1164        if current.total_file_size() + file.size_bytes <= target_size {
1165            current.add(file.add);
1166            current_last_ordinal = Some(file.stable_ordinal);
1167            continue;
1168        }
1169
1170        if let (Some(first), Some(last)) = (current_first_ordinal, current_last_ordinal) {
1171            planner_stats.max_bin_span_files =
1172                planner_stats.max_bin_span_files.max(last - first + 1);
1173        }
1174
1175        bins.push(current);
1176        current = MergeBin::from_file(file.add);
1177        current_first_ordinal = Some(file.stable_ordinal);
1178        current_last_ordinal = Some(file.stable_ordinal);
1179    }
1180
1181    if !current.is_empty() {
1182        if let (Some(first), Some(last)) = (current_first_ordinal, current_last_ordinal) {
1183            planner_stats.max_bin_span_files =
1184                planner_stats.max_bin_span_files.max(last - first + 1);
1185        }
1186        bins.push(current);
1187    }
1188
1189    (bins, planner_stats)
1190}
1191
1192async fn build_compaction_plan(
1193    log_store: &dyn LogStore,
1194    snapshot: &EagerSnapshot,
1195    filters: &[PartitionFilter],
1196    target_size: NonZeroU64,
1197) -> Result<(OptimizeOperations, Metrics, PlannerStats), DeltaTableError> {
1198    type PartitionFileEntry = (IndexMap<String, Scalar>, usize, Vec<OrderedFileCandidate>);
1199
1200    let mut metrics = Metrics::default();
1201    let mut planner_stats = PlannerStats::preserve_locality();
1202    let mut partition_files: HashMap<String, PartitionFileEntry> = HashMap::new();
1203
1204    let predicate = if filters.is_empty() {
1205        None
1206    } else {
1207        Some(Arc::new(to_kernel_predicate(
1208            filters,
1209            snapshot.schema().as_ref(),
1210        )?))
1211    };
1212
1213    // `file_views` returns active files with the newest first.
1214    // We use that order within each partition when building compact bins.
1215    let mut file_stream = snapshot.file_views(log_store, predicate);
1216    while let Some(file) = file_stream.next().await {
1217        let file = file?;
1218        metrics.total_considered_files += 1;
1219        let object_meta = ObjectMeta::try_from(&file)?;
1220        let partition_values = file
1221            .partition_values()
1222            .map(|v| {
1223                v.fields()
1224                    .iter()
1225                    .zip(v.values().iter())
1226                    .map(|(k, v)| (k.name().to_string(), v.clone()))
1227                    .collect::<IndexMap<_, _>>()
1228            })
1229            .unwrap_or_default();
1230        let partition_path = partition_values.hive_partition_path();
1231        let entry = partition_files
1232            .entry(partition_path)
1233            .or_insert_with(|| (partition_values, 0, vec![]));
1234        let stable_ordinal = entry.1;
1235        entry.1 += 1;
1236
1237        if object_meta.size > target_size.get() {
1238            metrics.total_files_skipped += 1;
1239            continue;
1240        }
1241
1242        entry.2.push(OrderedFileCandidate {
1243            add: file.to_add(),
1244            stable_ordinal,
1245            size_bytes: object_meta.size,
1246        });
1247    }
1248
1249    let mut operations: HashMap<String, (IndexMap<String, Scalar>, Vec<MergeBin>)> = HashMap::new();
1250    for (part, (partition, _, files)) in partition_files {
1251        let (merge_bins, partition_stats) =
1252            plan_compaction_bins_in_stable_order(files, target_size.get());
1253        planner_stats.absorb(&partition_stats);
1254
1255        operations.insert(part, (partition, merge_bins));
1256    }
1257
1258    // Prune merge bins with only 1 file, since they have no effect
1259    for (_, (_, bins)) in operations.iter_mut() {
1260        bins.retain(|bin| {
1261            if bin.len() == 1 {
1262                metrics.total_files_skipped += 1;
1263                false
1264            } else {
1265                true
1266            }
1267        });
1268        planner_stats.max_bin_span_files = planner_stats
1269            .max_bin_span_files
1270            .max(bins.iter().map(MergeBin::len).max().unwrap_or(0));
1271    }
1272    operations.retain(|_, (_, files)| !files.is_empty());
1273
1274    metrics.partitions_optimized = operations.len() as u64;
1275
1276    if operations.is_empty() {
1277        planner_stats.max_bin_span_files = 0;
1278    }
1279
1280    Ok((
1281        OptimizeOperations::Compact(operations),
1282        metrics,
1283        planner_stats,
1284    ))
1285}
1286
1287/// Validates that a z-order column path exists in the schema, supporting nested
1288/// struct fields via dot notation (e.g., "meta.field_a").
1289fn validate_zorder_column(schema: &StructType, column: &str) -> Result<(), DeltaTableError> {
1290    let mut segments = column.split('.').peekable();
1291    let mut current_struct = schema;
1292    while let Some(segment) = segments.next() {
1293        let field = current_struct.field(segment).ok_or_else(|| {
1294            DeltaTableError::Generic(format!(
1295                "Z-order column \"{column}\": field \"{segment}\" not found in schema"
1296            ))
1297        })?;
1298        if segments.peek().is_some() {
1299            match field.data_type() {
1300                DataType::Struct(inner) => current_struct = inner,
1301                _ => {
1302                    return Err(DeltaTableError::Generic(format!(
1303                        "Z-order column \"{column}\": \"{segment}\" is not a struct type"
1304                    )));
1305                }
1306            }
1307        }
1308    }
1309    Ok(())
1310}
1311
1312async fn build_zorder_plan(
1313    log_store: &dyn LogStore,
1314    zorder_columns: Vec<String>,
1315    snapshot: &EagerSnapshot,
1316    partition_keys: &[String],
1317    filters: &[PartitionFilter],
1318) -> Result<(OptimizeOperations, Metrics, PlannerStats), DeltaTableError> {
1319    if zorder_columns.is_empty() {
1320        return Err(DeltaTableError::Generic(
1321            "Z-order requires at least one column".to_string(),
1322        ));
1323    }
1324    let zorder_partition_cols = zorder_columns
1325        .iter()
1326        .filter(|col| partition_keys.contains(col))
1327        .collect_vec();
1328    if !zorder_partition_cols.is_empty() {
1329        return Err(DeltaTableError::Generic(format!(
1330            "Z-order columns cannot be partition columns. Found: {zorder_partition_cols:?}"
1331        )));
1332    }
1333    for col in &zorder_columns {
1334        validate_zorder_column(snapshot.schema().as_ref(), col)?;
1335    }
1336
1337    // For now, just be naive and optimize all files in each selected partition.
1338    let mut metrics = Metrics::default();
1339
1340    let mut partition_files: HashMap<String, (IndexMap<String, Scalar>, MergeBin)> = HashMap::new();
1341
1342    let predicate = if filters.is_empty() {
1343        None
1344    } else {
1345        Some(Arc::new(to_kernel_predicate(
1346            filters,
1347            snapshot.schema().as_ref(),
1348        )?))
1349    };
1350
1351    let mut file_stream = snapshot.file_views(log_store, predicate);
1352    while let Some(file) = file_stream.next().await {
1353        let file = file?;
1354        let partition_values = file
1355            .partition_values()
1356            .map(|v| {
1357                v.fields()
1358                    .iter()
1359                    .zip(v.values().iter())
1360                    .map(|(k, v)| (k.name().to_string(), v.clone()))
1361                    .collect::<IndexMap<_, _>>()
1362            })
1363            .unwrap_or_default();
1364        metrics.total_considered_files += 1;
1365        partition_files
1366            .entry(partition_values.hive_partition_path())
1367            .or_insert_with(|| (partition_values, MergeBin::new()))
1368            .1
1369            .add(file.to_add());
1370        debug!("partition_files inside the zorder plan: {partition_files:?}");
1371    }
1372
1373    let max_bin_span_files = partition_files
1374        .values()
1375        .map(|(_, bin)| bin.len())
1376        .max()
1377        .unwrap_or(0);
1378    let operation = OptimizeOperations::ZOrder(zorder_columns, partition_files);
1379    Ok((
1380        operation,
1381        metrics,
1382        PlannerStats::z_order(max_bin_span_files),
1383    ))
1384}
1385
1386#[cfg(test)]
1387mod compact_planner_tests {
1388    use super::*;
1389    use std::collections::HashMap;
1390
1391    fn candidate(stable_ordinal: usize, size_bytes: u64) -> OrderedFileCandidate {
1392        OrderedFileCandidate {
1393            add: Add {
1394                path: format!("part-{stable_ordinal}.parquet"),
1395                partition_values: HashMap::new(),
1396                size: size_bytes as i64,
1397                modification_time: stable_ordinal as i64,
1398                data_change: false,
1399                stats: None,
1400                tags: None,
1401                deletion_vector: None,
1402                base_row_id: None,
1403                default_row_commit_version: None,
1404                clustering_provider: None,
1405            },
1406            stable_ordinal,
1407            size_bytes,
1408        }
1409    }
1410
1411    fn ordinals(bin: &MergeBin) -> Vec<usize> {
1412        bin.iter()
1413            .map(|add| add.modification_time as usize)
1414            .collect::<Vec<_>>()
1415    }
1416
1417    #[test]
1418    fn test_ordered_compact_bins_are_contiguous() {
1419        let (bins, stats) = plan_compaction_bins_in_stable_order(
1420            vec![
1421                candidate(0, 6),
1422                candidate(1, 3),
1423                candidate(2, 6),
1424                candidate(3, 3),
1425            ],
1426            10,
1427        );
1428
1429        let planned_ordinals = bins.iter().map(ordinals).collect::<Vec<_>>();
1430
1431        assert_eq!(planned_ordinals, vec![vec![0, 1], vec![2, 3]]);
1432        assert_eq!(stats.max_bin_span_files, 2);
1433    }
1434
1435    #[test]
1436    fn test_ordered_compact_bins_do_not_merge_non_adjacent_files() {
1437        let (bins, _) = plan_compaction_bins_in_stable_order(
1438            vec![
1439                candidate(0, 8),
1440                candidate(1, 8),
1441                candidate(2, 2),
1442                candidate(3, 2),
1443            ],
1444            10,
1445        );
1446
1447        let planned_ordinals = bins.iter().map(ordinals).collect::<Vec<_>>();
1448
1449        assert_eq!(planned_ordinals, vec![vec![0], vec![1, 2], vec![3]]);
1450        assert!(
1451            planned_ordinals
1452                .iter()
1453                .all(|bin| { bin.windows(2).all(|window| window[1] == window[0] + 1) })
1454        );
1455    }
1456
1457    #[test]
1458    fn test_ordered_compact_bins_respect_ordinal_gaps() {
1459        let (bins, stats) =
1460            plan_compaction_bins_in_stable_order(vec![candidate(0, 3), candidate(2, 3)], 10);
1461
1462        let planned_ordinals = bins.iter().map(ordinals).collect::<Vec<_>>();
1463
1464        assert_eq!(planned_ordinals, vec![vec![0], vec![2]]);
1465        assert_eq!(stats.max_bin_span_files, 1);
1466    }
1467
1468    #[test]
1469    fn test_ordered_compact_bins_track_span_and_displacement() {
1470        let (_, stats) = plan_compaction_bins_in_stable_order(
1471            vec![
1472                candidate(0, 3),
1473                candidate(1, 3),
1474                candidate(2, 3),
1475                candidate(3, 9),
1476            ],
1477            10,
1478        );
1479
1480        assert_eq!(stats.planner_strategy, PlannerStrategy::PreserveLocality);
1481        assert!(stats.preserved_stable_order);
1482        assert_eq!(stats.max_bin_span_files, 3);
1483    }
1484
1485    #[test]
1486    fn test_optimize_input_target_size_must_fit_i64() {
1487        let input = OptimizeInput {
1488            target_size: std::num::NonZeroU64::new(i64::MAX as u64 + 1).unwrap(),
1489            predicate: None,
1490        };
1491
1492        let err = crate::protocol::DeltaOperation::try_from(input).unwrap_err();
1493        assert!(err.to_string().contains("optimize target_size"));
1494        assert!(err.to_string().contains("i64::MAX"));
1495    }
1496}
1497
1498pub(super) mod util {
1499    use super::*;
1500    use futures::Future;
1501    use tokio::task::JoinError;
1502
1503    pub async fn flatten_join_error<T, E>(
1504        future: impl Future<Output = Result<Result<T, E>, JoinError>>,
1505    ) -> Result<T, DeltaTableError>
1506    where
1507        E: Into<DeltaTableError>,
1508    {
1509        match future.await {
1510            Ok(Ok(result)) => Ok(result),
1511            Ok(Err(error)) => Err(error.into()),
1512            Err(error) => Err(DeltaTableError::GenericError {
1513                source: Box::new(error),
1514            }),
1515        }
1516    }
1517}
1518
1519/// Z-order utilities
1520pub(super) mod zorder {
1521    use super::*;
1522
1523    use arrow::buffer::{Buffer, OffsetBuffer, ScalarBuffer};
1524    use arrow_array::{Array, ArrayRef, BinaryArray};
1525    use arrow_buffer::bit_util::{get_bit_raw, set_bit_raw, unset_bit_raw};
1526    use arrow_row::{Row, RowConverter, SortField};
1527    use arrow_schema::ArrowError;
1528    // use arrow_schema::Schema as ArrowSchema;
1529
1530    pub use self::datafusion::ZOrderExecContext;
1531
1532    pub(super) mod datafusion {
1533        use super::*;
1534        use url::Url;
1535
1536        use ::datafusion::common::DataFusionError;
1537        use ::datafusion::logical_expr::{
1538            ColumnarValue, ScalarFunctionArgs, ScalarUDF, ScalarUDFImpl, Signature, TypeSignature,
1539            Volatility,
1540        };
1541        use ::datafusion::prelude::SessionContext;
1542        use arrow_schema::DataType;
1543        use itertools::Itertools;
1544        use std::any::Any;
1545
1546        pub const ZORDER_UDF_NAME: &str = "zorder_key";
1547
1548        pub struct ZOrderExecContext {
1549            pub columns: Arc<[String]>,
1550            pub ctx: SessionContext,
1551        }
1552
1553        impl ZOrderExecContext {
1554            pub fn new(
1555                columns: Vec<String>,
1556                session: SessionState,
1557                object_store_ref: ObjectStoreRef,
1558            ) -> Result<Self, DataFusionError> {
1559                let columns = columns.into();
1560
1561                let ctx = SessionContext::new_with_state(session);
1562                ctx.register_udf(ScalarUDF::from(datafusion::ZOrderUDF));
1563                ctx.register_object_store(&Url::parse("delta-rs://").unwrap(), object_store_ref);
1564                Ok(Self { columns, ctx })
1565            }
1566        }
1567
1568        // DataFusion UDF impl for zorder_key
1569        #[derive(Debug, Hash, PartialEq, Eq)]
1570        pub struct ZOrderUDF;
1571
1572        impl ScalarUDFImpl for ZOrderUDF {
1573            fn as_any(&self) -> &dyn Any {
1574                self
1575            }
1576
1577            fn name(&self) -> &str {
1578                ZORDER_UDF_NAME
1579            }
1580
1581            fn signature(&self) -> &Signature {
1582                static SIGNATURE: std::sync::LazyLock<Signature> =
1583                    std::sync::LazyLock::new(|| Signature {
1584                        type_signature: TypeSignature::VariadicAny,
1585                        volatility: Volatility::Immutable,
1586                        parameter_names: Some(vec![]),
1587                    });
1588                &SIGNATURE
1589            }
1590
1591            fn return_type(&self, _arg_types: &[DataType]) -> Result<DataType, DataFusionError> {
1592                Ok(DataType::Binary)
1593            }
1594
1595            fn invoke_with_args(
1596                &self,
1597                args: ScalarFunctionArgs,
1598            ) -> ::datafusion::common::Result<ColumnarValue> {
1599                zorder_key_datafusion(&args.args)
1600            }
1601        }
1602
1603        /// Datafusion zorder UDF body
1604        fn zorder_key_datafusion(
1605            columns: &[ColumnarValue],
1606        ) -> Result<ColumnarValue, DataFusionError> {
1607            debug!("zorder_key_datafusion: {columns:#?}");
1608            let length = columns
1609                .iter()
1610                .map(|col| match col {
1611                    ColumnarValue::Array(array) => array.len(),
1612                    ColumnarValue::Scalar(_) => 1,
1613                })
1614                .max()
1615                .ok_or(DataFusionError::NotImplemented(
1616                    "z-order on zero columns.".to_string(),
1617                ))?;
1618            let columns: Vec<ArrayRef> = columns
1619                .iter()
1620                .map(|col| col.clone().into_array(length))
1621                .try_collect()?;
1622            let array = zorder_key(&columns)?;
1623            Ok(ColumnarValue::Array(array))
1624        }
1625
1626        #[cfg(test)]
1627        mod tests {
1628            use super::*;
1629            use ::datafusion::assert_batches_eq;
1630            use arrow_array::{Int32Array, StringArray};
1631            use arrow_ord::sort::sort_to_indices;
1632            use arrow_schema::Field;
1633            use arrow_select::take::take;
1634            use rand::RngExt;
1635
1636            #[test]
1637            fn test_order() {
1638                let int: ArrayRef = Arc::new(Int32Array::from(vec![1, 2, 3, 4, 5]));
1639                let str: ArrayRef = Arc::new(StringArray::from(vec![
1640                    Some("a"),
1641                    Some("x"),
1642                    Some("a"),
1643                    Some("x"),
1644                    None,
1645                ]));
1646                let int_large: ArrayRef = Arc::new(Int32Array::from(vec![10000, 2000, 300, 40, 5]));
1647                let batch = RecordBatch::try_from_iter(vec![
1648                    ("int", int),
1649                    ("str", str),
1650                    ("int_large", int_large),
1651                ])
1652                .unwrap();
1653
1654                let expected_1 = vec![
1655                    "+-----+-----+-----------+",
1656                    "| int | str | int_large |",
1657                    "+-----+-----+-----------+",
1658                    "| 1   | a   | 10000     |",
1659                    "| 2   | x   | 2000      |",
1660                    "| 3   | a   | 300       |",
1661                    "| 4   | x   | 40        |",
1662                    "| 5   |     | 5         |",
1663                    "+-----+-----+-----------+",
1664                ];
1665                let expected_2 = vec![
1666                    "+-----+-----+-----------+",
1667                    "| int | str | int_large |",
1668                    "+-----+-----+-----------+",
1669                    "| 5   |     | 5         |",
1670                    "| 1   | a   | 10000     |",
1671                    "| 3   | a   | 300       |",
1672                    "| 2   | x   | 2000      |",
1673                    "| 4   | x   | 40        |",
1674                    "+-----+-----+-----------+",
1675                ];
1676                let expected_3 = vec![
1677                    "+-----+-----+-----------+",
1678                    "| int | str | int_large |",
1679                    "+-----+-----+-----------+",
1680                    "| 5   |     | 5         |",
1681                    "| 4   | x   | 40        |",
1682                    "| 2   | x   | 2000      |",
1683                    "| 3   | a   | 300       |",
1684                    "| 1   | a   | 10000     |",
1685                    "+-----+-----+-----------+",
1686                ];
1687
1688                let expected = [expected_1, expected_2, expected_3];
1689
1690                let indices = Int32Array::from(shuffled_indices().to_vec());
1691                let shuffled_columns = batch
1692                    .columns()
1693                    .iter()
1694                    .map(|c| take(c, &indices, None).unwrap())
1695                    .collect::<Vec<_>>();
1696                let shuffled_batch =
1697                    RecordBatch::try_new(batch.schema(), shuffled_columns).unwrap();
1698
1699                for i in 1..=batch.num_columns() {
1700                    let columns = (0..i)
1701                        .map(|idx| shuffled_batch.column(idx).clone())
1702                        .collect::<Vec<_>>();
1703
1704                    let order_keys = zorder_key(&columns).unwrap();
1705                    let indices = sort_to_indices(order_keys.as_ref(), None, None).unwrap();
1706                    let sorted_columns = shuffled_batch
1707                        .columns()
1708                        .iter()
1709                        .map(|c| take(c, &indices, None).unwrap())
1710                        .collect::<Vec<_>>();
1711                    let sorted_batch =
1712                        RecordBatch::try_new(batch.schema(), sorted_columns).unwrap();
1713
1714                    assert_batches_eq!(expected[i - 1], &[sorted_batch]);
1715                }
1716            }
1717            fn shuffled_indices() -> [i32; 5] {
1718                let mut rng = rand::rng();
1719                let mut array = [0, 1, 2, 3, 4];
1720                for i in (1..array.len()).rev() {
1721                    let j = rng.random_range(0..=i);
1722                    array.swap(i, j);
1723                }
1724                array
1725            }
1726
1727            #[tokio::test]
1728            async fn test_zorder_mixed_case() {
1729                use arrow_schema::Schema as ArrowSchema;
1730                let schema = Arc::new(ArrowSchema::new(vec![
1731                    Field::new("moDified", DataType::Utf8, true),
1732                    Field::new("ID", DataType::Utf8, true),
1733                    Field::new("vaLue", DataType::Int32, true),
1734                ]));
1735
1736                let batch = RecordBatch::try_new(
1737                    schema.clone(),
1738                    vec![
1739                        Arc::new(arrow::array::StringArray::from(vec![
1740                            "2021-02-01",
1741                            "2021-02-01",
1742                            "2021-02-02",
1743                            "2021-02-02",
1744                        ])),
1745                        Arc::new(arrow::array::StringArray::from(vec!["A", "B", "C", "D"])),
1746                        Arc::new(arrow::array::Int32Array::from(vec![1, 10, 20, 100])),
1747                    ],
1748                )
1749                .unwrap();
1750                // write some data
1751                let table = crate::DeltaTable::new_in_memory()
1752                    .write(vec![batch.clone()])
1753                    .with_save_mode(crate::protocol::SaveMode::Append)
1754                    .await
1755                    .unwrap();
1756
1757                let res = table
1758                    .optimize()
1759                    .with_type(OptimizeType::ZOrder(vec!["moDified".into()]))
1760                    .await;
1761                assert!(res.is_ok());
1762            }
1763
1764            /// Issue <https://github.com/delta-io/delta-rs/issues/2834>
1765            #[tokio::test]
1766            async fn test_zorder_space_in_partition_value() {
1767                use arrow_schema::Schema as ArrowSchema;
1768                let _ = pretty_env_logger::try_init();
1769                let schema = Arc::new(ArrowSchema::new(vec![
1770                    Field::new("modified", DataType::Utf8, true),
1771                    Field::new("country", DataType::Utf8, true),
1772                    Field::new("value", DataType::Int32, true),
1773                ]));
1774
1775                let batch = RecordBatch::try_new(
1776                    schema.clone(),
1777                    vec![
1778                        Arc::new(arrow::array::StringArray::from(vec![
1779                            "2021-02-01",
1780                            "2021-02-01",
1781                            "2021-02-02",
1782                            "2021-02-02",
1783                        ])),
1784                        Arc::new(arrow::array::StringArray::from(vec![
1785                            "Germany",
1786                            "China",
1787                            "Canada",
1788                            "Dominican Republic",
1789                        ])),
1790                        Arc::new(arrow::array::Int32Array::from(vec![1, 10, 20, 100])),
1791                        //Arc::new(arrow::array::StringArray::from(vec!["Dominican Republic"])),
1792                        //Arc::new(arrow::array::Int32Array::from(vec![100])),
1793                    ],
1794                )
1795                .unwrap();
1796                // write some data
1797                let table = DeltaTable::new_in_memory()
1798                    .write(vec![batch.clone()])
1799                    .with_partition_columns(vec!["country"])
1800                    .with_save_mode(crate::protocol::SaveMode::Overwrite)
1801                    .await
1802                    .unwrap();
1803
1804                let res = table
1805                    .optimize()
1806                    .with_type(OptimizeType::ZOrder(vec!["modified".into()]))
1807                    .await;
1808                assert!(res.is_ok(), "Failed to optimize: {res:#?}");
1809            }
1810
1811            #[tokio::test]
1812            async fn test_zorder_space_in_partition_value_garbage() {
1813                use arrow_schema::Schema as ArrowSchema;
1814                let _ = pretty_env_logger::try_init();
1815                let schema = Arc::new(ArrowSchema::new(vec![
1816                    Field::new("modified", DataType::Utf8, true),
1817                    Field::new("country", DataType::Utf8, true),
1818                    Field::new("value", DataType::Int32, true),
1819                ]));
1820
1821                let batch = RecordBatch::try_new(
1822                    schema.clone(),
1823                    vec![
1824                        Arc::new(arrow::array::StringArray::from(vec![
1825                            "2021-02-01",
1826                            "2021-02-01",
1827                            "2021-02-02",
1828                            "2021-02-02",
1829                        ])),
1830                        Arc::new(arrow::array::StringArray::from(vec![
1831                            "Germany", "China", "Canada", "USA$$!",
1832                        ])),
1833                        Arc::new(arrow::array::Int32Array::from(vec![1, 10, 20, 100])),
1834                    ],
1835                )
1836                .unwrap();
1837                // write some data
1838                let table = DeltaTable::new_in_memory()
1839                    .write(vec![batch.clone()])
1840                    .with_partition_columns(vec!["country"])
1841                    .with_save_mode(crate::protocol::SaveMode::Overwrite)
1842                    .await
1843                    .unwrap();
1844
1845                let res = table
1846                    .optimize()
1847                    .with_type(OptimizeType::ZOrder(vec!["modified".into()]))
1848                    .await;
1849                assert!(res.is_ok(), "Failed to optimize: {res:#?}");
1850            }
1851        }
1852    }
1853
1854    /// Creates a new binary array containing the zorder keys for the given columns
1855    ///
1856    /// Each value is 16 bytes * number of columns. Each column is converted into
1857    /// its row binary representation, and then the first 16 bytes are taken.
1858    /// These truncated values are interleaved in the array values.
1859    pub fn zorder_key(columns: &[ArrayRef]) -> Result<ArrayRef, ArrowError> {
1860        if columns.is_empty() {
1861            return Err(ArrowError::InvalidArgumentError(
1862                "Cannot zorder empty columns".to_string(),
1863            ));
1864        }
1865
1866        // length is length of first array or 1 if all scalars
1867        let out_length = columns[0].len();
1868
1869        if columns.iter().any(|col| col.len() != out_length) {
1870            return Err(ArrowError::InvalidArgumentError(
1871                "All columns must have the same length".to_string(),
1872            ));
1873        }
1874
1875        // We are taking 128 bits (16 bytes) from each value. Shorter values will be padded.
1876        let value_size: usize = columns.len() * 16;
1877
1878        // Initialize with zeros
1879        let mut out: Vec<u8> = vec![0; out_length * value_size];
1880
1881        for (col_pos, col) in columns.iter().enumerate() {
1882            set_bits_for_column(col.clone(), col_pos, columns.len(), &mut out)?;
1883        }
1884
1885        let offsets = (0..=out_length)
1886            .map(|i| (i * value_size) as i32)
1887            .collect::<Vec<i32>>();
1888
1889        let out_arr = BinaryArray::try_new(
1890            OffsetBuffer::new(ScalarBuffer::from(offsets)),
1891            Buffer::from_vec(out),
1892            None,
1893        )?;
1894
1895        Ok(Arc::new(out_arr))
1896    }
1897
1898    /// Given an input array, will set the bits in the output array
1899    ///
1900    /// Arguments:
1901    /// * `input` - The input array
1902    /// * `col_pos` - The position of the column. Used to determine position
1903    ///   when interleaving.
1904    /// * `num_columns` - The number of columns in the input array. Used to
1905    ///   determine offset when interleaving.
1906    fn set_bits_for_column(
1907        input: ArrayRef,
1908        col_pos: usize,
1909        num_columns: usize,
1910        out: &mut Vec<u8>,
1911    ) -> Result<(), ArrowError> {
1912        // Convert array to rows
1913        let converter = RowConverter::new(vec![SortField::new(input.data_type().clone())])?;
1914        let rows = converter.convert_columns(&[input])?;
1915
1916        for (row_i, row) in rows.iter().enumerate() {
1917            // How many bytes to get to this row's out position
1918            let row_offset = row_i * num_columns * 16;
1919            for bit_i in 0..128 {
1920                let bit = row.get_bit(bit_i);
1921                // Position of bit within the value. We place a value every
1922                // `num_columns` bits, offset by `col_pos` when interleaving.
1923                // So if there are 3 columns, and we are the second column, then
1924                // we place values at index: 1, 4, 7, 10, etc.
1925                let bit_pos = (bit_i * num_columns) + col_pos;
1926                let out_pos = (row_offset * 8) + bit_pos;
1927                // Safety: we pre-sized the output vector in the outer function
1928                if bit {
1929                    unsafe { set_bit_raw(out.as_mut_ptr(), out_pos) };
1930                } else {
1931                    unsafe { unset_bit_raw(out.as_mut_ptr(), out_pos) };
1932                }
1933            }
1934        }
1935
1936        Ok(())
1937    }
1938
1939    trait RowBitUtil {
1940        fn get_bit(&self, bit_i: usize) -> bool;
1941    }
1942
1943    impl RowBitUtil for Row<'_> {
1944        /// Get the bit at the given index, or just give false if the index is out of bounds
1945        fn get_bit(&self, bit_i: usize) -> bool {
1946            let byte_i = bit_i / 8;
1947            let bytes = self.as_ref();
1948            if byte_i >= bytes.len() {
1949                return false;
1950            }
1951            // Safety: we just did a bounds check above
1952            unsafe { get_bit_raw(bytes.as_ptr(), bit_i) }
1953        }
1954    }
1955
1956    #[cfg(test)]
1957    mod test {
1958        use arrow_array::{
1959            StringArray, UInt8Array, cast::as_generic_binary_array, new_empty_array,
1960        };
1961        use arrow_schema::DataType;
1962
1963        use super::*;
1964        use crate::ensure_table_uri;
1965
1966        #[test]
1967        fn test_rejects_no_columns() {
1968            let columns = vec![];
1969            let result = zorder_key(&columns);
1970            assert!(result.is_err());
1971        }
1972
1973        #[test]
1974        fn test_handles_no_rows() {
1975            let columns: Vec<ArrayRef> = vec![
1976                Arc::new(new_empty_array(&DataType::Int64)),
1977                Arc::new(new_empty_array(&DataType::Utf8)),
1978            ];
1979            let result = zorder_key(columns.as_slice());
1980            assert!(result.is_ok());
1981            let result = result.unwrap();
1982            assert_eq!(result.len(), 0);
1983        }
1984
1985        #[test]
1986        fn test_basics() {
1987            let columns: Vec<ArrayRef> = vec![
1988                // Small strings
1989                Arc::new(StringArray::from(vec![Some("a"), Some("b"), None])),
1990                // Strings of various sizes
1991                Arc::new(StringArray::from(vec![
1992                    "delta-rs: A native Rust library for Delta Lake, with bindings into Python",
1993                    "cat",
1994                    "",
1995                ])),
1996                Arc::new(UInt8Array::from(vec![Some(1), Some(4), None])),
1997            ];
1998            let result = zorder_key(columns.as_slice()).unwrap();
1999            assert_eq!(result.len(), 3);
2000            assert_eq!(result.data_type(), &DataType::Binary);
2001            assert_eq!(result.null_count(), 0);
2002
2003            let data: &BinaryArray = as_generic_binary_array(result.as_ref());
2004            assert_eq!(data.value_data().len(), 3 * 16 * 3);
2005            assert!(data.iter().all(|x| x.unwrap().len() == 3 * 16));
2006        }
2007
2008        #[tokio::test]
2009        async fn works_on_spark_table() {
2010            use tempfile::TempDir;
2011            // Create a temporary directory
2012            let tmp_dir = TempDir::new().expect("Failed to make temp dir");
2013            let table_name = "delta-1.2.1-only-struct-stats";
2014
2015            // Copy recursively from the test data directory to the temporary directory
2016            let source_path = format!("../test/tests/data/{table_name}");
2017            fs_extra::dir::copy(source_path, tmp_dir.path(), &Default::default()).unwrap();
2018
2019            let table_uri =
2020                ensure_table_uri(tmp_dir.path().join(table_name).to_str().unwrap()).unwrap();
2021            // Run optimize
2022            let (_, metrics) = DeltaTable::try_from_url(table_uri)
2023                .await
2024                .unwrap()
2025                .optimize()
2026                .await
2027                .unwrap();
2028
2029            // Verify it worked
2030            assert_eq!(metrics.num_files_added, 1);
2031        }
2032    }
2033}