zhuxiangyi opened a new pull request, #10098:
URL: https://github.com/apache/paimon/pull/10098

   ### Purpose
   
   Feature: convert an existing append table to a Data Evolution table **in 
place**, without rewriting any data file.
   
   #### Why
   
   `row-tracking.enabled` and `data-evolution.enabled` are immutable (`ALTER 
TABLE ... SET` is refused as soon as the table has a snapshot). Today the only 
way to get Data Evolution on a table that already holds data is to create a new 
table and `INSERT OVERWRITE` everything into it: a full rewrite, a new table 
location, and all history/tags/branches lost. Data Evolution is exactly the 
feature people want on *big* tables (partial-column updates, sub-field updates, 
blob columns), so this is the tables it is most painful for.
   
   Just flipping the two options is not enough, and is why they were made 
immutable: a Data Evolution table addresses rows by `firstRowId + position`, 
and the existing files have no `firstRowId`. Every code path that assumes the 
id (`nonNullFirstRowId`, the split generator, compaction, `MERGE INTO` on 
`_ROW_ID`) would fail or silently skip rows. Enabling on an existing table 
therefore needs three things: row ids for the existing files, a safe switch of 
the schema, and reads that tolerate the window in between.
   
   #### Usage
   
   ```sql
   -- Spark SQL
   CALL sys.enable_data_evolution(table => 'default.target_table', dry_run => 
true);  -- report only
   CALL sys.enable_data_evolution(table => 'default.target_table');
   
   -- Flink SQL
   CALL sys.enable_data_evolution(`table` => 'default.target_table', dry_run => 
true);
   CALL sys.enable_data_evolution('default.target_table');
   ```
   
   ```bash
   # Flink action jar
   <FLINK_HOME>/bin/flink run paimon-flink-action-<version>.jar 
enable_data_evolution \
       --warehouse <warehouse-path> --database default --table target_table 
[--dry_run true] \
       [--catalog_conf <key>=<value> ...]
   ```
   
   The call returns one row:
   
   | Result | Meaning |
   |---|---|
   | `Dry run. Enabling data evolution on table 'db.t' would do: ...` | the 
same report, nothing is changed |
   | `Success. Enabled data evolution on table 'db.t': schema 0 -> 1, snapshot 
2 -> 3, N file(s) with M row(s) assigned row ids, nextRowId=M.` | converted: 
one metadata-only `OVERWRITE` snapshot plus one new schema |
   | `Skipped. Table 'db.t' was not changed: data evolution is already 
enabled.` | idempotent |
   
   **What the table looks like afterwards**
   
   - Exactly two options are added to a new schema, everything else in the 
schema (fields, partition keys, comment, all other options) is carried over 
unchanged:
     ```
     row-tracking.enabled   = true
     data-evolution.enabled = true
     ```
     The result is the same schema a table gets from `CREATE TABLE ... 
TBLPROPERTIES ('row-tracking.enabled'='true', 
'data-evolution.enabled'='true')`. Nothing else is switched on: `bucket`, 
deletion vectors and `data-evolution.nested-field.enabled` stay as they were 
and can be changed afterwards with `ALTER TABLE` as on any Data Evolution table.
   - Every existing data file keeps its name, path (external paths included), 
statistics, embedded file index and deletion vector; it only gets a 
`firstRowId`. Row ids are contiguous per partition, in commit order, so 
`sys.reassign_row_id` has nothing to do afterwards. New writes continue at 
`nextRowId`.
   - The two options are written **after** every existing file has its row id 
(manifest rewrite first, schema second), so no writer can observe a schema that 
says "data evolution" over files that have no row id.
   
   **Requirements**: append table without primary key, `bucket = -1` (the 
default for an append table), `clustering.incremental` off. Tables of a REST 
catalog are rejected in this version.
   
   **What to expect around the conversion** (also in the docs):
   
   - A writer that loaded the table before the conversion, typically a running 
Flink streaming job, is refused on its next commit with a message asking to 
restart it; after the restart it assigns row ids as usual. Run the conversion 
at a write-traffic low.
   - Snapshots and tags from before the conversion stay readable; their files 
have no row id, so `_ROW_ID` reads as `NULL`. Rolling back to such a snapshot 
is refused.
   - Streaming reads that were started before the conversion continue past the 
conversion snapshot without replaying the table.
   
   #### What
   
   One PR, four commits, each self-contained:
   
   **1. `[core] Guard commits and rollbacks across a row-tracking switch`**
   
   Two guards that make the switch safe against concurrent writers and against 
history:
   
   - `FileStoreCommitImpl` refuses a commit from a writer that loaded the table 
*without* row tracking when the latest schema now has it enabled 
(`checkRowTrackingNotEnabledAfterLoad`). Such a writer would commit files 
without a row id into a row-tracking table. The message asks to restart the 
writer. Writers that loaded the table with row tracking, and ordinary schema 
changes (add column, set other options), are unaffected.
   - `AbstractFileStoreTable.rollbackTo` (snapshot and tag) and 
