Repository navigation
Conversation
## 🥞 Stacked PR Use this [link](https://github.com/delta-io/delta/pull/7246/files) to review incremental changes. - [stack/tw/flink-uc-delta-tables-api](delta-io#7229) [[Files changed](https://github.com/delta-io/delta/pull/7229/files)] [MERGED] - [stack/flink-uc-2-lookup-creds](delta-io#7244) [[Files changed](https://github.com/delta-io/delta/pull/7244/files)] [MERGED] - [stack/flink-uc-3-create](delta-io#7245) [[Files changed](https://github.com/delta-io/delta/pull/7245/files)] [MERGED] - [**stack/flink-uc-4-metrics**](delta-io#7246) [[Files changed](https://github.com/delta-io/delta/pull/7246/files)] ← _this PR_ - [stack/flink-uc-5-browsing](delta-io#7247) [[Files changed](https://github.com/delta-io/delta/pull/7247/files/0f5f1328fc39d1fe389d85eb480cc0742cdf67d0..f6b6d957cddc0ee36cfa48324b40b5d2d5298149)] - [stack/flink-uc-6-catalog-correctness](delta-io#7349) [[Files changed](https://github.com/delta-io/delta/pull/7349/files/f6b6d957cddc0ee36cfa48324b40b5d2d5298149..90df3fed25d136fe869f1a1034a5472997604b00)] --------- #### Which Delta project/connector is this regarding? - [ ] Spark - [ ] Standalone - [x] Flink - [ ] Kernel - [ ] Other (fill in here) ## What this PR enables Enables **post-commit metric reporting** over the Delta-Tables API, replacing the hand-written metrics REST client. ```text ENDPOINT OPERATION STATUS ---------------------------------------------------------------------- -------------------------- ------ GET /v1/catalogs/{catalog}/schemas/{schema}/tables/{table} Load table / table exists ✅ done POST /v1/catalogs/{catalog}/schemas/{schema}/tables/{table} Commit (updateTable) ✅ done GET /v1/catalogs/{catalog}/schemas/{schema}/tables/{table}/credentials Vend table storage creds ✅ done POST /v1/catalogs/{catalog}/schemas/{schema}/staging-tables Reserve staging table ✅ done POST /v1/catalogs/{catalog}/schemas/{schema}/tables Finalize (create) table ✅ done GET /v1/staging-tables/{table_id}/credentials Vend staging storage creds ✅ done POST /v1/catalogs/{catalog}/schemas/{schema}/tables/{table}/metrics Report commit metrics ⏳ THIS PR GET /v1/catalogs/{catalog}/schemas/{schema} List tables in schema ⛔ blocked ``` _Status legend: ✅ done (earlier PR in this stack) · ⏳ THIS PR · ⛔ blocked (no UC Delta client binding; stays on the OpenAPI client)._ ## Description This is the fourth PR in the stack. After each commit, Flink reports added and removed file counts, row and byte counts, and a file-size histogram to UC. Previously `ReportCommitMetricsListener` used a hand-written `DeltaMetricsApi` and a separate set of request and response classes. This PR converts the Kernel `TransactionReport` into `UCDeltaModels.CommitReport`, calls `UCDeltaClient.reportMetrics`, and removes the duplicated wire-format implementation. Metric reporting remains best-effort: a reporting failure is logged and does not fail the table commit. ## How was this patch tested? `build/sbt "flink / test"` passes: 245 total, 241 passed, 4 skipped. ## Does this PR introduce _any_ user-facing changes? No. Signed-off-by: Timothy Wang <timothy.art@gmail.com>
#### Which Delta project/connector is this regarding? - [ ] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [X] Other (Protocol) ## Description Update the interval types RFC to more accurately align with kernel reference implementation.
…a tree (delta-io#7366) ## Description A full checkpoint rewrite clusters the live files into one manifest per Spark partition and then writes a root listing one pointer per manifest. When a snapshot is small enough that the distribution yields a single manifest, that extra root adds a level of indirection describing nothing: readers follow a pointer to reach the only manifest there is. This promotes the lone manifest to the root in place instead. The `ContentRoot` points straight at what the executor already wrote, so no rename or extra I/O is needed, and no leaf pointers are returned. Trees with more than one manifest are unaffected. The promoted file keeps the name the executor gave it, which readers must already tolerate since manifest names are informational. ## How was this patch tested? Existing checkpoint suites, plus: - A test that a small full checkpoint promotes its single manifest to the root, and that the promoted root alone reconstructs exactly the live file set. - A test that a tree with no leaves stamps no back references, covering every checkpoint scenario: a full rewrite promotes its lone manifest, and an incremental keeps small entries root-resident. Reaching a leaf-resident entry needs enough files that the writer really packs leaves, so the per-leaf cap, the file count, and the expected leaf count now come from shared helpers on the test base rather than being restated in each suite. ## Does this PR introduce _any_ user-facing changes? No.
…mmit (delta-io#7315) ## Description This is the write-side counterpart to reading `LastManifestCommit` (LMC) from a `Snapshot`: it records the last manifest commit reference into both `CommitInfo` and the `VersionChecksum` (CRC) as commits are made, so that later reads can resolve it without extra I/O. Concretely: - `CommitInfo` gains a `lastManifestCommit: Option[LastManifestCommit]` field, and `VersionChecksum` carries it forward as well. Serialization/deserialization is updated accordingly (see `actions.scala`, `DeltaEncoders`, and the serializer suites). - `OptimisticTransaction` initializes the commit's `lastManifestCommit` by carrying forward the previous snapshot's reference; if the commit writes a new AMT checkpoint, the reference is updated to point at it. - `ConflictChecker` preserves/refreshes the reference when a transaction is logically reordered against a winning commit. - `Checksum` includes the reference in the CRC so it is available on the fast path. - Adds small `AMTUsageLogs` / `AMTUtils` helpers and a new `SnapshotLastManifestCommitSuite` that exercises how the reference is produced and resolved end to end. ## How was this patch tested? New `SnapshotLastManifestCommitSuite` plus updated `ActionSerializerSuite`, `CommitInfoSerializerSuite`, `InCommitTimestampSuite`, `AMTCheckpointTestBase`, `AMTSnapshotDiscoverySuite`, and `JsonUtilsSuite` cover serialization round-trips, carry-forward on commit, conflict-resolution refresh, and CRC population. ## Does this PR introduce _any_ user-facing change? No.
…ark (delta-io#7269) Surface Kernel accessors that already exist on the `*Impl` classes onto their public interfaces so connectors (Delta Spark) can use them without casting to internal types. ## Summary Abstract (implementers must provide; a throwing default would be a contract callers can't rely on): - `Scan.getScanFiles(Engine, boolean includeStats)` - `Snapshot.getProtocol()`, `getMetadata()`, `getCommitter()`, `getLogSegment()`, `getCreateCheckpointIterator(Engine)` Default, because the default is a real fallback rather than a throw: - `Snapshot.getDataPath()` (wraps `getPath()`), `getLogPath()`, `getCurrentCrcInfo()` (empty), `getLatestTransactionVersion(...)` (empty) Also `CommitRange.getDeltaFiles()`, and `DeltaHistoryManager`'s timestamp-resolution helpers widened from `SnapshotImpl` to `Snapshot`. **Breaking:** making the six methods abstract is source-breaking for `Scan`/`Snapshot` implementers outside this repo. Both interfaces are `@Evolving` and `kernelApi` is not under Mima. In-tree the only implementers are `ScanImpl`, `PaginatedScanImpl`, and `SnapshotImpl`, which already provide all of them. Return types that are still internal (`Protocol`, `Metadata`, `Path`, `CRCInfo`, `LogSegment`) match the existing pattern and are tracked by delta-io#4820 / delta-io#4821. `getClusteringColumnInfos()` keeps its throwing default from 4.3.0 — out of scope here. Follow-up Spark de-casts: delta-io#7273. ## Test plan - [x] `kernelApi` compile (main + test) + javafmt + checkstyle - [x] `CommitRangeBuilderSuite` (1161), `PostCommitSnapshotSuite` + `CreateCheckpointSuite` (36), `ScanSuite` (63) - [ ] CI ## Does this PR introduce any user-facing changes? Yes — new public Kernel API methods on `Snapshot`, `Scan`, and `CommitRange` (@SInCE 4.4.0). Six are abstract, so third-party implementers of `Scan`/`Snapshot` must add them.
…a-io#7343) ## Description <!-- - Describe what this PR changes. - Describe why we need the change. If this PR resolves an issue be sure to include "Resolves #XXX" to correctly link and close the issue upon merge. --> Enables the DSv2 batch write path to write partitioned column-mapped tables. Removes the reject guard, `Transaction.getWriteContext` physicalizes the partition directory and `AddFile` partition-value keys to `col-*` names. Identity transform for non-column-mapped tables. Tables that materialize partition columns into the Parquet body (`IcebergCompat`, `materializePartitionColumns`) are still rejected. ## How was this patch tested? <!-- If tests were added, say they were added here. Please make sure to test the changes thoroughly including negative and positive cases if possible. If the changes were tested in any way other than unit tests, please clarify how you tested step by step (ideally copy and paste-able, so that other reviewers can test and check, and descendants can verify in the future). If the changes were not tested, please explain why. --> New `V2WriteTest` cases: E2E `name`/`id`-mode partitioned column-mapped writes and then assert on-disk partition dirs use physical `col-*=value` names as well as round-trips on read. ## Does this PR introduce _any_ user-facing changes? <!-- If yes, please clarify the previous behavior and the change this PR proposes - provide the console output, description and/or an example to show the behavior difference if possible. If possible, please also clarify if this is a user-facing change compared to the released Delta Lake versions or within the unreleased branches such as master. If no, write 'No'. --> No
## Description <!-- - Describe what this PR changes. - Describe why we need the change. If this PR resolves an issue be sure to include "Resolves #XXX" to correctly link and close the issue upon merge. --> Adds a usage log event (`delta.timeTravel.atSyntaxUsage`) when the `@`-suffix time travel syntax is parsed for a path, recording whether a version (`@vN`) or timestamp (`@<yyyyMMddHHmmssSSS>`) was used, and its value. ## Does this PR introduce _any_ user-facing changes? <!-- If yes, please clarify the previous behavior and the change this PR proposes - provide the console output, description and/or an example to show the behavior difference if possible. If possible, please also clarify if this is a user-facing change compared to the released Delta Lake versions or within the unreleased branches such as master. If no, write 'No'. --> No
## 🥞 Stacked PR Use this [link](https://github.com/delta-io/delta/pull/7349/files) to review incremental changes. - [stack/flink-uc-5-browsing](delta-io#7247) [[Files changed](https://github.com/delta-io/delta/pull/7247/files)] [MERGED] - [**stack/flink-uc-6-catalog-correctness**](delta-io#7349) [[Files changed](https://github.com/delta-io/delta/pull/7349/files)] ← _this PR_ --------- ## Description Fix three correctness gaps in Flink SQL tables backed by Unity Catalog: - Map Delta `binary` columns to Flink `VARBINARY` so variable-length binary values can be written. - Preserve ordered partition keys from catalog column metadata and pass them to the catalog-backed Delta sink. - Return Flink's `UNKNOWN` statistics objects instead of rejecting statistics requests, allowing explicit target-column inserts to be planned. ## How was this patch tested? - `build/sbt "flink / javafmt"` - `build/sbt "flink / Test / javafmt"` - `build/sbt "flink / test"` — 248 tests, 0 failures - `build/sbt "flink / Compile / doc"` - Manual: four cross-engine integration workloads covering binary/null/boundary values, partitioned table creation with skewed and null keys, quoted columns with Unicode values, and schema evolution followed by writes from multiple engines. ## Does this PR introduce any user-facing changes? Yes. Flink SQL writes now handle variable-length binary data, preserve catalog table partitioning, and support plans that request unavailable catalog statistics.
…lify parallel commit reads (delta-io#7386) ## Description This patches a few rough edges of the `SingleCommit` API. - Adds `SingleCommit.getLogCommitActionsIteratorUnsafe` (and `DeltaFileProviderUtils.parallelReadAndParseLogCommitsAsSeqUnsafe`) for callers that must only read log commits. The iterator asserts that the commit is a log commit and throws an `IllegalStateException` if it encounters a manifest (checkpoint) commit. - Removes the unused `DeltaLog` argument from `DeltaFileProviderUtils`' `parallelReadAndParse*` helpers and from `parallelReadDeltaFiles`. The Hadoop configuration it built was never used, and the file-system safety check is already performed when the `DeltaLog` is initialized, so the argument was redundant. Callers in `HudiConverter`, `IcebergConverter`, and `IncrementalAMTWriter` are updated accordingly. - Renames `SingleCommit.sizeInBytes` to `commitFileSizeInBytes` and removes the now-unused `SingleCommit.path` accessor. ## How was this patch tested? Adds `AMTSingleCommitSuite`, which covers the new log-commit readers (including the assertion failure on a manifest commit and the parallel read path), and moves the existing `SingleCommit` test cases out of `DeltaLogSuite` into it. ## Does this PR introduce _any_ user-facing changes? No.
<!-- Thanks for sending a pull request! Here are some tips for you: 1. If this is your first time, please read our contributor guidelines: https://github.com/delta-io/delta/blob/master/CONTRIBUTING.md 2. If the PR is unfinished, add '[WIP]' in your PR title, e.g., '[WIP] Your PR title ...'. 3. Be sure to keep the PR description updated to reflect all changes. 4. Please write your PR title to summarize what this PR proposes. 5. If possible, provide a concise example to reproduce the issue for a faster review. 6. If applicable, include the corresponding issue number in the PR title and link it in the body. --> #### Which Delta project/connector is this regarding? <!-- Please add the component selected below to the beginning of the pull request title For example: [Spark] Title of my pull request --> - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description <!-- - Describe what this PR changes. - Describe why we need the change. If this PR resolves an issue be sure to include "Resolves #XXX" to correctly link and close the issue upon merge. --> `GenericColumnVector.getString` throws on a `null` entry, so `getPartitionRow` failed reading a `null` partition value. Check `isNullAt` first. ## How was this patch tested? <!-- If tests were added, say they were added here. Please make sure to test the changes thoroughly including negative and positive cases if possible. If the changes were tested in any way other than unit tests, please clarify how you tested step by step (ideally copy and paste-able, so that other reviewers can test and check, and descendants can verify in the future). If the changes were not tested, please explain why. --> ## Does this PR introduce _any_ user-facing changes? <!-- If yes, please clarify the previous behavior and the change this PR proposes - provide the console output, description and/or an example to show the behavior difference if possible. If possible, please also clarify if this is a user-facing change compared to the released Delta Lake versions or within the unreleased branches such as master. If no, write 'No'. --> No
<!-- Thanks for sending a pull request! Here are some tips for you: 1. If this is your first time, please read our contributor guidelines: https://github.com/delta-io/delta/blob/master/CONTRIBUTING.md 2. If the PR is unfinished, add '[WIP]' in your PR title, e.g., '[WIP] Your PR title ...'. 3. Be sure to keep the PR description updated to reflect all changes. 4. Please write your PR title to summarize what this PR proposes. 5. If possible, provide a concise example to reproduce the issue for a faster review. 6. If applicable, include the corresponding issue number in the PR title and link it in the body. --> #### Which Delta project/connector is this regarding? <!-- Please add the component selected below to the beginning of the pull request title For example: [Spark] Title of my pull request --> - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description A few small test suite refactors in the DSv2 Suites. ## How was this patch tested? Formatting only refactors. ## Does this PR introduce _any_ user-facing changes? No. --------- Co-authored-by: Murali Ramanujam <murali.ramanujam@databricks.com>
## Description This PR adds an optional `amtPassthrough` field to the `AddFile` action. It carries AMT-native fields (currently `spec_id`) alongside an `AddFile` so that they can round-trip through the Delta log without being lost. Changes: - Add `AMTPassthrough` case class and a new optional `amtPassthrough` field on `AddFile`. - Plumb the passthrough between `DataEntry` and `AddFile` (`AMTPassthrough.fromDataEntry`, `DataEntry.fromAddFile`), and add `fromRow` / `RowIndices` helpers to read the field out of an `InternalRow`. - Wire the new field through the affected read/write and clone paths. ## How was this patch tested? Added and updated unit tests: - `AMTSingleActionSerializerSuite`, `AMTSnapshotSuite`, `AMTSingleCommitSuite` - `ActionSerializerSuite`, `AddFileSuite`, `CheckpointsSuite` ## Does this PR introduce _any_ user-facing changes? No.
<!-- Thanks for sending a pull request! Here are some tips for you: 1. If this is your first time, please read our contributor guidelines: https://github.com/delta-io/delta/blob/master/CONTRIBUTING.md 2. If the PR is unfinished, add '[WIP]' in your PR title, e.g., '[WIP] Your PR title ...'. 3. Be sure to keep the PR description updated to reflect all changes. 4. Please write your PR title to summarize what this PR proposes. 5. If possible, provide a concise example to reproduce the issue for a faster review. 6. If applicable, include the corresponding issue number in the PR title and link it in the body. --> #### Which Delta project/connector is this regarding? <!-- Please add the component selected below to the beginning of the pull request title For example: [Spark] Title of my pull request --> - [ ] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [x] Other (build infrastructure) ## Description <!-- - Describe what this PR changes. - Describe why we need the change. If this PR resolves an issue be sure to include "Resolves #XXX" to correctly link and close the issue upon merge. --> MuleSoft was added in delta-io#4750 as a backup Maven repository, but it currently appears before the GCS Maven Central mirror and Maven Central. Coursier stops at the first repository that resolves a module, so common dependencies become tied to MuleSoft and builds fail when that repository returns HTTP 502. Move MuleSoft after the two Maven Central resolvers so it is used as the intended fallback. ## How was this patch tested? <!-- If tests were added, say they were added here. Please make sure to test the changes thoroughly including negative and positive cases if possible. If the changes were tested in any way other than unit tests, please clarify how you tested step by step (ideally copy and paste-able, so that other reviewers can test and check, and descendants can verify in the future). If the changes were not tested, please explain why. --> - Ran `git diff --check`. - Verified that Maven Central precedes MuleSoft in `build/sbt-config/repositories`. - CI exercises the resolver configuration. ## Does this PR introduce _any_ user-facing changes? <!-- If yes, please clarify the previous behavior and the change this PR proposes - provide the console output, description and/or an example to show the behavior difference if possible. If possible, please also clarify if this is a user-facing change compared to the released Delta Lake versions or within the unreleased branches such as master. If no, write 'No'. --> No.
## What changes were proposed in this pull request? This change centralizes creation of Kernel engines used by the Delta V2 connector. `SerializableReadOnlySnapshot` can reconstruct a scan with a caller-owned engine, allowing scan construction and file iteration to share the same engine. ## How was this patch tested? Verified with the Delta V2 scan test suite.
…o#7389) ## 🥞 Stacked PR Use this [link](https://github.com/delta-io/delta/pull/7389/files) to review incremental changes. - [stack/flink-uc-narrow-integer-types](delta-io#7389) [[Files changed](https://github.com/delta-io/delta/pull/7389/files)] --------- ## Description Unity Catalog table metadata represents TINYINT and SMALLINT as byte and short in type JSON. Map those values to Flink TinyIntType and SmallIntType while preserving column nullability. ## Testing - build/sbt "flink / javafmt" - build/sbt "flink / Test / javafmt" - build/sbt "flink / testOnly io.delta.flink.sink.sql.FlinkUnityCatalogTableTest" (3 passed) - build/sbt "flink / test" (249 total, 245 passed, 4 skipped) - Cross-engine integration matrix for TINYINT and SMALLINT values, schema discovery, and reader parity (36 passed) Signed-off-by: Timothy Wang <timothy.art@gmail.com>
## 🥞 Stacked PR Use this [link](https://github.com/delta-io/delta/pull/7392/files) to review incremental changes. - [stack/flink-uc-void-catalog-schema](delta-io#7392) [[Files changed](https://github.com/delta-io/delta/pull/7392/files)] --------- #### Which Delta project/connector is this regarding? - [ ] Spark - [ ] Standalone - [x] Flink - [ ] Kernel - [ ] Other ## Description A Delta table schema can contain a `VOID` column. For example, a column created from `CAST(NULL AS VOID)` is stored in the table schema even though Parquet files do not contain physical values for it. The Flink catalog schema parser did not recognize `void`, so Flink could not resolve the table. This change maps `void` to the existing Flink `NullType`. It changes only schema parsing; values in a `VOID` column remain null. ## How was this patch tested? Added a focused regression test that builds a catalog table with an `id` column and a `pending VOID` column, registers it in a Flink catalog, and resolves it through `TableEnvironment`. Passed: - `build/sbt "flink / javafmtCheckAll"` - `build/sbt "flink / testOnly io.delta.flink.sink.sql.FlinkUnityCatalogTableTest"` (4 tests) ## Does this PR introduce _any_ user-facing changes? Yes. Flink can now resolve a Delta catalog table whose schema contains a top-level `VOID` column. Signed-off-by: Timothy Wang <timothy.art@gmail.com>
…io#7405) <!-- Thanks for sending a pull request! Here are some tips for you: 1. If this is your first time, please read our contributor guidelines: https://github.com/delta-io/delta/blob/master/CONTRIBUTING.md 2. If the PR is unfinished, add '[WIP]' in your PR title, e.g., '[WIP] Your PR title ...'. 3. Be sure to keep the PR description updated to reflect all changes. 4. Please write your PR title to summarize what this PR proposes. 5. If possible, provide a concise example to reproduce the issue for a faster review. 6. If applicable, include the corresponding issue number in the PR title and link it in the body. --> #### Which Delta project/connector is this regarding? <!-- Please add the component selected below to the beginning of the pull request title For example: [Spark] Title of my pull request --> - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description V2 selected-file descriptors expose partition values through MapValue while their other metadata is value-owned. Store partition values in immutable maps for consistent ownership across Kernel and V1 AddFile inputs. Centralize partition-map construction and map-based InternalRow conversion while preserving null handling, physical-name lookup, and timezone-aware casts. Factor Kernel-to-Delta deletion-vector descriptor conversion through one field-for-field helper. <!-- - Describe what this PR changes. - Describe why we need the change. If this PR resolves an issue be sure to include "Resolves #XXX" to correctly link and close the issue upon merge. --> ## How was this patch tested? CI and - `build/sbt "sparkV2/testOnly io.delta.spark.internal.v2.utils.PartitionUtilsTest"` - `build/sbt "sparkV2/testOnly io.delta.spark.internal.v2.read.DeltaV2ScanTest -- *testSelectedFilesTracksPlannedAndRuntimeFilteredFiles*"` - `sparkV2 / Compile / javafmt <!-- If tests were added, say they were added here. Please make sure to test the changes thoroughly including negative and positive cases if possible. If the changes were tested in any way other than unit tests, please clarify how you tested step by step (ideally copy and paste-able, so that other reviewers can test and check, and descendants can verify in the future). If the changes were not tested, please explain why. --> ## Does this PR introduce _any_ user-facing changes? <!-- If yes, please clarify the previous behavior and the change this PR proposes - provide the console output, description and/or an example to show the behavior difference if possible. If possible, please also clarify if this is a user-facing change compared to the released Delta Lake versions or within the unreleased branches such as master. If no, write 'No'. --> No
…range has no file actions (delta-io#7404) #### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description A Delta Sharing streaming query started with `startingTimestamp` fails on its first batch when the requested version range contains no file actions: ``` [DELTA_TIMESTAMP_GREATER_THAN_COMMIT] The provided timestamp (...) is after the latest version available to this table (1970-01-01 00:00:00.499) ``` The root cause is that `startingTimestamp` is resolved twice. `DeltaFormatSharingSource` resolves it against the sharing server, but it also passes the raw user options to the wrapped `DeltaSource`, which resolves the timestamp again against the locally constructed delta log. That local log is synthetic: versions with no file actions get a commit `modificationTime` of 0, which the history manager monotonizes into small millisecond values and reports as the table's latest commit. When the range does contain file actions, real timestamps land in the log and the second lookup happens to succeed, which is why this only reproduces on an empty first batch. This change adds `parametersForDeltaSource`, which converts `startingTimestamp` into the version it resolves to before building the wrapped `DeltaSource`'s options. A concrete version only needs an existence check against commit filenames, never a timestamp lookup. The timestamp is converted rather than dropped, because with no starting option the wrapped `DeltaSource` would instead read a snapshot of the latest version. The behavior is gated by a new, default-off internal flag, `spark.sql.delta.sharing.streamingConvertStartingTimestampToVersion`. `startingVersion` is left untouched. ## How was this patch tested? Two new unit tests in `DeltaFormatSharingSourceSuite`, over a table whose last version is an `OPTIMIZE` (no `dataChange` actions, so the server returns no file actions for that range): - With the flag on, the first batch is empty and the offset advances past the range. - With the flag off, the query still fails with `DELTA_TIMESTAMP_GREATER_THAN_COMMIT`, pinning the regression. ## Does this PR introduce _any_ user-facing changes? No. The new behavior is gated behind a default-off internal configuration flag.
#### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other ## Description Make Spark 4.2.0 the default build version now that it is supported. This updates default and unsuffixed artifact publishing, Unity Catalog setup, and Scala example resolution while retaining suffixed builds for Spark 4.0 and 4.1. ## How was this patch tested? - python3.11 project/tests/test_cross_spark_publish.py - build/sbt exportSparkVersionsJson - Scala example dependency resolution - bash -n project/scripts/setup_unitycatalog_main.sh - git diff --check ## Does this PR introduce any user-facing changes? Yes. Unsuffixed Spark-dependent artifacts now target Spark 4.2.0. --------- Signed-off-by: Timothy Wang <timothy.art@gmail.com>
…et reader (delta-io#7378) <!-- Thanks for sending a pull request! Here are some tips for you: 1. If this is your first time, please read our contributor guidelines: https://github.com/delta-io/delta/blob/master/CONTRIBUTING.md 2. If the PR is unfinished, add '[WIP]' in your PR title, e.g., '[WIP] Your PR title ...'. 3. Be sure to keep the PR description updated to reflect all changes. 4. Please write your PR title to summarize what this PR proposes. 5. If possible, provide a concise example to reproduce the issue for a faster review. 6. If applicable, include the corresponding issue number in the PR title and link it in the body. --> #### Which Delta project/connector is this regarding? <!-- Please add the component selected below to the beginning of the pull request title For example: [Spark] Title of my pull request --> - [ ] Spark - [ ] Standalone - [ ] Flink - [x] Kernel - [ ] Other (fill in here) ## Description Kernel wraps its own `InputFile` before handing it to `parquet-mr`, so the wrapper is never a `parquet-mr` `HadoopInputFile`. Every entry point that takes an `InputFile` without a configuration therefore falls back to `new HadoopParquetConfiguration()`, which wraps a fresh `Configuration(loadDefaults=true)`. Hadoop loads such a configuration's properties lazily, so the first property read -- which `parquet-mr` issues immediately while building the read options -- scans the classpath for `core-default.xml` and friends under a JVM-global lock. With one of these per Parquet file, concurrent readers serialize on that lock. Pass the configuration the file is already being read with instead. Its properties are loaded once and memoized on the instance, so the scan disappears. This affects all three call sites that were constructing a fresh configuration: the footer read and the reader build in ParquetFileReader (both per Parquet file), and the footer read in ParquetStatsReader. ## How was this patch tested? Covered by existing tests ## Does this PR introduce _any_ user-facing changes? <!-- If yes, please clarify the previous behavior and the change this PR proposes - provide the console output, description and/or an example to show the behavior difference if possible. If possible, please also clarify if this is a user-facing change compared to the released Delta Lake versions or within the unreleased branches such as master. If no, write 'No'. -->
## 🥞 Stacked PR Use this [link](https://github.com/delta-io/delta/pull/7398/files) to review incremental changes. - [stack/flink-writer-completion-failures](delta-io#7398) [[Files changed](https://github.com/delta-io/delta/pull/7398/files)] --------- #### Which Delta project/connector is this regarding? - [ ] Spark - [ ] Standalone - [x] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description Flink keeps one file writer for each active table partition. At checkpoint time, it closes those writers and collects the new Delta files. Before this patch, a failure while closing a writer was sent through a Caffeine cache callback. Caffeine logged the exception but did not return it to the caller. Flink could therefore report a successful checkpoint with no data files, even though writing the rows had failed. This patch saves the first writer-completion failure and throws it immediately after all cached writers are closed. The same small check is applied to append and upsert mode. A failed writer stays failed, so a later checkpoint cannot silently skip its rows. ## How was this patch tested? - Added append and upsert regression cases that buffer a malformed row, then verify that `prepareCommit()` returns the exact file-completion failure instead of an empty result. - `build/sbt 'flink/testOnly io.delta.flink.sink.DeltaSinkWriterTest io.delta.flink.sink.DeltaSinkTest'` (27 tests passed, 1 existing test ignored) - `build/sbt 'flink/javafmt'` ## Does this PR introduce _any_ user-facing changes? Yes. A Flink job now fails its checkpoint when a cached file writer cannot finish. Previously, that failure could be logged while the checkpoint continued with an empty set of files. Signed-off-by: Timothy Wang <timothy.art@gmail.com>
…io#7391) ## What changes were proposed in this pull request? Rename the DSv2 snapshot manager interface from `DeltaSnapshotManager` to `DeltaV2SnapshotManager` for consistency with the V2 namespace convention. The old `DeltaSnapshotManager` abstract class is removed and its methods (`loadLatestSnapshot`, `loadSnapshotAt`) are pushed into the concrete implementations (`PathBasedSnapshotManager`, `UCManagedTableSnapshotManager`) which now implement `DeltaV2SnapshotManager` directly as a marker/contract interface. All callers across catalog, read, write, streaming, and changelog modules are updated to use the new type name. ## How was this patch tested? Existing tests updated to use the new name. No behavioral changes — purely a rename/restructure.
This adds an RFC, corresponding to issue delta-io#6081, to add a nanoseconds timestamps feature to the protocol. Proof of concept implementation, for `delta-rs`: * delta-io/delta-kernel-rs#1897 together with delta-io/delta-rs#4215 --------- Co-authored-by: Itamar Turner-Trauring <itamar@pythonspeed.com> Co-authored-by: OussamaSaoudi <45303303+OussamaSaoudi@users.noreply.github.com> Co-authored-by: Adam Reeve <adam.reeve@gr-oss.io>
## 🥞 Stacked PR Use this [link](https://github.com/delta-io/delta/pull/7402/files) to review incremental changes. - [stack/kernel-unpublished-conflict-retry](delta-io#7402) [[Files changed](https://github.com/delta-io/delta/pull/7402/files)] --------- <!-- Thanks for sending a pull request! Here are some tips for you: 1. If this is your first time, please read our contributor guidelines: https://github.com/delta-io/delta/blob/master/CONTRIBUTING.md 2. If this PR is unfinished, add '[WIP]' in the title. 3. Keep the description updated as the change evolves. --> #### Which Delta project/connector is this regarding? - [ ] Spark - [ ] Standalone - [ ] Flink - [x] Kernel - [ ] Other ## Description Two writers can start from the same table version. One writer may win the next version while the other writer is still trying to commit. For a catalog-managed table, the winning commit can be accepted before its JSON file appears in the Delta log. Kernel knew that a conflict happened, but it could not read the winner's actions from the log. It then failed with the internal error `No winning commits found`. Kernel cannot safely resolve a conflict without reading the winner's actions. This change reports a `ConcurrentWriteException` for this catalog-managed case. That tells the caller to load the latest table state and run the transaction again. Ordinary filesystem tables keep their previous invariant check because their winning JSON file must already exist. ## How was this patch tested? - Added a regression where a catalog-managed committer reports a conflict before the winner is available in the Delta log. It now returns `ConcurrentWriteException`. - Added the rejected path for a filesystem table, which must still report the missing winner as an invariant error. - Ran `kernelDefaults/testOnly io.delta.kernel.defaults.ConflictResolutionSuite`: 6 tests passed. - Ran a repeated three-writer append workload against a catalog-managed table: passed. ## Does this PR introduce _any_ user-facing changes? Yes. A connector that supports retrying `ConcurrentWriteException` can now reload and retry this catalog-managed conflict instead of receiving an internal `IllegalStateException`. Filesystem-managed tables are unchanged. Signed-off-by: Timothy Wang <timothy.art@gmail.com>
## 🥞 Stacked PR Use this [link](https://github.com/delta-io/delta/pull/7401/files) to review incremental changes. - [stack/flink-bounded-stream-final-commit](delta-io#7401) [[Files changed](https://github.com/delta-io/delta/pull/7401/files)] --------- #### Which Delta project/connector is this regarding? - [ ] Spark - [ ] Standalone - [x] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description A bounded Flink job can finish before its last rows are included in a normal checkpoint. Flink gives those final rows the number of the next checkpoint. The Delta pre-commit aggregator replaced that real checkpoint number with `Long.MAX_VALUE`. The committer waits for a real completed checkpoint number, so it never selected the final files. For example, a bounded job could report that it finished even though none of its 10 rows appeared in the table. This change keeps the checkpoint number that came from Flink. It also fails clearly if messages from two checkpoint numbers appear in one aggregation window. ## How was this patch tested? Added a regression test with a bounded source that finishes immediately instead of waiting for extra checkpoints. It checks that all 10 rows are committed. Added focused operator tests that prove checkpoint 42 stays checkpoint 42 through the final aggregation. They also verify that messages or barriers with conflicting checkpoint numbers are rejected with clear errors. The bounded-source test failed on the old behavior with 0 of 10 rows and passed after the fix. The focused tests passed: ```text ./build/sbt \ 'flink/Test/javafmtAll' \ 'flink/testOnly io.delta.flink.sink.DeltaWriterResultAggregatorTest io.delta.flink.sink.DeltaSinkTest' ``` ## Does this PR introduce _any_ user-facing changes? Yes. A bounded Flink streaming job now commits the rows produced at the end of its input instead of finishing successfully while leaving those rows out of the Delta table. Signed-off-by: Timothy Wang <timothy.art@gmail.com>
…elog (delta-io#7412) #### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description Extends CDC test coverage for the V2 changelog read path. Follow-up to delta-io#7314, which routed the generated Delete/Merge/Update CDC suites through `SELECT ... CHANGES` via `computeCDC`. This PR does the same for the `cdcRead` entry point used by the hand-written `DeltaCDCScalaSuite` / `DeltaCDCSQLSuite`. A new `ChangelogV2CdcReadMixin` overrides `cdcRead` to issue `<table> CHANGES FROM VERSION ... TO VERSION ...` (TableCatalog.loadChangelog -> DeltaV2ChangelogScan / DeltaV2ChangelogBatch) under `v2.enableMode=STRICT`, so the existing suite bodies validate the V2 changelog read path. Two concrete suites mix it in: - `DeltaCDCScalaChangelogV2Suite` - `DeltaCDCSQLChangelogV2Suite` Sister of `ChangelogV2CDCUtilMixin` from delta-io#7314. The mixin requires row tracking (the V2 changelog read path needs it), and scopes STRICT mode to the read only so table setup and DML still go through the default path. Tests that do not apply to the V2 read path are excluded by name: V1-specific error semantics, mixed timestamp/version boundaries (the `CHANGES` grammar only allows VERSION..VERSION or TIMESTAMP..TIMESTAMP), timestamp -> version resolution (unsupported on the JVM kernel backend today), open-ended reads (V1 resolves the end bound lazily at execution time, V2 eagerly at analysis time), and protocol assertions that do not expect the row-tracking / domainMetadata table features. The `SELECT ... CHANGES` clause only parses on Spark 4.2+, so tests cancel on older versions via the `ChangelogSyntaxSupportedShim`. ## How was this patch tested? This PR is tests only. Both new suites run the existing `DeltaCDCScalaSuite` / `DeltaCDCSQLSuite` bodies through the V2 changelog read path (25 tests each pass, with the inapplicable families ignored as described above). ## Does this PR introduce _any_ user-facing changes? No.
#### Which Delta project/connector is this regarding? - [X] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description This PR adds retries to snapshot reconstruction after a commit. This is necessary because the storage layer may occasionally return a cached listing result that does not yet include the newly committed files, effectively violating the expected storage guarantees. The retry logic is intended to mitigate these short-lived inconsistencies and prevent false-positive job failures caused by incomplete listing results. If the maximum number of retries is set to 0, these changes are effectively a no-op. Also, if the listing succeeds on the first attempt, the retry logic has no impact on behavior. So the changes are harmless. ## How was this patch tested? A new unit test. All existing tests are passed. ## Does this PR introduce _any_ user-facing changes? No. Co-authored-by: Aleksei Shishkin <aleksei.shishkin@databricks.com>
…ns (delta-io#7377) ## Description Adds `AMTIncrementalWriteSuite`, covering the incremental AMT write path (`IncrementalAMTWriter.writeIncremental`) end to end, and fixes a metric defect that writing it surfaced. ## How was this patch tested? This PR is almost entirely tests. ## Does this PR introduce _any_ user-facing changes? No
…7764) ## Description `DeltaV2OptimisticTransaction` receives `AddFile` partition values keyed by physical column names. Kernel's default write context treats those keys as logical names, which misinterprets column-mapped partitions. This change resolves partition data types by physical name and calls `getWriteContext` with `PartitionKeyType.PHYSICAL`. The append tests are extended to cover `name` and `id` column mapping with INT partition values, asserting that committed files preserve their physical partition keys and values. The mapped physical names carry the generated `col-` prefix; the existing unmapped null partition case remains. ## How was this patch tested? Extended the partition-append tests in `DeltaV2OptimisticTransactionSuite` to assert committed files preserve their physical partition keys/values, and added cases for both column mapping modes (`name` and `id`) alongside the existing INT, DATE, TIMESTAMP, DECIMAL, BOOLEAN, STRING and null partition cases. ## Does this PR introduce _any_ user-facing change? No.
…json expression" (delta-io#7765) Reverts delta-io#7756
…o#7755) ### Description When the latest commit file is removed, `CachedSnapshotManager` can reconstruct an older surviving version. It currently discards that result and retains the higher cached version. Reuse now requires equal table ID and version, so the recovered snapshot is installed, matching `DeltaLog`. The regression test deletes a real commit file with in-commit timestamps enabled and disabled, compares the recovered version and files with `DeltaLog`, and verifies cache replacement and subsequent reuse. `KernelContext.materializeHadoopConf` is also available throughout the internal `v2` package, allowing table managers to use the same active-session configuration materialization as Kernel I/O. ### How was this patch tested? `build/sbt 'sparkV2/testOnly io.delta.spark.internal.v2.tablemanager.CachedSnapshotManagerSuite -- -z "missing current commit"'`: both ICT variants pass, with production and test compilation. `build/sbt 'sparkV2/testOnly io.delta.spark.internal.v2.kernel.KernelContextSuite'`: all seven existing context tests pass after the visibility change, with production and test compilation. The manager and regression-test files pass the repository's Scala style rules. Module-wide `sparkV2/scalastyle` previously reported a `deltahadoopconfiguration` violation at the existing configuration-materialization call in `KernelContext.scala:40`. That call is unchanged; the visibility edit does not address this existing lint issue.
## 🥞 Stacked PR Use this [link](https://github.com/delta-io/delta/pull/7708/files) to review incremental changes. - [**stack/query-context-pr6-cache**](delta-io#7708) [[Files changed](https://github.com/delta-io/delta/pull/7708/files)] ← _this PR_ - [stack/query-context-pr4-streaming](delta-io#7706) [[Files changed](https://github.com/delta-io/delta/pull/7706/files/93ea3351f4a78070ef865fe5a7602df413c40c55..26318ace51549e9fd937fcfb0e38cfc5c5386342)] - [stack/query-context-pr5-write](delta-io#7707) [[Files changed](https://github.com/delta-io/delta/pull/7707/files/26318ace51549e9fd937fcfb0e38cfc5c5386342..25bfe5820e9124135ded36c6ea8c70b338ecd1cf)] --------- ## Summary - Builds on delta-io#7704. - Routes request-scoped query context through the cached Delta V2 snapshot manager. - Preserves cache identity while keeping per-operation catalog metadata out of retained state. - Adds focused coverage for legacy, request-scoped, and empty-context catalog routing. ## Testing - CachedSnapshotManagerSuite - Scala formatting checks - git diff --check --------- Signed-off-by: Vishnu Chandrashekhar <vishnu.c@databricks.com>
…e path (delta-io#7692)" (delta-io#7779) #### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description This reverts delta-io#7692 (b83b1d8) for now. A test issue was found after merge: the new `DeltaSinkSuite` shredding check looks up the variant column by its logical name, so it fails under column mapping (see delta-io#7778). The change will be resubmitted with that fix included. ## How was this patch tested? Clean `git revert`; the touched files are identical to their state before delta-io#7692. ## Does this PR introduce _any_ user-facing changes? Yes, within `master` only: it undoes delta-io#7692, so the DSv2 write path no longer applies `delta.enableVariantShredding` to its Parquet write options, as before delta-io#7692. This pull request and its description were written by Isaac. Signed-off-by: Liang-Chi Hsieh <viirya@gmail.com>
…o#7780) #### Which Delta project/connector is this regarding? - [x] Spark ## Description Collapsing file actions during replay can lose an AMT back reference when a later add or remove has no reference. Add an opt-in `retainFirstBackreference` mode to `InMemoryLogReplay` that retains the first non-empty reference for each unique file action and carries it onto the surviving add or tombstone. The default replay behavior is unchanged. Enable this mode in the minor-compaction test helper for AMT tables. Cover re-adds, remove-then-add and add-then-remove sequences, files without references, and compaction windows spanning a content root. ## How was this patch tested? Adds replay and minor-compaction coverage in `AMTBackReferenceSuite`, including tests across AMT checkpoint scenarios. `git diff --check` passes. Scala tests have not been run locally; CI validation is pending. ## Does this PR introduce _any_ user-facing changes? No. Co-authored-by: Yumingxuan Guo <171533467+yumingxuanguo-db@users.noreply.github.com>
…on conflict (delta-io#7784) ## Description `AMTWriterManager.handleLosingOptimizeCheckpoint` decides how to salvage the manifest tree an OPTIMIZE checkpoint already wrote when it loses its target version to concurrent commits. When the checkpoint lost only to log-only winners and every winner touched files created strictly after the losing tree's content-root version, the tree is still exact and is recommitted as-is. Otherwise a winner changed content the base tree describes, and previously every such case -- full or incremental -- raised `FullAMTWriteFailedWithConflict`, which unwinds out of the transaction so the checkpoint is regenerated from scratch in a new transaction. For a losing **incremental** checkpoint that round trip is unnecessary: its base tree (the prior checkpoint it extends) is unchanged and its fold window is small. This change recreates the incremental checkpoint from scratch against the new pre-commit log segment inside the same transaction (via `materialize(..., incremental = true)`), recording a new `treeOutcome` value `REBUILT_OPTIMIZE_CHECKPOINT_INCREMENTAL_TREE`. A losing **full** checkpoint keeps the existing behavior of signaling `FullAMTWriteFailedWithConflict` so the checkpoint is regenerated against a refreshed snapshot. The `REGENERATE_VIA_TXN_RETRY` outcome is renamed to `RETRY_VIA_NEW_TXN` for clarity. ## How was this patch tested? Updated `AMTConflictResolutionSuite`. ## Does this PR introduce _any_ user-facing changes? No.
## Description Populates a minimum sequence number for each manifest when writing adaptive metadata tree (AMT) checkpoints. Tracking the minimum sequence number across the live entries in a manifest lets readers reason about the lowest sequence number present without scanning every entry. Key changes: - Compute and record the manifest minimum sequence number in `AMTWriteHelper` when building manifest entries. - Carry the field through the AMT action model. ## How was this patch tested? Added coverage in `AMTCheckpointWriteSuite` and `AMTIncrementalWriteSpillSuite`. ## Does this PR introduce _any_ user-facing changes? No. Co-authored-by: Jinhua Song <jinhua.song@databricks.com>
… Snapshot primitives (delta-io#7782) ## Description Threads the Delta `Snapshot` through the DSv2 connector's write and metadata-only delete paths instead of passing low-level Kernel snapshot types around. Schema, protocol, and table size are read through the `Snapshot` facade, and the underlying Kernel snapshot is unwrapped only where a Kernel transaction requires it. Also: - Adds a small Scala `Option[Long]` -> `OptionalLong` conversion helper. - Reads table size in bytes via the CRC-backed scalar lookup when available. ## How was this patch tested? Existing and updated unit tests for the DSv2 write/delete paths (`DeltaV2WriteTest`, `DeltaV2BatchWriteTest`, `DeltaV2StreamingWriteTest`, `DeltaV2WriteContextTest`, `DeltaV2BatchWriteContextTest`, `ScalaUtilsTest`, and the metadata-only delete tests), which now load the snapshot through `PathBasedSnapshotManager`. ## Does this PR introduce _any_ user-facing change? No.
…SparkSession (delta-io#7775) #### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description Resolves delta-io#7776 `ReorgTableHelper.filterParquetFilesOnExecutors` reads Parquet footers inside a `mapPartitions` closure, which runs on executors. Since delta-io#5697, that closure gets its conf with `SparkSession.active.sessionState.conf`. Executors have no active or default `SparkSession`, so `REORG TABLE ... APPLY (PURGE)` fails on any cluster with separate executor JVMs: ``` [INTERNAL_ERROR] No active or default Spark session found SQLSTATE: XX000 at org.apache.spark.sql.SparkSessionCompanion.active(SparkSession.scala:1036) at org.apache.spark.sql.delta.commands.ReorgTableHelper.$anonfun$filterParquetFilesOnExecutors$1(ReorgTableHelper.scala:116) ``` The existing tests pass because they run in local mode, where task threads fall back to the JVM-wide default session. This PR reads the conf with `SQLConf.get` instead. On executors it returns the read-only conf that Spark sends from the driver with each task. This keeps the `ParquetToSparkSchemaConverter(SQLConf)` constructor from delta-io#5697, so the Spark 4.0.1 binary-compatibility fix is unchanged. `SQLConf.get` does not lose any session conf. `Dataset.collect()` runs under `SQLExecution.withSQLConfPropagated`, which copies every `spark.*` conf of the running session into the task's local properties, and `ReadOnlySQLConf` reads them from there. It is also more accurate in local mode: `SparkSession.active` on a task thread fell back to the JVM-wide default session, which is not necessarily the session running the `REORG`. `DeltaPurgeOperation` (`REORG ... APPLY (PURGE)`) and `DeltaRewriteTypeWideningOperation`, used when dropping the type widening feature, both use this helper, so both are fixed. ## How was this patch tested? - Added `DeltaReorgSuite` test "Purge DVs when tasks cannot see a SparkSession". It clears the default session before the purge, so the tasks behave as they would on a real executor. - Without the fix it fails with `No active or default Spark session found`. - With the fix it passes. - Ran the full `DeltaReorgSuite`; all tests pass. - Manual test on Kubernetes (Spark 4.2.0, 1 driver and 1 executor pod, Unity Catalog managed tables): - Before the fix, `REORG TABLE ... APPLY (PURGE)` failed on every table. - After the fix it succeeded on all 5 tables, and the table with a deletion vector was rewritten: `numDeletionVectorsRemoved: 1`, `numAddedFiles: 1`, `numRemovedFiles: 1`. ## Does this PR introduce _any_ user-facing changes? Yes. `REORG TABLE ... APPLY (PURGE)` works again on clusters with separate executors; it has failed there since 4.0.1. There are no API or config changes. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Signed-off-by: s-maling-celonis <s.maling@celonis.com>
…o#7272) #### Which Delta project/connector is this regarding? - [ ] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [x] Other (Protocol RFC) ## Description Propose `concurrentIdentityColumns`, a writer-only table feature that introduces a service-backed identity column. Instead of serializing identity generation through the commit, a service-backed identity column is bound to a catalog-hosted monotonic sequence. Writers reserve disjoint ranges and assign values locally, so many writers (and many tasks within a writer) can generate identity values in parallel. Design doc: [[DESIGN] Concurrent Identity Columns](https://docs.google.com/document/d/1QKLLPgF2TKWW0bBJPtMW_sXwr_CzY7kV_tfgshIhVDM/edit?tab=t.0#heading=h.h20oo4jjwcoq) The RFC specifies only the Delta log part. The sequence-allocation RPC itself stays a catalog concern: - New column-metadata key `delta.identity.concurrent.sequenceId`, an opaque pointer to the column's current sequence, mutually exclusive with `delta.identity.highWaterMark`. - Enablement requires `identityColumns` and `catalogManaged` in addition to `concurrentIdentityColumns` (Reader v3 / Writer v7). - `concurrentIdentityColumns` is a table-level mode: every identity column must be service-backed, no mixing with classic high-water-mark columns. - Value semantics: `start + k * step` for monotonically increasing k, never reused but not necessarily contiguous (gaps from unused reserved ranges are acceptable). - `SYNC IDENTITY`, `SET TBLPROPERTIES` (upgrade) and `DROP FEATURE` (downgrade) transitions, each seeded strictly past the existing extreme so no written value is ever regenerated. - Reader requirements: none beyond catalog-managed table support. - A non-normative appendix sketches the expected sequence-service calls (createSequence / reserveRange / dropSequence). This is a documentation-only change adding protocol_rfcs/concurrent-identity-columns.md and listing it in protocol_rfcs/README.md. see delta-io#7572 Implementation: delta-io#7618 --------- Co-authored-by: Johan Lasperas <johan.lasperas@databricks.com>
## Description Use the streaming test framework's configured timeout for the CDC AvailableNow read-limit test instead of a fixed 10-second wait. This keeps the wait bounded while avoiding failures on slower test environments. ## How was this patch tested? Local execution was attempted, but the SBT launcher could not be downloaded from the configured Maven repositories. CI will exercise the affected suite.
…elta-io#7783) #### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description Re-land of delta-io#7692, which was reverted in delta-io#7779. Compared with delta-io#7692, this adds two test-only fixes: - `DeltaSinkSuite#targetHasShreddedVariant` looked up the variant column in the Parquet footer by its logical name `v`. Under column mapping the footer carries the physical name (`col-<uuid>`), so the new streaming test failed with `v not found in message spark_schema` when run in the column-mapping subclasses of `DeltaSinkSuite`. It now resolves the physical name via `DeltaColumnMapping.getPhysicalName`. - A `V2WriteTest` javadoc is reworded to describe the `REORG ... APPLY (UNSHRED VARIANT)` case in engine-neutral terms. The rest of the change is identical to delta-io#7692. The DSv2 write path built its Parquet write options from the caller's write options alone (`DeltaV2WriteContext`), so `spark.sql.variant.inferShreddingSchema` was never among them. When that option is absent, `ParquetOptions.inferShreddingForVariant` falls back to the session default for the conf, so the writer shredded variant columns on tables that never enabled `delta.enableVariantShredding`. The V1 write path passes the table's value explicitly (`TransactionalWrite`). Two consequences, both reproduced by the new tests: - A table ends up with shredded data files while its protocol advertises neither `variantShredding` nor `variantShredding-preview`. Shredded variant requires that reader-writer feature, so a client honoring the protocol is not told it needs shredding support. - A per-table opt-out is silently undone. After the property is set to `false`, or after `REORG TABLE ... APPLY (UNSHRED VARIANT)` -- which unsets it and rewrites the files -- the next DSv2 write puts a shredded file back into the live snapshot. V1 does not. This is reachable where shredded writes are allowed in the session (`spark.sql.variant.writeShredding.enabled`) *and* the table opted out. **The fix** passes the table's `delta.enableVariantShredding` value into the Parquet write options, mirroring V1. It is resolved through `DeltaConfigs.ENABLE_VARIANT_SHREDDING` so alternate property keys and the default behave as they do on the V1 path. The option is set via `VariantShreddingShims` because the underlying conf exists only in Spark 4.1+ while `sparkV2` also cross-builds against 4.0, where the shim yields no entry and the writer keeps its default behavior. Any caller spelling of the option that differs from it only in case is dropped before the table-derived value is applied: `ParquetOptions` reads write options through a case-insensitive map, so leaving both in place would let map iteration order rather than the table property decide the effective value. **Streaming writes** need one more step. `V2Writes` rebuilds the write per micro-batch, but the sink relation -- and with it the snapshot the write state is built from -- is captured once for a non-transactional catalog, and turning shredding off changes neither the schema nor the protocol, since the `variantShredding` feature stays in place. The existing per-epoch schema and protocol guard therefore cannot see it, and the epoch would commit stale shredded output on top of the opt-out. `DeltaV2StreamingWrite#commit` now also compares the shredding decision against the reloaded snapshot and fails the epoch when it diverges. What it compares is the *effective* layout decision rather than the raw property: the version shim must supply the inference option, `spark.sql.variant.writeShredding.enabled` must be on, and the data schema must contain a variant column -- mirroring `ParquetOptions#inferShreddingForVariant` and `ParquetUtils#needShreddingInference`. A property toggle that cannot change any file therefore does not terminate the stream. The already-committed epoch check runs ahead of the layout guards, so replaying a committed epoch remains an idempotent no-op, as `StreamingWrite.commit` allows the engine to call it more than once per epoch. ## How was this patch tested? New unit tests. **Batch layout follows the property** -- `V2WriteTest#variantWriteFollowsTableShreddingProperty`, parameterized over both values of `delta.enableVariantShredding`. It writes the same row through V1 and through DSv2 into two tables carrying the same property, then reads the Parquet footers: - property `false`: neither writer may emit a `typed_value` child on the variant column. This is the case that fails without the fix (DSv2 shredded, V1 did not). - property `true`: both writers must still emit it, so the fix does not disable shredding outright. - Either layout must read back identically through both connectors (`variant_get` on two fields). **Feature present, property not enabled** -- the state `REORG TABLE ... APPLY (UNSHRED VARIANT)` leaves behind, where the feature stays in the protocol while the property no longer enables shredding, so the test cannot conflate property gating with protocol gating -- `V2WriteTest#variantWriteFollowsPropertyWhenFeaturePresentButDisabled` builds that state directly at creation (`delta.feature.variantShredding = supported` with the property explicitly `false`), since the `REORG` syntax is not in this repo's SQL parser, and asserts that neither a DSv2 nor a V1 write shreds. Shredding is judged from the files in the table's *current snapshot* (`input_file_name()`) rather than from a directory listing, so files that are on disk but no longer referenced do not count. **The option merge** is asserted directly rather than through the resulting file layout, because the layout cannot show which of two colliding spellings won -- `DeltaV2WriteContextTest#mergeVariantShreddingOptionsOverridesCallerSpellings` passes two case-variant spellings and asserts exactly one key survives, carrying the table-derived value. `DeltaV2WriteContextTest#variantLayoutFollowsPropertyMatchesShreddingEligibility` checks that the layout is marked property-sensitive only for schemas the writer can actually shred. **The streaming guard** -- tests in `DeltaV2StreamingWriteTest`, on a real variant table: - `testCommit_failsWhenShreddingPropertyChangeAffectsLayout`: property on at stream start, set to `false`, the commit fails naming the property and appends nothing. - `testCommit_ignoresShreddingPropertyChangeThatCannotAffectLayout`: the same toggle on a table with no variant column, where the layout cannot change; the epoch commits normally. - `testCommit_replayOfCommittedEpochIsIdempotentAfterPropertyChange`: commit an epoch, change the property, replay the same epoch with freshly written files; the replay adds no data and no new table version. - `testCommit_failsWhenShreddingPropertyEnabledMidStream`: the opposite direction, enabling the property mid-stream on a table that already supports the feature; the commit fails. - `testCommit_streamingWriteShredsWhenPropertyEnabled` / `testCommit_streamingWriteDoesNotShredWhenPropertyDisabled`: the committed file's Parquet footer has (or lacks) a `typed_value` child, so the property reaches the executor-side writer. **End to end through the streaming sink** -- `DeltaSinkSuite`'s `streaming write shreds variant columns only when the table enables it` runs a streaming query against tables with the property on and off, checks the file layout, and reads the data back. `DeltaV2SinkSuite` lists it as passing, so it also runs on the V2 sink. With the physical-name lookup it also passes in `DeltaSinkIdColumnMappingSuite` and `DeltaSinkNameColumnMappingSuite` when run there. **Version gating.** The cases that require shredding to actually happen are gated on the production shim yielding an inference option, which it does not on Spark 4.0; the property-`false` coverage still runs there. The gate is taken from the shim rather than from a version string so it cannot drift from what the shim does. ## Does this PR introduce _any_ user-facing changes? Yes, within `master` -- the DSv2 write path now honors `delta.enableVariantShredding`. Previously, with shredded writes allowed in the session, a DSv2 write shredded variant columns regardless of the table property: a table with `delta.enableVariantShredding = false` (or with the property absent after `REORG TABLE ... APPLY (UNSHRED VARIANT)`) still received shredded data files, and its protocol did not advertise the `variantShredding` reader-writer feature those files require. After this change the DSv2 write path shreds exactly when the table opted in, matching the V1 write path. Additionally, a streaming DSv2 write now fails the epoch, asking for a query restart, if the property changes mid-query in a way that would change the file layout. This pull request and its description were written by Isaac. Signed-off-by: Liang-Chi Hsieh <viirya@gmail.com> Co-authored-by: Isaac <no-reply@databricks.com>
…a-io#7793) ## Description Carry the shared `org.apache.spark.sql.delta.Snapshot` facade through the Delta V2 read path instead of eagerly unwrapping it to the Kernel snapshot. `DeltaV2ScanBuilder` now passes the `Snapshot` facade directly into `DeltaV2Scan`, and `DeltaV2Batch` holds the facade as a field. Path, version, and cheap file-size statistics are read through the facade's public API (`path()`, `version()`, `dataPath()`, `sizeInBytesIfKnown`); the underlying Kernel snapshot is unwrapped only inside the small helpers that still require Kernel-specific APIs (deletion-vector support checks and CRC-based table size). This is a representation-only refactor — scan planning and statistics behavior are unchanged. ## How was this patch tested? Existing unit tests covering the Delta V2 scan, batch, and snapshot paths exercise this change; it is a pure refactor with no behavioral difference, so no new tests are added. ## Does this PR introduce _any_ user-facing changes? No.
## Description Now that Delta 4.4.0 has been released and branch-4.4 has been cut, bump the master development version to 4.5.0-SNAPSHOT. This also regenerates python/delta/version.py so the Python package version remains consistent with version.sbt. Fixes delta-io#7652 ## How was this patch tested? - build/sbt sparkV1/generatePythonVersion - Focused version consistency check (version.sbt matches python/delta/version.py) - git diff --check Co-authored-by: Vishnu Chandrashekhar <seewishnew@users.noreply.github.com>
…ion (delta-io#7795) ## What changes were proposed in this pull request? Reorder commit preparation so resolved file actions are validated and the adaptive metadata tree is written before UniForm/Iceberg metadata generation. Iceberg metadata generation does not change the resolved file actions consumed by the tree writer. The final action assembly still starts from the transaction state returned by metadata generation. ## How was this patch tested? Existing focused tests for metadata publication, DML, conflict retries, row-ID allocation, and snapshot retention passed. No new tests are included in this patch. An independent-reader end-to-end test was not run. Co-authored-by: Kaiqi Jin <22572475+KaiqiJinWow@users.noreply.github.com>
…to back-reference check (delta-io#7797) #### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description Two small fixes in the Adaptive Metadata Tree (AMT) write path: 1. `AMTWriteHelper.writeLeavesDistributed` now converts each `AddFile` into a `DataEntry` leaf row with `mapPartitions` instead of a per-row typed `map`. The conversion and the `baseRowId` requirement are unchanged. 2. The test-only `verifyCommitBackReferences` check now receives the transaction's catalog table and passes it to `deltaLog.getChanges` when listing intermediate commits. The commit listing then resolves the table the same way the committing transaction does (e.g. for catalog-managed tables) instead of by path only. ## How was this patch tested? Existing unit tests. ## Does this PR introduce _any_ user-facing changes? No.
…o#7746) ## Description Extends the row-index filter provider path so that manifest-backed row-index filters can be resolved portably during a scan. Manifest filters do not carry a `DeletionVectorDescriptor` (or a roaring bitmap) the way deletion-vector filters do, so this change threads the manifest filter information through `RowIndexFilterProvider` and the Parquet reader path and resolves the correct filter for manifest-backed entries. Key changes: - Add manifest support to `RowIndexFilterProvider` and construct the appropriate filter when a manifest bitmap is present rather than a deletion vector. - Thread the provider through `DeltaLogFileIndex`, the manifest-aware Parquet file format, and the AMT checkpoint provider. ## How was this patch tested? Added/updated coverage in `DeltaParquetFileFormatSuite`, `RowIndexMarkingFiltersSuite`, and `AMTInheritanceReadSuite`. ## Does this PR introduce _any_ user-facing changes? No. Co-authored-by: Jinhua Song <jinhua.song@databricks.com>
<!-- Thanks for sending a pull request! Here are some tips for you: 1. If this is your first time, please read our contributor guidelines: https://github.com/delta-io/delta/blob/master/CONTRIBUTING.md 2. If the PR is unfinished, add '[WIP]' in your PR title, e.g., '[WIP] Your PR title ...'. 3. Be sure to keep the PR description updated to reflect all changes. 4. Please write your PR title to summarize what this PR proposes. 5. If possible, provide a concise example to reproduce the issue for a faster review. 6. If applicable, include the corresponding issue number in the PR title and link it in the body. --> #### Which Delta project/connector is this regarding? <!-- Please add the component selected below to the beginning of the pull request title For example: [Spark] Title of my pull request --> - [ ] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [x] Other (Iceberg) ## Description This bumps the Iceberg version from 1.11.0 to 1.12.0 <!-- - Describe what this PR changes. - Describe why we need the change. If this PR resolves an issue be sure to include "Resolves #XXX" to correctly link and close the issue upon merge. --> ## How was this patch tested? existing tests <!-- If tests were added, say they were added here. Please make sure to test the changes thoroughly including negative and positive cases if possible. If the changes were tested in any way other than unit tests, please clarify how you tested step by step (ideally copy and paste-able, so that other reviewers can test and check, and descendants can verify in the future). If the changes were not tested, please explain why. --> ## Does this PR introduce _any_ user-facing changes? no <!-- If yes, please clarify the previous behavior and the change this PR proposes - provide the console output, description and/or an example to show the behavior difference if possible. If possible, please also clarify if this is a user-facing change compared to the released Delta Lake versions or within the unreleased branches such as master. If no, write 'No'. -->
…ath (delta-io#7774) ## Description File selection in the Delta V2 scan goes through `snapshot.filesForScan` and returns a `DeltaScan`. Previously the scan eagerly built physical `PartitionedFile`s — resolving paths and serializing deletion-vector metadata — and threaded them through `DeltaV2Batch`, so statistics and validation consumers paid for physical materialization they never used. This change removes the redundant `DeltaScanFile` wrapper along with the `ensurePlanned` / `getSelectedFiles` coupling. `DeltaV2Batch` now retains the selected `DeltaScan` and defers physical `PartitionedFile` construction to `PartitionUtils.planInputPartitions`, which materializes partition rows only when a batch actually needs physical input partitions. Aggregate scan statistics are read directly from the `DeltaScan`. Statistics and runtime-filtering behavior are unchanged, and the metadata aggregation over selected `AddFile`s is preserved. ## How was this patch tested? Existing `DeltaV2ScanTest` coverage — runtime-filtering, row-count, copy-state, scan-report, and failed-materialization scenarios — passes unchanged. The scan-statistics assertions were updated to read aggregates from the `DeltaScan` directly. ## Does this PR introduce _any_ user-facing changes? No.
…delta-io#7803) #### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description Removes superfluous `SET DeltaSQLConf.DELTA_MERGE_PRESERVE_NULL_SOURCE_STRUCTS=true` from test setup as that conf already defaults to `true`. ## How was this patch tested? Test-only change, CI. ## Does this PR introduce _any_ user-facing changes? No Co-authored-by: Matthis Gördel <matthis.goerdel@databricks.com>
## Description Preserve row tracking domain metadata when a Delta V2 optimistic transaction commits through Kernel. This keeps the row tracking high-water mark available during conflict resolution, including transactions that do not add files. ## Tests - Added coverage for row tracking domain metadata conversion and Delta V2 transaction commits. - Added concurrent commit coverage for row tracking high-water mark handling.
#### Which Delta project/connector is this regarding? - [ ] Spark - [ ] Standalone - [ ] Flink - [x] Kernel - [ ] Other (fill in here) ## Description `UCCatalogManagedClient` throws `IllegalArgumentException` when a requested snapshot version or commit-range boundary exceeds UC's latest ratified version. CDF streaming with `startingVersion=latest` requests `latest + 1`, so consumers need to distinguish an unavailable version from malformed arguments and other failures. Add `VersionToLoadAfterLatestCommitException` and use it for snapshot versions and explicit start/end versions above the catalog maximum, preserving the diagnostic messages. Invalid argument combinations, negative snapshot versions, and reversed ranges retain `IllegalArgumentException`. This supplies the precise UC failure type needed by delta-io#7754. Tests cover both the Kernel client and Spark snapshot-manager boundaries. This change shares the exception definition proposed in delta-io#7055, with contributor credit preserved, but does not include or depend on that PR's path-based or truncated-history changes. Either PR can merge first; the other can then reuse the shared class. Related to delta-io#6745. ## How was this patch tested? - Regression tests cover snapshot, range-start, and range-end versions at `latest + 1` and farther ahead, including exact exception type and diagnostic messages in the Kernel tests. The new future-version cases failed with the old `IllegalArgumentException` before their production changes. - All 52 tests passed across `UCCatalogManagedClientSuite`, `UCCatalogManagedClientCommitRangeSuite`, and `UcLoadSnapshotTelemetrySuite`. - All 10 `UCManagedTableSnapshotManagerSuite` tests passed on Spark 4.0, including propagation of the typed exception for future snapshot and range-boundary versions. - Earlier validation passed all 51 `SnapshotManagerSuite` tests without the path-based changes from delta-io#7055; that code is unchanged by the commit-range update. - Build-enforced Java/Scala formatting and Checkstyle, plus `git diff --check`, passed. Commands (Java 17): ```sh ./build/sbt 'kernelUnityCatalog/testOnly io.delta.kernel.unitycatalog.UCCatalogManagedClientSuite io.delta.kernel.unitycatalog.UCCatalogManagedClientCommitRangeSuite io.delta.kernel.unitycatalog.UcLoadSnapshotTelemetrySuite' ./build/sbt 'kernelApi/testOnly io.delta.kernel.internal.SnapshotManagerSuite' ./build/sbt -DsparkVersion=4.0 'sparkV2/testOnly io.delta.spark.internal.v2.snapshot.unitycatalog.UCManagedTableSnapshotManagerSuite' ``` ## Does this PR introduce _any_ user-facing changes? Yes. A UC snapshot request or explicit commit-range start/end version above the latest ratified version now throws `VersionToLoadAfterLatestCommitException` (a `KernelException`) instead of `IllegalArgumentException`. The messages are unchanged. Callers handling this condition should catch the specific type. --------- Signed-off-by: Vishnu Chandrashekhar <vishnu.c.9965@gmail.com> Co-authored-by: GabrielBBaldez <130607246+GabrielBBaldez@users.noreply.github.com>
…elta-io#7798) #### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description `DeltaV2MicroBatchStream` currently has two initial-snapshot implementations: a driver-side `InitialSnapshotCache` that materializes the full file list, and a flag-gated DataFrame-based snapshot cache (`streaming.distributedInitialSnapshot`). The V1 `DeltaSource` already handles initial-snapshot ordering and lifecycle through `DeltaSourceSnapshot`. This PR makes `DeltaV2MicroBatchStream` use `DeltaSourceSnapshot` exclusively: - Lazily create one `DeltaSourceSnapshot` for the initial snapshot version and reuse it while the stream reads that snapshot. - Follow the same lifecycle as the V1 source: release the snapshot when planning moves past the initial snapshot and when the stream stops. - Use the same snapshot path for initial change-data-feed reads. - Remove the driver-side `InitialSnapshotCache`, the DataFrame snapshot cache routing, and their helpers from `DeltaV2MicroBatchStream`. - Add `KernelActionUtils.toKernelScanAddFile`, which converts a V1 `AddFile` to the Kernel scan `AddFile`. It keeps deletion vectors, row-tracking fields, tags, and stats. - Remove the initial-snapshot file-count limit (`streaming.initialSnapshotMaxFiles`) and its `DELTA_STREAMING_INITIAL_SNAPSHOT_TOO_LARGE` error. `DeltaSourceSnapshot` does not collect the whole file list on the driver, so the limit no longer applies. ## How was this patch tested? - New `KernelActionUtilsSuite` test round-trips a V1 `AddFile` through `toKernelScanAddFile` and `addFileFromKernel`, covering partition values, tags, deletion vector, row-tracking fields, and stats. - Existing V2 streaming suites cover initial snapshot and initial CDF reads. ## Does this PR introduce _any_ user-facing changes? No. The removed config and error class were internal and applied only to the V2 streaming source.
…elta-io#7796) #### Which Delta project/connector is this PR targeting? Spark #### Description The convenience accessors on `SnapshotStateManager` (`sizeInBytes`, `numOfFiles`, `domainMetadata`, the deletion-vector metrics, and the `*IfKnown` getters) each carried their own copy of the "read from the checksum, otherwise fall back to state reconstruction" logic. This PR routes them all through a single `checksumOptForState`, gated on `fastQueryPath.enabled`, and records which accessor forced a reconstruction via `recordComputedStateAccess`. Two `Snapshot` subclasses need to keep their prior values now that the accessors live on the base trait: - `DummySnapshot` seeds the incremental checksum of the first commit, which may enable deletion vectors. It now serves its already-computed initial-state DV metrics directly instead of gating them on the pre-commit protocol's DV readability -- otherwise the first checksum stores `None` and the next commit reconstructs state to fill it. - `DeltaV2Snapshot` serves domain metadata from Kernel without state reconstruction, so it reports `domainMetadatasIfKnown` as always known. See delta-io#7507 for the readable-vs-writable deletion-vector metric predicate divergence referenced in the comments. #### How was this patch tested? Existing `ChecksumSuite`, `DeltaV2SnapshotSuite`, and optimistic-transaction suites, plus a new `ChecksumSuite` case asserting that disabling `fastQueryPath.enabled` bypasses the checksum-backed snapshot state. #### Does this PR introduce _any_ user-facing changes? No.
…delta-io#7804) #### Which Delta project/connector is this regarding? - [x] Spark ## Description Preserve literal string operands when translating overwrite filters. Translating `startswith`, `endswith`, and `contains` into `LIKE` patterns incorrectly treats percent signs, underscores, and backslashes as wildcards or escapes. This can remove unrelated rows or reject valid replacement rows. Add `replaceWhere.literalStringPredicates.enabled`, an internal session configuration that defaults to `true`. When enabled, use Catalyst `StartsWith`, `EndsWith`, and `Contains` expressions. When disabled, retain the unchanged legacy `LIKE` translation. Add 42 regression cases: seven scenarios with both flag values through SQL `INSERT ... REPLACE WHERE`, `DataFrameWriterV2.overwrite`, and the `replaceWhere` writer option. Cover literal percent, underscore, and backslash operands; empty strings and nulls; mixed incoming rows and unchanged transaction versions after rejected writes; brackets and braces; and explicit `LIKE` wildcard semantics. ## How was this patch tested? - Added the regression matrix to `DeltaDataFrameWriterV2Suite`. - Verified the complete test block matches the intended helpers and all seven scenarios. - `git diff --check` passed. - Ran `spark/testOnly org.apache.spark.sql.delta.DeltaDataFrameWriterV2Suite -- -z "literal string predicates"` with Java 17 and Spark 4.1: all 42 tests passed. ## Does this PR introduce _any_ user-facing changes? Yes. Overwrite filters preserve literal percent signs, underscores, and backslashes by default. Setting the new configuration to `false` restores the previous behavior. The `replaceWhere` writer option retains its existing literal matching behavior with either flag value.
…-io#7809) #### Which Delta project/connector is this regarding? - [x] Spark ## Description Extract `filterFileList` and `rewritePartitionFilters` into a top-level `DeltaLogUtils` object in `DeltaLog.scala`. Keep the existing methods in `DeltaLog` as forwarding methods, preserving their signatures and default arguments. Update `DeltaSourceSnapshot` to use `DeltaLogUtils` so snapshot readers can use the partition-filtering helpers without depending on the `DeltaLog` companion. The filtering implementations are unchanged; no new source file is needed. ## How was this patch tested? - `git diff --check` passed. - No new test suite was added because this is a behavior-preserving refactor. The existing data-skipping and Delta source tests have not been run locally; automated test results are pending CI. ## Does this PR introduce _any_ user-facing changes? No.
#### Which Delta project/connector is this regarding? - [x] Spark - [ ] Standalone - [ ] Flink - [ ] Kernel - [ ] Other (fill in here) ## Description The AMT manifest writers (the distributed leaf writer and the single-file root writer) built a plain `ParquetFileFormat` and then patched the job with `Checkpoints.configureIcebergManifestParquetWrite` to produce Iceberg-V4 manifests (int64 `TIMESTAMP_MICROS` timestamps, plus `DeltaParquetWriteSupport` so nested list/map field ids are written). This PR moves that write configuration into `AMTParquetFileFormat`, so the format owns both its read and write behavior, in the same way `DeltaParquetFileFormat.prepareWrite` applies its Iceberg write settings. - `AMTParquetFileFormat`: override `prepareWrite` to always set `TIMESTAMP_MICROS` and `DeltaParquetWriteSupport`, since AMT manifests are always Iceberg-V4 manifests. - `AMTWriteHelper`: the leaf writer and the root writer now use `AMTParquetFileFormat`. - `Checkpoints.writeAtomicCheckpointParquetFile`: replace the `writeAsIcebergManifest: Boolean` flag with a `format: ParquetFileFormat` parameter. It defaults to a plain `ParquetFileFormat`, so classic checkpoint writes are unchanged. The now-unused `configureIcebergManifestParquetWrite` is removed. - Tests: update the two `writeAtomicCheckpointParquetFile` call sites. The written manifests use the same Parquet settings as before. ## How was this patch tested? Existing tests, which already check nested field ids and `TIMESTAMP_MICROS` on written AMT manifests and round-trip the manifests: `AMTCheckpointWriteSuite`, `AMTSnapshotSuite`, and the other AMT write/read suites. ## Does this PR introduce _any_ user-facing changes? No.
Adds a metadata-only command that removes the partitioning of a Delta table without rewriting data files: ALTER TABLE t REPLACE PARTITIONED BY WITH CLUSTER BY NONE ALTER TABLE t REPLACE PARTITIONED BY WITH CLUSTER BY (c1, c2) The command verifies that every data file physically stores the partition columns (written when delta.writePartitionColumnsToParquet is enabled), then re-adds all files in a single commit with empty partitionValues, dataChange=false and stats extended with the former partition values, so data skipping keeps working. With NONE the table becomes a plain unpartitioned table; with columns it becomes a clustered table. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
1 task done
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Which Delta project/connector is this regarding?
Description
Adds a metadata-only way to remove the partitioning of a Delta table, optionally replacing it with liquid clustering:
Today the only way to drop partitioning is to rewrite the whole table (CTAS / REPLACE TABLE). For large tables that are over-partitioned (e.g. partitioned by a high-cardinality ingestion id), that is very expensive. However, since
delta.writePartitionColumnsToParquet(delta-io#6525, default on) /materializePartitionColumns(delta-io#5636), data files already contain the partition column values physically, so the partitioning can be dropped by only rewriting the log.The syntax pairs
REPLACE PARTITIONED BYwith the existingCLUSTER BY (...) | CLUSTER BY NONEclause rather than introducing a new standalone command, so it composes with the existing clustering DDL.Semantics
DELTA_REPLACE_PARTITIONED_BY_PARTITION_COLUMNS_NOT_MATERIALIZED, reporting the count and example paths. (Internal confspark.databricks.delta.alterTable.replacePartitionedBy.verifyMaterializedPartitionColumnscan disable this check.)commitLarge, like RESTORE/CLONE, so any concurrent commit aborts it):metaDatawithpartitionColumns = [](+ clustering table properties when columns are given);partitionValues = {},dataChange = false, and stats extended withminValues/maxValues/nullCountfor the former partition columns (string/integral/decimal/date; null count only for other types), so data skipping on former partition columns continues to work;CLUSTER BY NONEdoes not change the protocol.CLUSTER BY (cols)enables the clustering table feature, same asALTER TABLE ... CLUSTER BY.baseRowId,defaultRowCommitVersion) and column mapping are preserved. Time travel to pre-conversion versions keeps working.Follow-ups / limitations
nullCountstats only.dataChange=falsecommit (no new data).How was this patch tested?
New
ReplacePartitionedByWithClusterBySuite(parser, no-rewrite NONE conversion, stats + data skipping, writes/UPDATE/DELETE/MERGE/OPTIMIZE after conversion, conversion to clustered, catalog tables with null partitions, column mapping, DV + row tracking, error cases). Also ranDeltaSqlParserSuite,DeltaErrorsSuite,DeltaThrowableSuiteand scalastyle.Does this PR introduce any user-facing changes?
Yes, new SQL syntax
ALTER TABLE ... REPLACE PARTITIONED BY WITH CLUSTER BY (...) | NONE.Follow-up (stacked): #18 rewrites data files that do not store the partition columns instead of failing.