buoyant_kernel/lib.rs
1//! # Delta Kernel
2//!
3//! Delta-kernel-rs is an experimental [Delta](https://github.com/delta-io/delta/) implementation
4//! focused on interoperability with a wide range of query engines. It supports reads and
5//! (experimental) writes (only blind appends in the write path currently). This library defines a
6//! number of traits which must be implemented to provide a working delta implementation. They are
7//! detailed below. There is a provided "default engine" that implements all these traits and can
8//! be used to ease integration work. See [`DefaultEngine`](engine/default/index.html) for more
9//! information.
10//!
11//! A full `rust` example for reading table data using the default engine can be found in the
12//! [read-table-single-threaded] example (and for a more complex multi-threaded reader see the
13//! [read-table-multi-threaded] example). An example for reading the table changes for a table
14//! using the default engine can be found in the [read-table-changes] example. The [write-table]
15//! example demonstrates how to write data to a Delta table using the default engine.
16//!
17//! [read-table-single-threaded]:
18//! https://github.com/delta-io/delta-kernel-rs/tree/main/kernel/examples/read-table-single-threaded
19//! [read-table-multi-threaded]:
20//! https://github.com/delta-io/delta-kernel-rs/tree/main/kernel/examples/read-table-multi-threaded
21//! [read-table-changes]:
22//! https://github.com/delta-io/delta-kernel-rs/tree/main/kernel/examples/read-table-changes
23//! [write-table]:
24//! https://github.com/delta-io/delta-kernel-rs/tree/main/kernel/examples/write-table
25//!
26//! # Engine trait
27//!
28//! The [`Engine`] trait allows connectors to bring their own implementation of functionality such
29//! as reading parquet files, listing files in a file system, parsing a JSON string etc. This
30//! trait exposes methods to get sub-engines which expose the core functionalities customizable by
31//! connectors.
32//!
33//! ## Expression handling
34//!
35//! Expression handling is done via the [`EvaluationHandler`], which in turn allows the creation of
36//! [`ExpressionEvaluator`]s. These evaluators are created for a specific predicate [`Expression`]
37//! and allow evaluation of that predicate for a specific batch of data.
38//!
39//! ## File system interactions
40//!
41//! Delta Kernel needs to perform some basic operations against file systems like listing and
42//! reading files. These interactions are encapsulated in the [`StorageHandler`] trait.
43//! Implementers must take care that all assumptions on the behavior of the functions - like sorted
44//! results - are respected.
45//!
46//! ## Reading log and data files
47//!
48//! Delta Kernel requires the capability to read and write json files and read parquet files, which
49//! is exposed via the [`JsonHandler`] and [`ParquetHandler`] respectively. When reading files,
50//! connectors are asked to provide the context information they require to execute the actual
51//! operation. This is done by invoking methods on the [`StorageHandler`] trait.
52
53#![cfg_attr(all(doc, NIGHTLY_CHANNEL), feature(doc_cfg))]
54#![warn(
55 unreachable_pub,
56 trivial_numeric_casts,
57 unused_extern_crates,
58 rust_2018_idioms,
59 rust_2021_compatibility,
60 clippy::unwrap_used,
61 clippy::expect_used,
62 clippy::panic
63)]
64// we re-allow panics in tests
65#![cfg_attr(test, allow(clippy::unwrap_used, clippy::expect_used, clippy::panic))]
66
67/// This `extern crate` declaration allows the macro to reliably refer to
68/// `delta_kernel::schema::DataType` no matter which crate invokes it. Without that, `delta_kernel`
69/// cannot invoke the macro because `delta_kernel` is an unknown crate identifier (you have to use
70/// `crate` instead). We could make the macro use `crate::schema::DataType` instead, but then the
71/// macro is useless outside the `delta_kernel` crate.
72// TODO: when running `cargo package -p delta_kernel` this gives 'unused' warning - #1095
73#[allow(unused_extern_crates)]
74extern crate self as delta_kernel;
75
76use std::any::Any;
77use std::fs::DirEntry;
78use std::sync::Arc;
79use std::time::SystemTime;
80use std::{cmp::Ordering, ops::Range};
81
82use bytes::Bytes;
83use url::Url;
84
85use self::schema::{DataType, SchemaRef};
86
87mod action_reconciliation;
88pub mod actions;
89pub mod checkpoint;
90pub mod committer;
91// Public under test-utils so integration tests can inspect CRC state via Snapshot::get_current_crc_if_loaded_for_testing.
92#[cfg(feature = "test-utils")]
93pub mod crc;
94#[cfg(not(feature = "test-utils"))]
95pub(crate) mod crc;
96pub mod engine_data;
97pub mod error;
98pub mod expressions;
99mod log_compaction;
100mod log_path;
101mod log_reader;
102pub mod metrics;
103pub mod partition;
104pub mod scan;
105pub mod schema;
106pub mod snapshot;
107pub mod table_changes;
108pub mod table_configuration;
109pub mod table_features;
110pub mod table_properties;
111pub mod transaction;
112pub mod transforms;
113
114pub use log_path::LogPath;
115
116mod row_tracking;
117
118pub(crate) mod clustering;
119
120mod arrow_compat;
121#[cfg(any(feature = "arrow-57", feature = "arrow-58"))]
122pub use arrow_compat::*;
123
124#[cfg(feature = "internal-api")]
125pub mod column_trie;
126#[cfg(not(feature = "internal-api"))]
127pub(crate) mod column_trie;
128pub mod kernel_predicates;
129pub(crate) mod utils;
130
131#[cfg(feature = "internal-api")]
132pub use utils::try_parse_uri;
133
134// for the below modules, we cannot introduce a macro to clean this up. rustfmt doesn't follow into
135// macros, and so will not format the files associated with these modules if we get too clever. see:
136// https://github.com/rust-lang/rustfmt/issues/3253
137
138#[cfg(feature = "internal-api")]
139pub mod path;
140#[cfg(not(feature = "internal-api"))]
141pub(crate) mod path;
142
143#[cfg(feature = "internal-api")]
144pub mod log_replay;
145#[cfg(not(feature = "internal-api"))]
146pub(crate) mod log_replay;
147
148#[cfg(feature = "internal-api")]
149pub mod log_segment;
150#[cfg(not(feature = "internal-api"))]
151pub(crate) mod log_segment;
152
153#[cfg(feature = "internal-api")]
154pub mod last_checkpoint_hint;
155#[cfg(not(feature = "internal-api"))]
156pub(crate) mod last_checkpoint_hint;
157
158pub(crate) mod log_segment_files;
159
160#[cfg(feature = "internal-api")]
161pub mod history_manager;
162#[cfg(not(feature = "internal-api"))]
163pub(crate) mod history_manager;
164
165#[cfg(feature = "internal-api")]
166pub mod parallel;
167#[cfg(not(feature = "internal-api"))]
168pub(crate) mod parallel;
169
170pub use action_reconciliation::{ActionReconciliationIterator, ActionReconciliationIteratorState};
171pub use delta_kernel_derive;
172use delta_kernel_derive::internal_api;
173pub use engine_data::{
174 EngineData, FilteredEngineData, FilteredRowVisitor, GetData, RowIndexIterator, RowVisitor,
175};
176pub use error::{DeltaResult, Error};
177pub use expressions::{Expression, ExpressionRef, Predicate, PredicateRef};
178pub use log_compaction::{should_compact, LogCompactionWriter};
179pub use metrics::MetricsReporter;
180pub use snapshot::Snapshot;
181pub use snapshot::SnapshotRef;
182
183use expressions::literal_expression_transform;
184use expressions::Scalar;
185use schema::{StructField, StructType};
186
187#[cfg(any(
188 feature = "default-engine-native-tls",
189 feature = "default-engine-rustls",
190 feature = "arrow-conversion"
191))]
192pub mod engine;
193
194/// Delta table version is 8 byte unsigned int
195pub type Version = u64;
196
197/// Sentinel version indicating a pre-commit state (table does not exist yet).
198/// Used for create-table transactions before the first commit.
199pub const PRE_COMMIT_VERSION: Version = u64::MAX;
200
201pub type FileSize = u64;
202pub type FileIndex = u64;
203
204/// A specification for a range of bytes to read from a file location
205pub type FileSlice = (Url, Option<Range<FileIndex>>);
206
207/// Data read from a Delta table file and the corresponding scan file information.
208pub type FileDataReadResult = (FileMeta, Box<dyn EngineData>);
209
210/// An iterator of data read from specified files
211pub type FileDataReadResultIterator =
212 Box<dyn Iterator<Item = DeltaResult<Box<dyn EngineData>>> + Send>;
213
214/// The metadata that describes an object.
215#[derive(Debug, Clone, PartialEq, Eq)]
216pub struct FileMeta {
217 /// The fully qualified path to the object
218 pub location: Url,
219 /// The last modified time as milliseconds since unix epoch
220 pub last_modified: i64,
221 /// The size in bytes of the object
222 pub size: FileSize,
223}
224
225impl Ord for FileMeta {
226 fn cmp(&self, other: &Self) -> Ordering {
227 self.location.cmp(&other.location)
228 }
229}
230
231impl PartialOrd for FileMeta {
232 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
233 Some(self.cmp(other))
234 }
235}
236
237impl TryFrom<DirEntry> for FileMeta {
238 type Error = Error;
239
240 fn try_from(ent: DirEntry) -> DeltaResult<FileMeta> {
241 let metadata = ent.metadata()?;
242 let last_modified = metadata
243 .modified()?
244 .duration_since(SystemTime::UNIX_EPOCH)
245 .map_err(|_| Error::generic("Failed to convert file timestamp to milliseconds"))?;
246 let location = Url::from_file_path(ent.path())
247 .map_err(|_| Error::generic(format!("Invalid path: {:?}", ent.path())))?;
248 let last_modified = last_modified.as_millis().try_into().map_err(|_| {
249 Error::generic(format!(
250 "Failed to convert file modification time {:?} into i64",
251 last_modified.as_millis()
252 ))
253 })?;
254 Ok(FileMeta {
255 location,
256 last_modified,
257 size: metadata.len(),
258 })
259 }
260}
261
262impl FileMeta {
263 /// Create a new instance of `FileMeta`
264 pub fn new(location: Url, last_modified: i64, size: u64) -> Self {
265 Self {
266 location,
267 last_modified,
268 size,
269 }
270 }
271}
272
273/// Extension trait that makes it easier to work with traits objects that implement [`Any`],
274/// implemented automatically for any type that satisfies `Any`, `Send`, and `Sync`. In particular,
275/// given some `trait T: Any + Send + Sync`, it allows upcasting `T` to `dyn Any + Send + Sync`,
276/// which in turn allows downcasting the result to a concrete type.
277///
278/// For example, the following code will compile:
279///
280/// ```
281/// # use buoyant_kernel as delta_kernel;
282/// # use delta_kernel::AsAny;
283/// # use std::any::Any;
284/// # use std::sync::Arc;
285/// trait Foo : AsAny {}
286/// struct Bar;
287/// impl Foo for Bar {}
288///
289/// let f: Arc<dyn Foo> = Arc::new(Bar);
290/// let a: Arc<dyn Any + Send + Sync> = f.as_any();
291/// let b: Arc<Bar> = a.downcast().unwrap();
292/// ```
293///
294/// In contrast, very similar code that relies only on `Any` would fail to compile:
295///
296/// ```fail_compile
297/// # use std::any::Any;
298/// # use std::sync::Arc;
299/// trait Foo: Any + Send + Sync {}
300///
301/// struct Bar;
302/// impl Foo for Bar {}
303///
304/// let f: Arc<dyn Foo> = Arc::new(Bar);
305/// let b: Arc<Bar> = f.downcast().unwrap(); // `Arc::downcast` method not found
306/// ```
307///
308/// As would this:
309///
310/// ```fail_compile
311/// # use std::any::Any;
312/// # use std::sync::Arc;
313/// trait Foo: Any + Send + Sync {}
314///
315/// struct Bar;
316/// impl Foo for Bar {}
317///
318/// let f: Arc<dyn Foo> = Arc::new(Bar);
319/// let a: Arc<dyn Any + Send + Sync> = f; // trait upcasting coercion is not stable rust
320/// let f: Arc<Bar> = a.downcast().unwrap();
321/// ```
322///
323/// NOTE: `AsAny` inherits the `Send + Sync` constraint from [`Arc::downcast`].
324pub trait AsAny: Any + Send + Sync {
325 /// Obtains a `dyn Any` reference to the object:
326 ///
327 /// ```
328 /// # use buoyant_kernel as delta_kernel;
329 /// # use delta_kernel::AsAny;
330 /// # use std::any::Any;
331 /// # use std::sync::Arc;
332 /// trait Foo : AsAny {}
333 /// struct Bar;
334 /// impl Foo for Bar {}
335 ///
336 /// let f: &dyn Foo = &Bar;
337 /// let a: &dyn Any = f.any_ref();
338 /// let b: &Bar = a.downcast_ref().unwrap();
339 /// ```
340 fn any_ref(&self) -> &(dyn Any + Send + Sync);
341
342 /// Obtains an `Arc<dyn Any>` reference to the object:
343 ///
344 /// ```
345 /// # use buoyant_kernel as delta_kernel;
346 /// # use delta_kernel::AsAny;
347 /// # use std::any::Any;
348 /// # use std::sync::Arc;
349 /// trait Foo : AsAny {}
350 /// struct Bar;
351 /// impl Foo for Bar {}
352 ///
353 /// let f: Arc<dyn Foo> = Arc::new(Bar);
354 /// let a: Arc<dyn Any + Send + Sync> = f.as_any();
355 /// let b: Arc<Bar> = a.downcast().unwrap();
356 /// ```
357 fn as_any(self: Arc<Self>) -> Arc<dyn Any + Send + Sync>;
358
359 /// Converts the object to `Box<dyn Any>`:
360 ///
361 /// ```
362 /// # use buoyant_kernel as delta_kernel;
363 /// # use delta_kernel::AsAny;
364 /// # use std::any::Any;
365 /// # use std::sync::Arc;
366 /// trait Foo : AsAny {}
367 /// struct Bar;
368 /// impl Foo for Bar {}
369 ///
370 /// let f: Box<dyn Foo> = Box::new(Bar);
371 /// let a: Box<dyn Any> = f.into_any();
372 /// let b: Box<Bar> = a.downcast().unwrap();
373 /// ```
374 fn into_any(self: Box<Self>) -> Box<dyn Any + Send + Sync>;
375
376 /// Convenient wrapper for [`std::any::type_name`], since [`Any`] does not provide it and
377 /// [`Any::type_id`] is useless as a debugging aid (its `Debug` is just a mess of hex digits).
378 fn type_name(&self) -> &'static str;
379}
380
381// Blanket implementation for all eligible types
382impl<T: Any + Send + Sync> AsAny for T {
383 fn any_ref(&self) -> &(dyn Any + Send + Sync) {
384 self
385 }
386 fn as_any(self: Arc<Self>) -> Arc<dyn Any + Send + Sync> {
387 self
388 }
389 fn into_any(self: Box<Self>) -> Box<dyn Any + Send + Sync> {
390 self
391 }
392 fn type_name(&self) -> &'static str {
393 std::any::type_name::<Self>()
394 }
395}
396
397/// Extension trait that facilitates object-safe implementations of `PartialEq`.
398pub trait DynPartialEq: AsAny {
399 fn dyn_eq(&self, other: &dyn Any) -> bool;
400}
401
402// Blanket implementation for all eligible types
403impl<T: PartialEq + AsAny> DynPartialEq for T {
404 fn dyn_eq(&self, other: &dyn Any) -> bool {
405 other.downcast_ref::<T>().is_some_and(|other| self == other)
406 }
407}
408
409/// Trait for implementing an Expression evaluator.
410///
411/// It contains one Expression which can be evaluated on multiple ColumnarBatches.
412/// Connectors can implement this trait to optimize the evaluation using the
413/// connector specific capabilities.
414pub trait ExpressionEvaluator: AsAny {
415 /// Evaluate the expression on a given EngineData.
416 ///
417 /// Produces one value for each row of the input.
418 /// The data type of the output is same as the type output of the expression this evaluator is using.
419 fn evaluate(&self, batch: &dyn EngineData) -> DeltaResult<Box<dyn EngineData>>;
420}
421
422/// Trait for implementing a Predicate evaluator.
423///
424/// It contains one Predicate which can be evaluated on multiple ColumnarBatches.
425/// Connectors can implement this trait to optimize the evaluation using the
426/// connector specific capabilities.
427pub trait PredicateEvaluator: AsAny {
428 /// Evaluate the predicate on a given EngineData.
429 ///
430 /// Produces one boolean value for each row of the input.
431 fn evaluate(&self, batch: &dyn EngineData) -> DeltaResult<Box<dyn EngineData>>;
432}
433
434/// Provides expression evaluation capability to Delta Kernel.
435///
436/// Delta Kernel can use this handler to evaluate a predicate on partition filters,
437/// fill up partition column values, and any computation on data using Expressions.
438pub trait EvaluationHandler: AsAny {
439 /// Create an [`ExpressionEvaluator`] that can evaluate the given [`Expression`]
440 /// on columnar batches with the given [`Schema`] to produce data of [`DataType`].
441 ///
442 /// If the provided output type is a struct, its fields describe the columns of output produced
443 /// by the evaluator. Otherwise, the output schema is a single column named "output" of the
444 /// specified `output_type`. In all cases, the output schema is only used for its names (all
445 /// field names will be updated to match) and nullability (non-nullable columns can be converted
446 /// to nullable). Any mismatch in types (including number of columns) will produce an error.
447 ///
448 /// # Parameters
449 ///
450 /// - `input_schema`: Schema of the input data.
451 /// - `expression`: Expression to evaluate.
452 /// - `output_type`: Expected result data type.
453 ///
454 /// [`Schema`]: crate::schema::StructType
455 /// [`DataType`]: crate::schema::DataType
456 fn new_expression_evaluator(
457 &self,
458 input_schema: SchemaRef,
459 expression: ExpressionRef,
460 output_type: DataType,
461 ) -> DeltaResult<Arc<dyn ExpressionEvaluator>>;
462
463 /// Create a [`PredicateEvaluator`] that can evaluate the given [`Predicate`] on columnar
464 /// batches with the given [`Schema`] to produce a column of boolean results.
465 ///
466 /// The output schema is a single nullable boolean column named "output".
467 ///
468 /// # Parameters
469 ///
470 /// - `input_schema`: Schema of the input data.
471 /// - `predicate`: Predicate to evaluate.
472 ///
473 /// [`Schema`]: crate::schema::StructType
474 fn new_predicate_evaluator(
475 &self,
476 input_schema: SchemaRef,
477 predicate: PredicateRef,
478 ) -> DeltaResult<Arc<dyn PredicateEvaluator>>;
479
480 /// Create a single-row all-null-value [`EngineData`] with the schema specified by
481 /// `output_schema`.
482 // NOTE: we should probably allow DataType instead of SchemaRef, but can expand that in the
483 // future.
484 fn null_row(&self, output_schema: SchemaRef) -> DeltaResult<Box<dyn EngineData>>;
485
486 /// Create a multi-row [`EngineData`] by applying the given schema to multiple rows of values.
487 ///
488 /// Each element in `rows` represents one row of data, where each row is a slice of structured
489 /// scalar values (one scalar per top-level field in the schema).
490 ///
491 /// # Parameters
492 ///
493 /// - `schema`: Schema describing the structure of each row.
494 /// - `rows`: Slice of rows, where each row contains one structured scalar per top-level schema
495 /// field.
496 ///
497 /// # Returns
498 ///
499 /// A multi-row `EngineData` containing all rows.
500 ///
501 /// # Errors
502 ///
503 /// Returns an error if any row has a number of scalars that does not match the number of
504 /// top-level fields in `schema`, or if any scalar value cannot be appended to its corresponding
505 /// field's builder (e.g. due to a type mismatch).
506 ///
507 /// # Example
508 ///
509 /// For a schema with fields `[add: Struct, remove: Struct]`, each row should contain exactly 2
510 /// scalars: one for the `add` field and one for the `remove` field.
511 fn create_many(
512 &self,
513 schema: SchemaRef,
514 rows: &[&[Scalar]],
515 ) -> DeltaResult<Box<dyn EngineData>>;
516}
517
518/// Internal trait to allow us to have a private `create_one` API that's implemented for all
519/// EvaluationHandlers.
520// For some reason rustc doesn't detect it's usage so we allow(dead_code) here...
521#[allow(dead_code)]
522#[internal_api]
523trait EvaluationHandlerExtension: EvaluationHandler {
524 /// Create a single-row [`EngineData`] by applying the given schema to the leaf-values given in
525 /// `values`.
526 // Note: we will stick with a Schema instead of DataType (more constrained can expand in
527 // future)
528 fn create_one(&self, schema: SchemaRef, values: &[Scalar]) -> DeltaResult<Box<dyn EngineData>> {
529 // just get a single int column (arbitrary)
530 let null_row_schema = Arc::new(StructType::new_unchecked(vec![StructField::nullable(
531 "null_col",
532 DataType::INTEGER,
533 )]));
534 let null_row = self.null_row(null_row_schema.clone())?;
535
536 // Convert schema and leaf values to an expression
537 let row_expr = literal_expression_transform(schema.as_ref(), values)?;
538
539 let eval =
540 self.new_expression_evaluator(null_row_schema, row_expr.into(), schema.into())?;
541 eval.evaluate(null_row.as_ref())
542 }
543}
544
545// Auto-implement the extension trait for all EvaluationHandlers
546impl<T: EvaluationHandler + ?Sized> EvaluationHandlerExtension for T {}
547
548/// A trait that allows converting a type into (single-row) EngineData
549///
550/// This is typically used with the `#[derive(IntoEngineData)]` macro
551/// which leverages the traits `ToDataType` and `Into<Scalar>` for struct fields
552/// to convert a struct into EngineData.
553///
554/// # Example
555/// ```ignore
556/// # use std::sync::Arc;
557/// # use delta_kernel_derive::{Schema, IntoEngineData};
558///
559/// #[derive(Schema, IntoEngineData)]
560/// struct MyStruct {
561/// a: i32,
562/// b: String,
563/// }
564///
565/// let my_struct = MyStruct { a: 42, b: "Hello".to_string() };
566/// // typically used with ToSchema
567/// let schema = Arc::new(MyStruct::to_schema());
568/// // single-row EngineData
569/// let engine = todo!(); // create an engine
570/// let engine_data = my_struct.into_engine_data(schema, engine);
571/// ```
572#[internal_api]
573pub(crate) trait IntoEngineData {
574 /// Consume this type to produce a single-row EngineData using the provided schema.
575 fn into_engine_data(
576 self,
577 schema: SchemaRef,
578 engine: &dyn Engine,
579 ) -> DeltaResult<Box<dyn EngineData>>;
580}
581
582/// Provides file system related functionalities to Delta Kernel.
583///
584/// Delta Kernel uses this handler whenever it needs to access the underlying
585/// file system where the Delta table is present. Connector implementation of
586/// this trait can hide filesystem specific details from Delta Kernel.
587pub trait StorageHandler: AsAny {
588 /// List the paths in the same directory that are lexicographically greater than
589 /// (UTF-8 sorting) the given `path`. The result should also be sorted by the file name.
590 ///
591 /// If the path is directory-like (ends with '/'), the result should contain
592 /// all the files in the directory.
593 fn list_from(&self, path: &Url)
594 -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<FileMeta>>>>;
595
596 /// Read data specified by the start and end offset from the file.
597 fn read_files(
598 &self,
599 files: Vec<FileSlice>,
600 ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<Bytes>>>>;
601
602 /// Copy a file atomically from source to destination. If the destination file already exists,
603 /// it must return Err(Error::FileAlreadyExists).
604 fn copy_atomic(&self, src: &Url, dest: &Url) -> DeltaResult<()>;
605
606 /// Write data to the specified path.
607 ///
608 /// If `overwrite` is false and the file already exists, this must return
609 /// `Err(Error::FileAlreadyExists)`.
610 fn put(&self, path: &Url, data: Bytes, overwrite: bool) -> DeltaResult<()>;
611
612 /// Perform a HEAD request for the given file at a Url, returning the file metadata.
613 ///
614 /// If the file does not exist, this must return an `Err` with [`Error::FileNotFound`].
615 fn head(&self, path: &Url) -> DeltaResult<FileMeta>;
616}
617
618/// Provides JSON handling functionality to Delta Kernel.
619///
620/// Delta Kernel can use this handler to parse JSON strings into Row or read content from JSON files.
621/// Connectors can leverage this trait to provide their best implementation of the JSON parsing
622/// capability to Delta Kernel.
623pub trait JsonHandler: AsAny {
624 /// Parse the given json strings and return the fields requested by output schema as columns in [`EngineData`].
625 /// json_strings MUST be a single column batch of engine data, and the column type must be string
626 fn parse_json(
627 &self,
628 json_strings: Box<dyn EngineData>,
629 output_schema: SchemaRef,
630 ) -> DeltaResult<Box<dyn EngineData>>;
631
632 /// Read and parse the JSON format file at given locations and return the data as EngineData with
633 /// the columns requested by physical schema. Note: The [`FileDataReadResultIterator`] must emit
634 /// data from files in the order that `files` is given. For example if files ["a", "b"] is provided,
635 /// then the engine data iterator must first return all the engine data from file "a", _then_ all
636 /// the engine data from file "b". Moreover, for a given file, all of its [`EngineData`] and
637 /// constituent rows must be in order that they occur in the file. Consider a file with rows
638 /// (1, 2, 3). The following are legal iterator batches:
639 /// iter: [EngineData(1, 2), EngineData(3)]
640 /// iter: [EngineData(1), EngineData(2, 3)]
641 /// iter: [EngineData(1, 2, 3)]
642 /// The following are illegal batches:
643 /// iter: [EngineData(3), EngineData(1, 2)]
644 /// iter: [EngineData(1), EngineData(3, 2)]
645 /// iter: [EngineData(2, 1, 3)]
646 ///
647 /// Additionally, engines may not merge engine data across file boundaries.
648 ///
649 /// # Parameters
650 ///
651 /// - `files` - File metadata for files to be read.
652 /// - `physical_schema` - Select list of columns to read from the JSON file.
653 /// - `predicate` - Optional push-down predicate hint (engine is free to ignore it).
654 fn read_json_files(
655 &self,
656 files: &[FileMeta],
657 physical_schema: SchemaRef,
658 predicate: Option<PredicateRef>,
659 ) -> DeltaResult<FileDataReadResultIterator>;
660
661 /// Atomically (!) write a single JSON file. Each row of the input data should be written as a
662 /// new JSON object appended to the file. this write must:
663 /// (1) serialize the data to newline-delimited json (each row is a json object literal)
664 /// (2) write the data to storage atomically (i.e. if the file already exists, fail unless the
665 /// overwrite flag is set)
666 ///
667 /// For example, the JSON data should be written as { "column1": "val1", "column2": "val2", .. }
668 /// with each row on a new line.
669 ///
670 /// NOTE: Null columns should not be written to the JSON file. For example, if a row has columns
671 /// ["a", "b"] and the value of "b" is null, the JSON object should be written as
672 /// { "a": "..." }. Note that including nulls is technically valid JSON, but would bloat the
673 /// log, therefore we recommend omitting them.
674 ///
675 /// # Parameters
676 ///
677 /// - `path` - URL specifying the location to write the JSON file
678 /// - `data` - Iterator of EngineData to write to the JSON file. Each row should be written as
679 /// a new JSON object appended to the file. (that is, the file is newline-delimited JSON, and
680 /// each row is a JSON object on a single line)
681 /// - `overwrite` - If true, overwrite the file if it exists. If false, the call must fail if
682 /// the file exists.
683 fn write_json_file(
684 &self,
685 path: &Url,
686 data: Box<dyn Iterator<Item = DeltaResult<FilteredEngineData>> + Send + '_>,
687 overwrite: bool,
688 ) -> DeltaResult<()>;
689}
690
691/// Reserved field IDs for metadata columns in Delta tables.
692///
693/// These field IDs are reserved and should not be used for regular table columns.
694/// They are used to provide file-level metadata as virtual columns during reads.
695pub mod reserved_field_ids {
696 /// Reserved field ID for the file name metadata column (`_file`).
697 /// This column provides the name of the Parquet file that contains each row.
698 pub const FILE_NAME: i64 = 2147483646;
699}
700
701/// Metadata from a Parquet file footer.
702///
703/// This struct contains metadata extracted from a Parquet file's footer, including the schema.
704/// It is designed to be extensible for future additions such as row group statistics.
705#[derive(Debug, Clone)]
706pub struct ParquetFooter {
707 /// The schema of the Parquet file, converted to Delta Kernel's schema format.
708 pub schema: SchemaRef,
709}
710
711/// Provides Parquet file related functionalities to Delta Kernel.
712///
713/// Connectors can leverage this trait to provide their own custom
714/// implementation of Parquet data file functionalities to Delta Kernel.
715pub trait ParquetHandler: AsAny {
716 /// Read and parse the Parquet file at given locations and return the data as EngineData with
717 /// the columns requested by physical schema. The ParquetHandler _must_ return exactly the
718 /// columns specified in `physical_schema`, and they _must_ be in schema order.
719 ///
720 /// # Resolving Parquet schema to the physical schema
721 ///
722 /// When reading the Parquet file, the columns are resolved from the Parquet schema to the
723 /// kernel's `physical_schema`. To do so, the parquet reader must match each Parquet column
724 /// to a [`StructField`] in the `physical_schema`. All columns in the returned `EngineData`
725 /// must be in the same order as specified in `physical_schema`.
726 ///
727 /// Parquet columns are matched to `physical_schema` [`StructField`]s using the following rules:
728 /// 1. **Field ID**: If a [`StructField`] in `physical_schema` contains a field ID
729 /// (specified in [`ColumnMetadataKey::ParquetFieldId`] metadata), use the ID to
730 /// match the Parquet column's field id
731 /// 2. **Field Name**: If no field ID is present in the `physical_schema`'s [`StructField`] or no matching parquet field ID is found,
732 /// fall back to matching by column name
733 ///
734 /// # Metadata Columns
735 ///
736 /// The ParquetHandler must support virtual metadata columns that provide additional information
737 /// about each row. These columns are not stored in the Parquet file but are generated at read time.
738 ///
739 /// ## Row Index Column
740 ///
741 /// When a column in `physical_schema` is marked as a row index metadata column (via
742 /// [`StructField::create_metadata_column`] with [`schema::MetadataColumnSpec::RowIndex`]), the
743 /// ParquetHandler must populate it with the 0-based row position within the Parquet file:
744 ///
745 /// - **Column name**: User-specified (commonly `"row_index"` or `"_metadata.row_index"`)
746 /// - **Type**: `LONG` (non-nullable)
747 /// - **Values**: Sequential integers starting at 0 for each file
748 /// - **Use case**: Track row positions for downstream processing, or internally used to compute Row IDs
749 ///
750 /// Example: A file with 5 rows would have row_index values `[0, 1, 2, 3, 4]`.
751 ///
752 /// ## File Name Column (Reserved Field ID)
753 ///
754 /// When a column in `physical_schema` has the reserved field ID
755 /// [`reserved_field_ids::FILE_NAME`] (2147483646), the ParquetHandler must populate it
756 /// with the file path/name:
757 ///
758 /// - **Column name**: `"_file"`
759 /// - **Type**: `STRING` (non-nullable)
760 /// - **Field ID**: 2147483646 (reserved)
761 /// - **Values**: The file path/URL (e.g., `"s3://bucket/path/file.parquet"`)
762 /// - **Use case**: Track which file each row came from in multi-file reads
763 ///
764 /// Example: All rows from the same file would have the same `_file` value.
765 ///
766 /// ## Metadata Column Examples
767 ///
768 /// ```rust,ignore
769 /// use delta_kernel::schema::{StructType, StructField, DataType, MetadataColumnSpec};
770 ///
771 /// // Example 1: Schema with row_index metadata column
772 /// let schema_with_row_index = StructType::try_new([
773 /// StructField::nullable("id", DataType::INTEGER),
774 /// StructField::create_metadata_column("row_index", MetadataColumnSpec::RowIndex),
775 /// StructField::nullable("value", DataType::STRING),
776 /// ])?;
777 ///
778 /// // Example 2: Schema with _file metadata column (using reserved field ID)
779 /// let schema_with_file_path = StructType::try_new([
780 /// StructField::nullable("id", DataType::INTEGER),
781 /// StructField::create_metadata_column("_file", MetadataColumnSpec::FilePath),
782 /// StructField::nullable("value", DataType::STRING),
783 /// ])?;
784 /// ```
785 ///
786 /// ---
787 ///
788 /// If no matching Parquet column is found, `NULL` values are returned
789 /// for nullable columns in `physical_schema`. For non-nullable columns, an error is returned.
790 ///
791 ///
792 /// ## Column Matching Examples
793 ///
794 /// Consider a `physical_schema` with the following fields:
795 /// - Column 0: `"i_logical"` (integer, non-null) with field ID 1 (via [`ColumnMetadataKey::ParquetFieldId`])
796 /// - Column 1: `"s"` (string, nullable) with no field ID metadata
797 /// - Column 2: `"i2"` (integer, nullable) with no field ID metadata
798 ///
799 /// [`ColumnMetadataKey::ParquetFieldId`]: crate::schema::ColumnMetadataKey::ParquetFieldId
800 ///
801 /// And a Parquet file containing these columns:
802 /// - Column 0: `"i2"` (integer, nullable) with field ID 3
803 /// - Column 1: `"i"` (integer, non-null) with field ID 1
804 /// - No `"s"` column present
805 ///
806 /// The column matching would work as follows:
807 /// - `"i_logical"` matches `"i"` by field ID (both have ID 1)
808 /// - `"i2"` matches `"i2"` by column name (no field ID to match on)
809 /// - `"s"` has no matching Parquet column, so NULL values are returned
810 ///
811 /// The returned data will contain exactly 3 columns in physical schema order:
812 /// `{i_logical: parquet[1], s: NULL.., i2: parquet[0]}`
813 ///
814 /// # Parameters
815 ///
816 /// - `files` - File metadata for files to be read.
817 /// - `physical_schema` - Select list and order of columns to read from the Parquet file.
818 /// - `predicate` - Optional push-down predicate hint (engine is free to ignore it).
819 ///
820 /// # Returns
821 /// A [`DeltaResult`] containing a [`FileDataReadResultIterator`].
822 /// Each element of the iterator is a [`DeltaResult`] of [`EngineData`]. The [`EngineData`]
823 /// has the contents of `files` and must match the provided `physical_schema`.
824 ///
825 /// Note: The [`FileDataReadResultIterator`] must emit data from files in the order that `files`
826 /// is given. For example if files ["a", "b"] is provided, then the engine data iterator must
827 /// first return all the engine data from file "a", _then_ all the engine data from file "b".
828 /// Moreover, for a given file, all of its [`EngineData`] and constituent rows must be in order
829 /// that they occur in the file. Consider a file with rows
830 /// (1, 2, 3). The following are legal iterator batches:
831 /// iter: [EngineData(1, 2), EngineData(3)]
832 /// iter: [EngineData(1), EngineData(2, 3)]
833 /// iter: [EngineData(1, 2, 3)]
834 /// The following are illegal batches:
835 /// iter: [EngineData(3), EngineData(1, 2)]
836 /// iter: [EngineData(1), EngineData(3, 2)]
837 /// iter: [EngineData(2, 1, 3)]
838 ///
839 /// Additionally, engines must not merge engine data across file boundaries.
840 ///
841 /// [`ColumnMetadataKey::ParquetFieldId`]: crate::schema::ColumnMetadataKey
842 fn read_parquet_files(
843 &self,
844 files: &[FileMeta],
845 physical_schema: SchemaRef,
846 predicate: Option<PredicateRef>,
847 ) -> DeltaResult<FileDataReadResultIterator>;
848
849 /// Write data to a Parquet file at the specified URL.
850 ///
851 /// This method writes the provided `data` to a Parquet file at the given `url`.
852 ///
853 /// This will overwrite the file if it already exists. For filesystem-backed
854 /// implementations, the parent directories must be created if they do not exist.
855 ///
856 /// # Parameters
857 ///
858 /// - `url` - The full URL path where the Parquet file should be written
859 /// (e.g., `s3://bucket/path/file.parquet`).
860 /// - `data` - An iterator of engine data to be written to the Parquet file.
861 ///
862 /// # Returns
863 ///
864 /// A [`DeltaResult`] indicating success or failure.
865 fn write_parquet_file(
866 &self,
867 location: url::Url,
868 data: Box<dyn Iterator<Item = DeltaResult<Box<dyn EngineData>>> + Send>,
869 ) -> DeltaResult<()>;
870
871 /// Read the footer metadata from a Parquet file without reading the data.
872 ///
873 /// This method reads only the Parquet file footer (metadata section), which is useful for
874 /// schema inspection, compatibility checking, and determining whether parsed statistics
875 /// columns are present and compatible with the current table schema.
876 ///
877 /// # Parameters
878 ///
879 /// - `file` - File metadata for the Parquet file whose footer should be read. The `size` field
880 /// should contain the actual file size to enable efficient footer reads without additional
881 /// I/O operations.
882 ///
883 /// # Returns
884 ///
885 /// A [`DeltaResult`] containing a [`ParquetFooter`] with the Parquet file's metadata, including
886 /// the schema converted to Delta Kernel's format.
887 ///
888 /// # Field IDs
889 ///
890 /// If the Parquet file contains field IDs (written when column mapping is enabled), they are
891 /// preserved in each [`StructField`]'s metadata. Callers can access field IDs via
892 /// [`StructField::get_config_value`] with [`ColumnMetadataKey::ParquetFieldId`].
893 ///
894 /// # Errors
895 ///
896 /// Returns an error if:
897 /// - The file cannot be accessed or does not exist
898 /// - The file is not a valid Parquet file
899 /// - The footer cannot be read or parsed
900 /// - The schema cannot be converted to Delta Kernel's format
901 ///
902 /// [`StructField`]: crate::schema::StructField
903 /// [`StructField::get_config_value`]: crate::schema::StructField::get_config_value
904 /// [`ColumnMetadataKey::ParquetFieldId`]: crate::schema::ColumnMetadataKey::ParquetFieldId
905 fn read_parquet_footer(&self, file: &FileMeta) -> DeltaResult<ParquetFooter>;
906}
907
908/// The `Engine` trait encapsulates all the functionality an engine or connector needs to provide
909/// to the Delta Kernel in order to read the Delta table.
910///
911/// Engines/Connectors are expected to pass an implementation of this trait when reading a Delta
912/// table.
913pub trait Engine: AsAny {
914 /// Get the connector provided [`EvaluationHandler`].
915 fn evaluation_handler(&self) -> Arc<dyn EvaluationHandler>;
916
917 /// Get the connector provided [`StorageHandler`]
918 fn storage_handler(&self) -> Arc<dyn StorageHandler>;
919
920 /// Get the connector provided [`JsonHandler`].
921 fn json_handler(&self) -> Arc<dyn JsonHandler>;
922
923 /// Get the connector provided [`ParquetHandler`].
924 fn parquet_handler(&self) -> Arc<dyn ParquetHandler>;
925
926 /// Get the connector provided [`MetricsReporter`] for metrics collection.
927 ///
928 /// Returns an optional reporter that will receive metric events from Delta operations.
929 /// The default implementation returns None (no metrics reporting).
930 fn get_metrics_reporter(&self) -> Option<Arc<dyn MetricsReporter>> {
931 None
932 }
933}
934
935// we have an 'internal' feature flag: default-engine-base, which is actually just the shared
936// pieces of default-engine-native-tls and default-engine-rustls. the crate can't compile with _only_
937// default-engine-base, so we give a friendly error here.
938#[cfg(all(
939 feature = "default-engine-base",
940 not(any(
941 feature = "default-engine-native-tls",
942 feature = "default-engine-rustls",
943 ))
944))]
945compile_error!(
946 "The default-engine-base feature flag is not meant to be used directly. \
947 Please use either default-engine-native-tls or default-engine-rustls."
948);
949
950// Rustdoc's documentation tests can do some things that regular unit tests can't. Here we are
951// using doctests to test macros. Specifically, we are testing for failed macro invocations due
952// to invalid input, not the macro output when the macro invocation is successful (which can/should be
953// done in unit tests). This module is not exclusively for macro tests only so other doctests can also be added.
954// https://doc.rust-lang.org/rustdoc/write-documentation/documentation-tests.html#include-items-only-when-collecting-doctests
955#[cfg(doctest)]
956mod doctests;