`rollbackToAsLatest` refuse to roll back to a snapshot whose schema has no row 
tracking while the current schema has it: that snapshot's files would become 
unaddressable.
   - `replaceManifestList` gets an overload that takes the schema id and 
`nextRowId` explicitly (used by commit 3).
   
   **2. `[core] Read data-evolution files that have no first row id as plain 
files`**
   
   A Data Evolution table can now hold files without `firstRowId` (files 
written before the switch, snapshots and tags from before the conversion, or a 
table where only the schema was flipped). They are treated as complete-row 
files:
   
   - `DataEvolutionUtils.splitByRowIdPresence` separates them; 
`DataEvolutionSplitGenerator` bin-packs them into raw-convertible splits 
*after* the row-id-range groups, so a file with an id and a file without never 
share a split.
   - `DataEvolutionSplitReadProvider` does not match such splits; 
`AppendOnlyFileStoreTable.newRead` always registers the plain append raw-file 
provider after the Data Evolution one, so they are read as ordinary append 
files with `_ROW_ID = NULL`.
   - `DataEvolutionFileStoreScan` groups them as singletons and applies file 
statistics to them individually (predicate pushdown keeps working on them).
   - `DataFileMeta.nonNullFirstRowId` and `DataEvolutionCompactRangePlanner` 
name the file in their message and point at `sys.enable_data_evolution`, 
instead of a bare `firstRowId must not be null`.
   
   **3. `[core] Add DataEvolutionEnabler to convert an append table in place`**
   
   - `SchemaChange.enableDataEvolution()`: a dedicated schema change (JSON 
action `enableDataEvolution`, also added to the REST OpenAPI contract) that 
sets both options in a new schema. `SchemaManagerUtils` lets it through the 
immutable-option check; `validateTableSchema` still enforces every other Data 
Evolution constraint (no primary key, `bucket = -1`, no 
`clustering.incremental`). `ALTER TABLE SET ('row-tracking.enabled' = 'true')` 
stays refused.
   - `DataEvolutionEnabler(catalog, identifier).run(dryRun)`:
     1. **validate**: append table, candidate schema passes 
`validateTableSchema`; REST catalog tables are rejected in this version (the 
manifest rewrite is a direct `replaceManifestList` on the file store, which the 
REST protocol does not expose yet).
     2. **assign row ids**: metadata only. Live `ADD` entries without id are 
ordered per partition in manifest/commit order and given contiguous ids 
starting at `nextRowId` (or 0); only the manifests that hold converted files 
are rewritten, and the snapshot is replaced with `replaceManifestList` carrying 
the new `nextRowId`. On a conflict (another commit landed) it re-plans from the 
new latest snapshot, bounded by the table's commit retry settings.
     3. **flip the schema**: `catalog.alterTable(identifier, 
enableDataEvolution())`.
     4. **repair**: a writer that loaded the table between 2 and 3 may still 
commit files without id (the guard from commit 1 only sees the schema after 3). 
Up to five repair rounds assign ids to anything that slipped in.
     - Dry run reports the plan without changing anything. Idempotent: a 
converted table reports `Skipped.`
     - Row ids are contiguous per partition in commit order, so a converted 
table needs no `reassign_row_id` afterwards.
   
   **4. `[spark][flink] Add the enable_data_evolution procedure and action`**
   
   - Spark `CALL sys.enable_data_evolution(table => 't' [, dry_run => true])`, 
Flink `CALL sys.enable_data_evolution('db.t' [, dry_run])` and the 
`enable_data_evolution` action jar.
   - Spark `MERGE INTO` on a Data Evolution table refuses a target that still 
has files without row id, with a message pointing at the procedure, instead of 
silently skipping those rows.
   - Spark `BaseScan` reads `_ROW_ID` / `_SEQUENCE_NUMBER` as nullable so such 
files read as `NULL` (a non-nullable read type turned it into `0`).
   - Two `copy(...)` → `copyWithoutTimeTravel(...)` changes in 
`MergeIntoPaimonDataEvolutionTable` (Spark 3 and the Spark 4.0 copy) and 
`DataEvolutionPaimonWriter`. This fixes a **pre-existing bug on native Data 
Evolution tables too**: the merge pins its scan with `scan.snapshot-id`, and 
`copy` re-applies that time travel and adopts the schema the snapshot was 
committed with. After a schema-only change (`ALTER TABLE` creates no snapshot: 
enabling `data-evolution.nested-field.enabled`, adding a column) the files 
written by the merge were stamped with the *old* schema id, and reading them 
back returned the stale sub-field value or failed in `TableSchema.project`. 
`copyWithoutTimeTravel` keeps the current schema while still pinning the 
snapshot.
   
   #### Benefit
   
   - A table of any size gets Data Evolution with one metadata-only call: no 
