1use 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#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
70#[serde(rename_all = "camelCase")]
71pub enum PlannerStrategy {
72 #[default]
74 UnknownLegacy,
75 PreserveLocality,
77 ZOrder,
79}
80
81#[derive(Default, Debug, PartialEq, Clone, Serialize, Deserialize)]
83#[serde(rename_all = "camelCase", from = "MetricsSerde")]
84pub struct Metrics {
85 pub num_files_added: u64,
87 pub num_files_removed: u64,
89 #[serde(
91 serialize_with = "serialize_metric_details",
92 deserialize_with = "deserialize_metric_details"
93 )]
94 pub files_added: MetricDetails,
95 #[serde(
97 serialize_with = "serialize_metric_details",
98 deserialize_with = "deserialize_metric_details"
99 )]
100 pub files_removed: MetricDetails,
101 pub partitions_optimized: u64,
103 pub num_batches: u64,
105 pub total_considered_files: usize,
107 pub total_files_skipped: usize,
109 pub preserve_insertion_order: bool,
111 pub planner_strategy: PlannerStrategy,
113 pub preserved_stable_order: bool,
115 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
167fn 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
175fn 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#[derive(Debug, PartialEq, Clone, Serialize, Deserialize)]
187#[serde(rename_all = "camelCase")]
188pub struct MetricDetails {
189 pub avg: f64,
191 pub max: i64,
193 pub min: i64,
195 pub total_files: usize,
197 pub total_size: i64,
199}
200
201impl MetricDetails {
202 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 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)]
220pub struct PartialMetrics {
222 pub num_files_added: u64,
224 pub num_files_removed: u64,
226 pub files_added: MetricDetails,
228 pub files_removed: MetricDetails,
230 pub num_batches: u64,
232}
233
234impl Metrics {
235 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#[derive(Debug)]
266pub enum OptimizeType {
267 Compact,
269 ZOrder(Vec<String>),
271}
272
273pub struct OptimizeBuilder<'a> {
278 snapshot: Option<EagerSnapshot>,
280 log_store: LogStoreRef,
282 filters: &'a [PartitionFilter],
284 target_size: Option<NonZeroU64>,
286 writer_properties: Option<WriterProperties>,
288 commit_properties: CommitProperties,
290 max_concurrent_tasks: usize,
292 optimize_type: OptimizeType,
294 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 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 pub fn with_type(mut self, optimize_type: OptimizeType) -> Self {
331 self.optimize_type = optimize_type;
332 self
333 }
334
335 pub fn with_filters(mut self, filters: &'a [PartitionFilter]) -> Self {
337 self.filters = filters;
338 self
339 }
340
341 pub fn with_target_size(mut self, target: NonZeroU64) -> Self {
343 self.target_size = Some(target);
344 self
345 }
346
347 pub fn with_writer_properties(mut self, writer_properties: WriterProperties) -> Self {
349 self.writer_properties = Some(writer_properties);
350 self
351 }
352
353 pub fn with_commit_properties(mut self, commit_properties: CommitProperties) -> Self {
355 self.commit_properties = commit_properties;
356 self
357 }
358
359 #[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 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 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 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 pub fn with_session_state(mut self, session: Arc<dyn Session>) -> Self {
396 self.session = Some(session);
397 self
398 }
399
400 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
502fn create_remove(add: &Add) -> Action {
504 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#[derive(Debug)]
527enum OptimizeOperations {
528 Compact(HashMap<String, (IndexMap<String, Scalar>, Vec<MergeBin>)>),
533 ZOrder(
535 Vec<String>,
536 HashMap<String, (IndexMap<String, Scalar>, MergeBin)>,
537 ),
538 }
540
541impl Default for OptimizeOperations {
542 fn default() -> Self {
543 OptimizeOperations::Compact(HashMap::new())
544 }
545}
546
547#[derive(Debug)]
548pub struct MergePlan {
550 operations: OptimizeOperations,
551 metrics: Metrics,
553 planner_stats: PlannerStats,
555 task_parameters: Arc<MergeTaskParameters>,
557 read_table_version: Version,
559 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#[derive(Debug)]
594pub struct MergeTaskParameters {
595 file_schema: SchemaRef,
597 writer_properties: WriterProperties,
599 input_parameters: OptimizeInput,
601 num_indexed_cols: DataSkippingNumIndexedCols,
603 stats_columns: Option<Vec<String>>,
605}
606
607type 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 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 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 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 let writer_config = PartitionWriterConfig::try_new(
691 task_parameters.file_schema.clone(),
692 partition_values.clone(),
693 Some(task_parameters.writer_properties.clone()),
694 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 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 #[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 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 let mut actions = vec![];
922
923 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#[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#[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 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 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
1287fn 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 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
1519pub(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 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 #[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 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 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 #[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 ],
1794 )
1795 .unwrap();
1796 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 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 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 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 let value_size: usize = columns.len() * 16;
1877
1878 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 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 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 let row_offset = row_i * num_columns * 16;
1919 for bit_i in 0..128 {
1920 let bit = row.get_bit(bit_i);
1921 let bit_pos = (bit_i * num_columns) + col_pos;
1926 let out_pos = (row_offset * 8) + bit_pos;
1927 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 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 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 Arc::new(StringArray::from(vec![Some("a"), Some("b"), None])),
1990 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 let tmp_dir = TempDir::new().expect("Failed to make temp dir");
2013 let table_name = "delta-1.2.1-only-struct-stats";
2014
2015 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 let (_, metrics) = DeltaTable::try_from_url(table_uri)
2023 .await
2024 .unwrap()
2025 .optimize()
2026 .await
2027 .unwrap();
2028
2029 assert_eq!(metrics.num_files_added, 1);
2031 }
2032 }
2033}