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 ofdelete,insert,update_preimage, orupdate_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
inCommitTimestampfield 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_newrequires Change Data Feed to remain enabled and requires exact schema equality.TableChanges::try_new_row_tracking_cdf_listingrequires 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
impl TableChanges
Sourcepub fn try_new(
table_root: Url,
engine: &dyn Engine,
start_version: Version,
end_version: Option<Version>,
) -> DeltaResult<Self>
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_logfolder is located)engine: Implementation ofEngineapis.start_version: The start version of the change data feedend_version: The end version (inclusive) of the change data feed. If this is none, this defaults to the newest table version.
Sourcepub 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.
pub fn try_new_row_tracking_cdf_listing( table_root: Url, engine: &dyn Engine, start_version: Version, end_version: Option<Version>, ) -> DeltaResult<Self>
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.
Sourcepub fn start_version(&self) -> Version
pub fn start_version(&self) -> Version
The start version of the TableChanges.
Sourcepub fn end_version(&self) -> Version
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.
Sourcepub fn schema(&self) -> &Schema
pub fn schema(&self) -> &Schema
The logical schema of the change data feed. For details on the shape of the schema, see
TableChanges.
Sourcepub fn table_root(&self) -> &Url
pub fn table_root(&self) -> &Url
Path to the root of the table that is being read.
Sourcepub fn materialized_row_id_column_name(&self) -> DeltaResult<&str>
Available on crate feature internal-api only.
pub fn materialized_row_id_column_name(&self) -> DeltaResult<&str>
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.
Sourcepub fn materialized_row_commit_version_column_name(&self) -> DeltaResult<&str>
Available on crate feature internal-api only.
pub fn materialized_row_commit_version_column_name(&self) -> DeltaResult<&str>
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.
Sourcepub fn scan_builder(self: Arc<Self>) -> TableChangesScanBuilder
pub fn scan_builder(self: Arc<Self>) -> TableChangesScanBuilder
Create a TableChangesScanBuilder for an Arc<TableChanges>.
Sourcepub fn into_scan_builder(self) -> TableChangesScanBuilder
pub fn into_scan_builder(self) -> TableChangesScanBuilder
Consume this TableChanges to create a TableChangesScanBuilder
Sourcepub 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.
pub fn scan_file_listing( self: Arc<Self>, engine: Arc<dyn Engine>, mode: TableChangesListingMode, ) -> DeltaResult<impl Iterator<Item = DeltaResult<TableChangesFileAction>>>
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§
Auto Trait Implementations§
impl !RefUnwindSafe for TableChanges
impl !UnwindSafe for TableChanges
impl Freeze for TableChanges
impl Send for TableChanges
impl Sync for TableChanges
impl Unpin for TableChanges
impl UnsafeUnpin for TableChanges
Blanket Implementations§
Source§impl<T> AsAny for T
impl<T> AsAny for T
Source§fn any_ref(&self) -> &(dyn Any + Send + Sync + 'static)
fn any_ref(&self) -> &(dyn Any + Send + Sync + 'static)
dyn Any reference to the object: Read moreSource§fn as_any(self: Arc<T>) -> Arc<dyn Any + Send + Sync> ⓘ
fn as_any(self: Arc<T>) -> Arc<dyn Any + Send + Sync> ⓘ
Arc<dyn Any> reference to the object: Read moreSource§fn into_any(self: Box<T>) -> Box<dyn Any + Send + Sync>
fn into_any(self: Box<T>) -> Box<dyn Any + Send + Sync>
Box<dyn Any>: Read moreSource§fn type_name(&self) -> &'static str
fn type_name(&self) -> &'static str
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> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> FoldWithOption for T
impl<T> FoldWithOption for T
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
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 moreSource§impl<T> PolicyExt for Twhere
T: ?Sized,
impl<T> PolicyExt for Twhere
T: ?Sized,
Source§impl<KernelType, ArrowType> TryIntoArrow<ArrowType> for KernelTypewhere
ArrowType: TryFromKernel<KernelType>,
impl<KernelType, ArrowType> TryIntoArrow<ArrowType> for KernelTypewhere
ArrowType: TryFromKernel<KernelType>,
Source§fn try_into_arrow(self) -> Result<ArrowType, ArrowError>
fn try_into_arrow(self) -> Result<ArrowType, ArrowError>
arrow-conversion and (crate features arrow-conversion or declarative-plans or default-engine-base) only.Source§impl<KernelType, ArrowType> TryIntoKernel<KernelType> for ArrowTypewhere
KernelType: TryFromArrow<ArrowType>,
impl<KernelType, ArrowType> TryIntoKernel<KernelType> for ArrowTypewhere
KernelType: TryFromArrow<ArrowType>,
Source§fn try_into_kernel(self) -> Result<KernelType, ArrowError>
fn try_into_kernel(self) -> Result<KernelType, ArrowError>
arrow-conversion and (crate features arrow-conversion or declarative-plans or default-engine-base) only.