Skip to main content

TableChanges

Struct TableChanges 

Source
pub struct TableChanges { /* private fields */ }
Expand description

Represents a call to read the Change Data Feed (CDF) between two versions of a table. The schema of TableChanges will be the schema of the table at the end version with three additional columns:

  • _change_type: String representing the type of change that for that commit. This may be one of delete, insert, update_preimage, or update_postimage.

  • _commit_version: Long representing the commit the change occurred in.

  • _commit_timestamp: Time at which the commit occurred. The timestamp is retrieved from the file modification time of the log file. No timezone is associated with the timestamp.

    Currently, in-commit timestamps (ICT) is not supported. In the future when ICT is enabled, the timestamp will be retrieved from the inCommitTimestamp field of the CommitInfo` action. See issue #559 For details on In-Commit Timestamps, see the Protocol.

Three properties must hold for the entire CDF range:

  • Reading must be supported for every commit in the range: every enabled reader feature must be supported by the kernel. The supported read features will be expanded in the future to cover more delta table features.
  • TableChanges::try_new requires Change Data Feed to remain enabled and requires exact schema equality.
  • TableChanges::try_new_row_tracking_cdf_listing requires row tracking to remain enabled. It allows additive nullable columns and relaxed nullability, but rejects datatype changes.

Construction validates the range boundaries. Intermediate metadata and protocol updates are validated when the transaction log is replayed by the scan or listing operation.

§Examples

Get TableChanges for versions 0 to 1 (inclusive)

let url = delta_kernel::try_parse_uri(path)?;
let table_changes = TableChanges::try_new(url, &engine, 0, Some(1))?;

For more details, see the following sections of the protocol:

Implementations§

Source§

impl TableChanges

Source

pub fn try_new( table_root: Url, engine: &dyn Engine, start_version: Version, end_version: Option<Version>, ) -> DeltaResult<Self>

Creates a new TableChanges instance for the given version range. This function checks these properties:

  • The change data feed table feature must be enabled in both the start or end versions.
  • Every enabled reader feature must be supported by the kernel.
  • The schemas at the start and end versions are the same.

Note that this does not check that change data feed is enabled for every commit in the range. It also does not check that the schema remains the same for the entire range.

§Parameters
  • table_root: url pointing at the table root (where _delta_log folder is located)
  • engine: Implementation of Engine apis.
  • start_version: The start version of the change data feed
  • end_version: The end version (inclusive) of the change data feed. If this is none, this defaults to the newest table version.
Source

pub fn try_new_row_tracking_cdf_listing( table_root: Url, engine: &dyn Engine, start_version: Version, end_version: Option<Version>, ) -> DeltaResult<Self>

Available on crate feature internal-api only.

Creates a listing-only change feed from row-tracking metadata.

This path requires delta.enableRowTracking and ignores _change_data files. TableChanges::scan_file_listing returns the add and remove actions whose data must be reconciled by row ID.

Construction validates the range boundaries. TableChanges::scan_file_listing validates intermediate metadata and protocol updates while replaying the range. Every enabled reader feature must be supported by Kernel, and each schema must be readable using the end-version logical schema without datatype widening.

§Parameters
  • table_root: URL of the table root containing _delta_log.
  • engine: Engine used to read the transaction log.
  • start_version: First version in the change feed.
  • end_version: The end version (inclusive) of the change data feed. If this is none, this defaults to the newest table version.
§Errors

Returns an error if the range cannot be loaded or a boundary has unavailable row tracking, unsupported reader features, or an incompatible schema. Errors from intermediate versions are returned by TableChanges::scan_file_listing.

Source

pub fn start_version(&self) -> Version

The start version of the TableChanges.

Source

pub fn end_version(&self) -> Version

The end version (inclusive) of the TableChanges. If no end_version was specified in TableChanges::try_new, this returns the newest version as of the call to try_new.

Source

pub fn schema(&self) -> &Schema

The logical schema of the change data feed. For details on the shape of the schema, see TableChanges.

Source

pub fn table_root(&self) -> &Url

Path to the root of the table that is being read.

Source

pub fn materialized_row_id_column_name(&self) -> DeltaResult<&str>

Available on crate feature internal-api only.

Returns the physical Parquet column that stores materialized row IDs.

§Errors

Returns an error unless this value was created by TableChanges::try_new_row_tracking_cdf_listing.

Source

pub fn materialized_row_commit_version_column_name(&self) -> DeltaResult<&str>

Available on crate feature internal-api only.

Returns the physical Parquet column that stores materialized row commit versions.

§Errors

Returns an error unless this value was created by TableChanges::try_new_row_tracking_cdf_listing.

Source

pub fn scan_builder(self: Arc<Self>) -> TableChangesScanBuilder

Create a TableChangesScanBuilder for an Arc<TableChanges>.

Source

pub fn into_scan_builder(self) -> TableChangesScanBuilder

Consume this TableChanges to create a TableChangesScanBuilder

Source

pub fn scan_file_listing( self: Arc<Self>, engine: Arc<dyn Engine>, mode: TableChangesListingMode, ) -> DeltaResult<impl Iterator<Item = DeltaResult<TableChangesFileAction>>>

Available on crate feature internal-api only.

Lists the files required to reconstruct a row-tracking change feed without reading data files.

Each TableChangesFileAction contains an add side, a remove side, or both. Read each side using its own deletion vector. For every selected physical row, reconstruct the row ID as coalesce(materialized_row_id, base_row_id + physical_row_index). Reconstruct the row commit version as coalesce(materialized_row_commit_version, default_row_commit_version). The materialized column names are returned by TableChanges::materialized_row_id_column_name and TableChanges::materialized_row_commit_version_column_name. Assign physical row indexes before applying deletion vectors.

Rows present only on the add side are inserts, rows present only on the remove side are deletes, and matching row IDs with different row commit versions are updates. Matching both the row ID and row commit version identifies a row carried forward without change. AddCDCFile actions are ignored.

Both listing modes buffer the selected actions before returning. The iterator therefore reports preparation errors immediately and uses memory proportional to the range’s action count. This method requires a TableChanges created by TableChanges::try_new_row_tracking_cdf_listing.

§Errors

Returns an error if this value was not created by TableChanges::try_new_row_tracking_cdf_listing or if the log actions cannot be replayed into a valid listing.

Trait Implementations§

Source§

impl Debug for TableChanges

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> AsAny for T
where T: Any + Send + Sync,

Source§

fn any_ref(&self) -> &(dyn Any + Send + Sync + 'static)

Obtains a dyn Any reference to the object: Read more
Source§

fn as_any(self: Arc<T>) -> Arc<dyn Any + Send + Sync>

Obtains an Arc<dyn Any> reference to the object: Read more
Source§

fn into_any(self: Box<T>) -> Box<dyn Any + Send + Sync>

Converts the object to Box<dyn Any>: Read more
Source§

fn type_name(&self) -> &'static str

Convenient wrapper for std::any::type_name, since Any does not provide it and Any::type_id is useless as a debugging aid (its Debug is just a mess of hex digits).
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> FoldWithOption for T

Source§

fn fold_with<U>(self, opt: Option<U>, f: impl FnOnce(Self, U) -> Self) -> Self

Available on crate feature internal-api only.
Applies an optional fold operation f to self if opt is Some; otherwise returns self unchanged. Read more
Source§

fn try_fold_with<U, E>( self, opt: Option<U>, f: impl FnOnce(Self, U) -> Result<Self, E>, ) -> Result<Self, E>

Available on crate feature internal-api only.
Fallible fold_with: applies Result-returning f to self if opt is Some, otherwise returns self unchanged (wrapped in Ok).
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<KernelType, ArrowType> TryIntoArrow<ArrowType> for KernelType
where ArrowType: TryFromKernel<KernelType>,

Source§

fn try_into_arrow(self) -> Result<ArrowType, ArrowError>

Available on crate feature arrow-conversion and (crate features arrow-conversion or declarative-plans or default-engine-base) only.
Source§

impl<KernelType, ArrowType> TryIntoKernel<KernelType> for ArrowType
where KernelType: TryFromArrow<ArrowType>,

Source§

fn try_into_kernel(self) -> Result<KernelType, ArrowError>

Available on crate feature arrow-conversion and (crate features arrow-conversion or declarative-plans or default-engine-base) only.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more