data rewrite, same location, history/tags/branches kept, old snapshots stay 
readable.
   - The switch is safe: stale writers are refused instead of corrupting row 
ids, rollbacks across the boundary are refused, files that slipped in are 
repaired, and every row-id-dependent path reports the offending file and the 
fix.
   - Sub-field / added-column merges after `ALTER TABLE` on existing Data 
Evolution tables now write files with the right schema id.
   
   ### Tests
   
   **paimon-core**
   
   - `RowTrackingEnableGuardTest` (6): stale writer refused after the switch, 
stale writer still commits after an ordinary schema change, writer on a 
row-tracking table never refused, rollback across the boundary refused 
(snapshot and tag), plain append table unaffected, `replaceManifestList` with 
schema id.
   - `DataEvolutionFilesWithoutRowIdTest` (5): full read / projection / 
predicate on files without id (`_ROW_ID` NULL, statistics prune), mixed files 
with and without id never share a split, time travel to a pre-conversion 
snapshot, compaction refuses with a message naming the file, 
`nonNullFirstRowId` message.
   - `DataEvolutionEnablerTest` (17): empty table only flips the schema; 
contiguous ids in sequence order; per-partition ranges; only manifests holding 
converted files are rewritten; every other file attribute (external path, 
embedded file index, stats, ...) is kept byte-for-byte; dry run changes 
nothing; converted table supports partial-column write and Data Evolution 
compaction; converts one branch only; rejects primary-key / bucketed / 
incremental-clustering tables and REST catalog tables (mock 
`RESTCatalogServer`, dry run included); `ALTER TABLE` still cannot enable row 
tracking; concurrent append and concurrent compaction before the row-id commit 
are assigned on retry; a stale writer committing between the row-id commit and 
the schema change is repaired; gives up when snapshots keep moving.
   - `SchemaManagerTest#testEnableDataEvolutionChange`, `RESTApiJsonTest` 
round-trip of the new action.
   - All existing `DataEvolution*` suites, `AppendOnlySimpleTableTest`, 
`BlobTableTest` pass.
   
   **paimon-spark** (Spark 3.5)
   
   - `EnableDataEvolutionProcedureTest` (7): basic conversion + dry run + 
idempotence + `MERGE INTO` / `UPDATE` / `INSERT` afterwards (partial-column 
files, row ids continue after the assigned range); partition-contiguous ids and 
`reassign_row_id` reports `Skipped.`; deletion vectors kept and `DELETE` / 
merge-delete keep working; sub-field data evolution on a converted table; 
pre-conversion snapshot and tag readable with `_ROW_ID` NULL; refusals (primary 
key, bucketed, `ALTER TABLE`); `MERGE INTO` refuses files without id, then the 
procedure repairs the table.
   - `RowTrackingTestBase`: new `merge after a schema-only change writes files 
with the current schema` (fails on master: the written file carries schema id 
0). `RowTrackingTest`, `BlobUpdateTest`, `DataEvolutionDeletionTest`, 
`UpdateTableTest`, `ReassignRowIdProcedureTest`, `NestedSubfieldMergeIntoTest`, 
`DataEvolutionUpdateSnapshotTest` pass.
   
   **paimon-flink**
   
   - `EnableDataEvolutionProcedureITCase` (5): procedure (ids read through 
`t$row_tracking`), refusal of a primary-key table, action jar, writer loaded 
before the conversion is refused, streaming reads started before the conversion 
continue past the conversion snapshot without replaying the table (default and 
`streaming-read-append-overwrite`).
   
   `spotless:check` and `checkstyle:check` pass on paimon-api, paimon-core, 
paimon-spark-common, paimon-spark-4.0, paimon-spark-ut and paimon-flink-common.
   
   ### API and Format
   
   - New `SchemaChange.enableDataEvolution()` / JSON action 
`enableDataEvolution` (added to the REST OpenAPI contract). The REST catalog 
does not support the conversion itself yet; the procedure reports this.
   - Data Evolution tables may now contain data files with `firstRowId = null`; 
they are read as complete-row files with `_ROW_ID = NULL`. The file format 
itself is unchanged.
   - New Spark procedure `sys.enable_data_evolution`, Flink procedure 
`sys.enable_data_evolution` and action `enable_data_evolution`.
   
   ### Documentation
   
   - `multimodal-table/data-evolution.mdx`: new section "Enable on an Existing 
Append Table" (what the procedure does, what to expect around the conversion: 
stale writers, old snapshots, rollback, nested-field option).
   - `append-table/row-tracking.md`, `spark/procedures.md`, 
`spark/procedures/indexes.md`, `flink/procedures.md`, 
`flink/procedures/compaction.md`, `flink/action-jars.md`: procedure and action 
entries.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to