[core][spark][flink] Enable data evolution on an existing append table in place - #10098
zhuxiangyi wants to merge 21 commits into
Conversation
23b7b30 to
383e19c
Compare
Row tracking derives every row id from the first row id that the commit assigns to a new file, so once a table with snapshots enables it, two things must not happen: a writer that loaded the table before the switch keeps committing files that never get a first row id, and the table is rolled back to a snapshot whose files have none. Both would leave the table unreadable as a data-evolution table. Refuse a commit whose writer has row tracking disabled while the latest schema has it enabled, so a stale writer restarts and reloads the table. The check reuses the latest-schema lookup the commit already does; it is hoisted before the manifests are written so a refused commit leaves nothing behind. Refuse a rollback to a snapshot committed under a schema without row tracking while the current schema has it, on both the table and the commit rollback paths. Also let replaceManifestList reference an explicit schema id, so a metadata-only commit can be the first snapshot of a new schema.
…iles A data-evolution table that enabled row tracking after it already had snapshots holds files without a first row id: everything committed before the switch, until sys.enable_data_evolution assigns ids to them, and every snapshot, tag or time-travel target from before the switch forever. Such a file is a complete-row file, so let it be read as one instead of failing in nonNullFirstRowId. The scan groups those files as singletons next to the row-id-range groups and prunes them with plain per-file statistics, the split generator packs them like an append table does, and the data-evolution split read declines a split containing them so that the raw file read, now the fallback provider of a data-evolution table, serves it with a NULL _ROW_ID. Operations that need row ids keep refusing such files, with a message that names the file and points at the procedure: compaction planning and every remaining nonNullFirstRowId caller.
Converting an existing append table to a data-evolution table needs two things that ALTER TABLE cannot give: a first row id for every data file that is already in the table, and the two immutable options flipped. DataEvolutionEnabler does both without rewriting any data file. It first rewrites the manifests of the latest snapshot so that every live file without a first row id gets one, contiguous per partition and in commit order, and commits them as a metadata-only snapshot; a concurrent commit makes it plan again on the new latest snapshot. It then issues the new SchemaChange.EnableDataEvolution through the catalog, so catalog metadata stays in sync while the row-tracking constraints are still validated, and finally assigns ids to whatever a writer on the previous schema committed in between. The procedure is idempotent, has a dry run, and refuses primary-key, bucketed and incrementally clustered tables as well as REST catalog tables, whose server does not know the change yet.
Expose DataEvolutionEnabler as `sys.enable_data_evolution(table [, dry_run])` in Spark and Flink, and as the `enable_data_evolution` Flink action, and document the conversion next to the data-evolution create-table section. Three engine-side changes make a converted table behave: the Spark scan reads the row-tracking fields as nullable so that files from before the conversion yield a NULL _ROW_ID instead of 0; the Spark MERGE refuses files that still have no first row id, since a partial-column write cannot address their rows; and the merge pins its scan snapshot without adopting that snapshot's schema and writes with the current one. The last change also fixes a MERGE right after any schema-only change, such as enabling data-evolution.nested-field.enabled or adding a column, which stamped the written files with the previous schema id and left a sub-field file undecodable.
A run that committed the row ids and then failed before the schema change leaves the table with row ids but without the options. Running again finishes the job, but reported nextRowId=null, because the field was taken from this run's assignment instead of the table. Take it from the planned snapshot, so only an empty table reports null. Same for the dry run. Also covers the two claims the procedure exists for: a resumed run, and a converted table building a global index, which a plain append table cannot have at all.
d354161 to
e3d2721
Compare
| * row tracking enabled would leave the table unreadable as a data-evolution table, so refuse | ||
| * it; roll the schema back first if the conversion really has to be undone. | ||
| */ | ||
| private void checkRollbackKeepsRowTracking(Snapshot target) { |
There was a problem hiding this comment.
This needs to check the latest persisted schema rather than this table object's cached coreOptions() . A FileStoreTable loaded before enable_data_evolution continues to report rowTrackingEnabled() == false ; after another handle completes conversion, that stale object can therefore call rollbackTo(1) / rollbackTo("before") and bypass this guard, making a pre-conversion snapshot current while the latest schema still has row tracking enabled. The stale-writer commit guard does not cover these direct rollback paths, and rollbackToAsLatest has the same issue. Please derive the source-side state from schemaManager().latest() in every rollback path, and add a regression test that retains the pre-conversion table handle, enables row tracking through another handle, then verifies snapshot, tag, and as-latest rollback are all refused and leave the latest snapshot unchanged.
fix: Base both rollback guards on the latest persisted schema, not the table object's cached options.
There was a problem hiding this comment.
Fixed in 35110ec. Both guards now go by the latest persisted schema instead of the options of the table object:
AbstractFileStoreTable.checkRollbackKeepsRowTracking, used byrollbackTo(long)androllbackTo(String), readsschemaManager().latest().FileStoreCommitImpl.rollbackToAsLatestreadsschemaManager.latestOrThrow(...).
Regression test: RowTrackingEnableGuardTest#testRollbackThroughTableLoadedBeforeRowTrackingIsRefused. It keeps the table object loaded before row tracking was enabled, enables row tracking through another handle, and checks that snapshot, tag and as-latest rollback are all refused, that the latest snapshot and the tag are unchanged, and that a target committed with row tracking still rolls back through the same object. Reverting either guard makes the test fail.
| beforeRowIdCommit.run(); | ||
| return commit.replaceManifestList( | ||
| assignment.snapshot, | ||
| table.schema().id(), |
There was a problem hiding this comment.
Please do not stamp this snapshot with the schema cached when the enabler initially loaded the table. Schema commits do not move the latest snapshot, so an ADD COLUMN can complete between plan() and this CAS without making replaceManifestList fail; this line will then create a newer conversion snapshot referencing the older schema, and time travel to it loses a schema change that committed before the snapshot. The problem persists across retries because the same table object is reused. Fetch the latest schema ID for each commit attempt (after beforeRowIdCommit , consistent with the normal commit path), and add a hook-based regression test that commits a schema change before the replacement snapshot and asserts that snapshot references the new schema and exposes the new column under time travel.
fix: Resolve the latest schema ID inside each assignment commit attempt, after the test/concurrency hook and immediately before replacement.
There was a problem hiding this comment.
Fixed in 37ff2bd. commitAssignment now reads the schema id with table.schemaManager().latestOrThrow(...) inside every commit attempt, after the beforeRowIdCommit hook and right before replaceManifestList, the same way the normal commit path does. The cached table.schema() is no longer used for the snapshot.
Regression test: DataEvolutionEnablerTest#testRowIdCommitReferencesTheLatestSchema. It commits ADD COLUMN c from the hook, between planning and the replacement, and asserts that the row id snapshot references the new schema, that time travel to it exposes c, and that the final schema has c with data evolution enabled. Reverting the change makes the test fail.
| * issued by the {@code sys.enable_data_evolution} procedure, which first assigns a first row id | ||
| * to every existing data file. | ||
| */ | ||
| static SchemaChange enableDataEvolution() { |
There was a problem hiding this comment.
This conversion-only action should not be exposed through generic REST schema alteration while DataEvolutionEnabler explicitly rejects REST catalogs. As added here and in the OpenAPI schema, a client can POST enableDataEvolution through alterTable ; SchemaManagerUtils then flips both options without assigning IDs to existing files. The table is left in the transient mixed state, but REST has no supported conversion procedure to repair it, so operations such as MERGE remain unusable. Please either keep this action out of public REST deserialization/OpenAPI until an atomic server-side conversion endpoint exists, or make the REST server reject it. Add a REST regression test on a nonempty append table verifying that direct alteration cannot change either option.
applies for the line 85 in this same file; also applies to docs/static/rest-catalogn-open-api.yaml and SchemaManagerUtils.java files
There was a problem hiding this comment.
Fixed in 7eb1237 (REST) and 37ff2bd (the other catalogs).
REST: the action is no longer part of the protocol.
- It is removed from the
SchemaChangeJSON subtypes and fromrest-catalog-open-api.yaml, which is back to master's version, so a REST server cannot parse it. RESTCatalog.alterTablerejects it before sending, with the same message the procedure gives.- Tests:
RESTApiJsonTestasserts that{"action":"enableDataEvolution"}does not parse.DataEvolutionEnablerTest#testRejectsRESTCatalognow also sends a directalterTableon a REST table that has data: it is refused and the options stay unchanged.
Other catalogs: FileSystemSchemaManager.commitChanges accepts EnableDataEvolution on a table without row tracking only while every live data file has a row id. That means the table has no snapshot yet, or its latest snapshot is a row id commit of the procedure.
- Such a commit is marked with the snapshot property
data-evolution.row-ids-assigned-snapshot-id, whose value is the snapshot's own id. A later snapshot that copies its base's properties therefore does not carry the mark, the same pattern as the reassign plan in [core] Persist row ID reassignment plans and mark snapshots #10008. SchemaManagerUtilsstill skips the immutable-option check for this change, but only after this precondition has been checked.- Tests:
testSchemaChangeAloneIsRefusedWhileFilesHaveNoRowId: a table with data is refused and keeps its options.testSchemaChangeAloneIsAcceptedOnTableWithoutSnapshottestSchemaChangeNeedsTheLatestSnapshotToBeAMarkedOne: a snapshot that only copied the mark is refused.
JingsongLi
left a comment
There was a problem hiding this comment.
The use case has strong end-to-end value: converting an existing large append table in place avoids a full data rewrite and preserves its history, while enabling Data Evolution and global indexes. The normal read/split fallback and Spark/Flink procedure paths look coherent. I cannot recommend this conversion for production yet because its safety conditions are not enforced at all interleavings:
- P1: An in-flight stale writer can commit after the procedure reports success.
FileStoreCommitImpl.tryCommitOncereads and checks the latest schema at lines 1047-1050, then does manifest work before the snapshot CAS at line 1315. If that writer pauses after the check,DataEvolutionEnablercan change the schema and finish its last repair scan (DataEvolutionEnabler.java:142-171); the writer then commits an old-schema file withoutfirstRowId. Snapshot CAS checks the snapshot UUID, not the schema, and no subsequent repair occurs. This leaves an ID-less live file in a supposedly converted table. Coordinate the schema switch and snapshot commit or otherwise fence stale commits atomically; add a deterministic paused-writer test that releases the writer after the procedure's final scan. - P1: Rollback guards use cached options. A table/commit handle opened before conversion still has
rowTrackingEnabled() == false, soAbstractFileStoreTable.checkRollbackKeepsRowTrackingandFileStoreCommitImpl.rollbackToAsLatestcan restore a pre-conversion snapshot after the latest persisted schema enabled row tracking. Akash3121 already reported the table rollback paths. Check the latest persisted schema in all rollback entry points and test snapshot, tag, and as-latest rollback through a retained old handle. - P2: The assignment snapshot can carry a stale schema ID.
DataEvolutionEnabler.commitAssignmentusestable.schema().id()at line 351. That schema is cached on the loaded table. A concurrent schema-only change does not move the snapshot CAS, so this commit can create a newer conversion snapshot referring to the older schema. Akash3121 already reported this; resolve the latest schema for each commit attempt and test anADD COLUMNbetween planning and replacement. - P1: Public schema alteration bypasses assignment.
SchemaChange.enableDataEvolution()is exposed through generic catalog/REST alteration andSchemaManagerUtilsapplies it without proving existing files have IDs. A client can enable both options on a populated table without the manifest rewrite, although the procedure itself rejects REST catalogs. Akash3121 already raised this. Restrict the action to a verified conversion path or reject direct alteration on nonempty tables; test direct REST/catalog alteration on populated data.
Production rollout also needs an explicit requirement to stop older Paimon writer binaries before conversion: the stale-writer guard exists only in this PR's runtime, so an older running binary will not enforce it after the finite repair loop.
Verification: RowTrackingEnableGuardTest (6/6), DataEvolutionFilesWithoutRowIdTest (5/5), RESTApiJsonTest (36/36), and the new SchemaManagerTest case (1/1) passed locally. The local enabler suite was blocked mainly by the existing CodeGenerator service-loading failure; one REST test hit this sandbox's bind restriction. The Flink SQL procedure IT failed at its initial write because of a local Avro DataFileWriter.setEncoder linkage error. The PR head's JDK 8/11, Flink 1, Spark, and E2E CI checks pass; Flink 2 CI is still running. These CI results do not exercise the unsafe interleavings above.
The guards that refuse a rollback across the switch to row tracking read the options of the table object. A table loaded before the switch still reports row tracking as disabled, so rollbackTo(snapshot), rollbackTo(tag) and rollbackToAsLatest through such an object could restore a snapshot whose files have no row id. Decide by the latest persisted schema instead.
SchemaChange.enableDataEvolution() switches on row tracking without assigning row ids; only sys.enable_data_evolution may issue it, after assigning them, and that procedure does not support REST catalogs. Exposed through the REST protocol, a client could flip the options over files without row id and leave a table no tool can repair. Remove the action from the JSON subtypes and the OpenAPI spec, and reject it in RESTCatalog.alterTable.
A writer that passed the commit schema check on the previous schema could still commit a file without row id after the procedure reported success: the schema change creates no snapshot, so the writer's snapshot CAS succeeded and no repair followed. After the schema change, commit an empty fence snapshot through the normal commit path, then repair what was committed before it. A writer reads its base snapshot before it checks the schema, so it either committed before the fence or loses the snapshot id to it and is refused on retry. The schema change alone is accepted only while every live data file has a row id: the table has no snapshot, or its latest snapshot is a row id commit of the procedure, marked with a property holding its own id so that snapshots copying their base's properties do not inherit the mark. The procedure assigns again and retries when a writer commits between the assignment and the schema change. Each row id commit reads the latest schema right before the replacement: a concurrent schema change does not move the snapshot and would otherwise be missing from the row id snapshot.
…ommit A plain append writer numbers its rows, so a file written in one commit of many rows has sequence numbers far above the snapshot ids of later commits. A row-tracking commit stamps its files with its snapshot id, and a data-evolution read takes each column from the file with the highest sequence number. After the conversion, a partial-column write or a MERGE INTO on such a file was therefore ignored: the converted file kept winning. Stamp each file that gets a row id with the id of the snapshot that assigns it, like any row-tracking commit does.
|
@JingsongLi thanks for the review. All four points were real. I reproduced each one with a deterministic test before changing anything, and they are fixed in new commits on top of the branch. Akash's three inline threads are answered individually with the code references. 1. In-flight stale writer (P1), 37ff2bd A writer reads its base snapshot before it checks the schema. So a writer that passed the check on the old schema has only two outcomes:
Tests pause a real writer right after its schema check, with a
2. Rollback guards on cached options (P1), 35110ec 3. Stale schema id on the row id snapshot (P2), 37ff2bd 4. Direct schema change (P1), 7eb1237 and 37ff2bd
Older writer binaries: agreed, this cannot be enforced in code. The conversion section of the docs now requires stopping writers of older Paimon versions before converting. One more bug, found while checking the tests (not raised in the review), ff3eeba Also
Verification
|
|
Re-reviewed head P1 — the public direct catalog action can still enable Data Evolution without a fence. Please make the generic action unavailable to external callers or ensure every allowed entry point performs the same fence-and-repair protocol before reporting success. Add the direct-alter in-flight-writer regression. The documented rollout must also require older writer binaries to stop, since they cannot enforce the new stale-writer check. This is a production merge blocker; keep the PR open for the fix. |
…tion on The procedure fences off writers that checked the previous schema and repairs what they committed, but a direct catalog.alterTable() with the public SchemaChange.enableDataEvolution() did neither: on an empty table, or on a snapshot whose files all had row ids, a writer in flight could still commit a file without row id after the switch. Remove the public schema change. The change is now DataEvolutionEnabler.EnableDataEvolution with a private constructor, so only the procedure can issue it; it still goes through catalog.alterTable to keep catalog locks and metastore sync. With no other caller left, the schema manager precondition and the snapshot mark are no longer needed: a writer that commits between the row id commit and the switch is repaired after the fence.
…irst snapshot row-tracking.enabled and data-evolution.enabled are immutable only once a table has snapshots, so ALTER TABLE could still switch them on for an empty table. A writer that loaded the table without row tracking may be committing the first snapshot at that moment, and its files get no row id: the race the procedure closes with its fence. Refuse the switch before the first snapshot as well, pointing to CREATE TABLE or sys.enable_data_evolution. Switching the options off and setting them when creating a table are unaffected.
|
@JingsongLi thanks, confirmed with the same probe. Fixed by leaving the procedure as the only way to switch the options on:
Your probe is added as a regression test ( |
|
Fixed three additional issues in
Added seven regression tests. 215 core tests passed, along with Checkstyle and Spotless. Spark/Flink suites were not rerun locally. |
|
Reviewed exact head
Validation: 139 focused Core tests passed with normal JDK 8 Maven checks; all 24 selected conversion/row-tracking/sub-field Spark tests passed on each of Spark 3.5.8 and Spark 4.1.2, including actual SQL conversion, update/merge/delete, historical reads and incremental BTree index use. The 6 procedure/action/streaming tests also passed on each of Flink 1.20.1 and Flink 2.2.0 after aligning the local reactor test classpath from Avro 1.11.3 to the Paimon compiler's 1.11.4; the initial mismatch failed before conversion. The additional actual-table probes above reproduce both P1s on this head. Exact-head CI is green. Older-version writers still must be stopped as documented; these findings also occur with this version's supported rollback/recovery APIs. |
…ution conversion Rolling a row-tracking-only table back to a snapshot or tag from before sys.enable_data_evolution restored the row-count sequence numbers that the conversion normalized, so later partial-column updates were silently hidden. The rollback guards only compared row-tracking.enabled. Check the data-evolution boundary as well, in one helper shared by rollbackTo and rollbackToAsLatest. A committer loaded after the conversion could still commit compaction output restored from a checkpoint of the row-tracking-only append compactor. That output has no first row id and gets none on commit, so it replaced converted files and dropped their row ids. Refuse such files in data-evolution commits before anything is published.
…utionEnablerTest SchemaManager#listAllIds lists the schema directory without sorting, so the order depends on the file system.
…ata evolution conversion A writer that loaded the table before sys.enable_data_evolution numbers the rows of all files it rolls in one sequence, and the commit keeps the numbers of a file that does not start at 0. A committer loaded after the conversion can still commit such files when a job restores them from a checkpoint, and their sequence numbers then hide later column updates. Stamp them with the commit's snapshot id, like the files of a data-evolution writer. The fence of the conversion is committed on the data-evolution schema before the files of writers that committed during the conversion are repaired, so the schema check of a rollback lets it through. Check the files that a rollback of a converted table makes live again as well, unless the latest snapshot holds them in the same state already.
A copy-on-write UPDATE, DELETE or MERGE INTO on a row-tracking table rewrites a file with the row ids of its rows stored in it and no first row id. The conversion assigned such a file a new first row id range, which contradicts the stored ids: the rows still read with their old ids, and a later column update by those ids added phantom rows instead of updating them. The stored ids need not be contiguous, so no first row id can describe them: refuse the conversion, before changing anything, and ask to rewrite those rows first. A rollback no longer counts such files as unconverted, since no commit or conversion assigns them a first row id.
…ution # Conflicts: # paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkProcedures.java
|
Thanks for the detailed probes. Both issues reproduced as described and are fixed on
While verifying I also fixed three related issues:
Flink and Spark procedure tests and the relevant core suites pass. |
JingsongLi
left a comment
There was a problem hiding this comment.
Requirement fit: SUPPORTED. Metadata-only conversion has substantial value for existing append tables and global indexing. Implementation: FINDINGS: the latest rollback and restored-compaction fixes pass, but a current-version concurrent copy-on-write update can leave an unsupported converted table and the retry then reports Skipped. Validation: 148 formal core, 9 Spark SQL procedure and 6 Flink SQL procedure cases pass. A source-free actual file commit/conversion/read-planning probe exposes the missing interleaving (with a serial rejection control). Current code/E2E CI passes; licensing failed before build on Maven Central connection reset. Actual Spark copy-file controls also reproduce overlapping row-ID allocation, including silent loss of copied rows after a successful conversion.
| resetSequences.add(entry.identifier()); | ||
| } | ||
| } | ||
| long start = latest.nextRowId() == null ? 0L : latest.nextRowId(); |
There was a problem hiding this comment.
[P1] Reserve existing row-ID ranges before assigning missing IDs
A legal copy_files from a row-tracking source into an existing ordinary append target retains file.firstRowId, while that target still has nextRowId=0. The enabler preserves these existing IDs but assigns ordinary appended files starting at zero too. I reproduced this through the actual public Spark copy workflow (CopySchemaOperator/ListDataFilesOperator/CopyDataFilesOperator/CopyFilesCommitOperator, physical Parquet), then appended rows normally. Copy two tracked rows + append two ordinary rows reads four rows before conversion; conversion reports Success, but reads only the two appended rows afterwards because both files now own [0,2) and the newer sequence wins. Copy two + append one instead makes ordinary reads fail with overlapping-range error. The same public flow with a plain source is a passing control (all rows retained). Derive and validate allocation bounds from the live ranges, or reassign/reject retained IDs when converting a non-row-tracking table; do not assume its nextRowId accounted for copied file metadata. Include copies with unequal/equal row counts and copied-only tables.
There was a problem hiding this comment.
Fixed in bbadbc6. Missing row ids now start after the largest row id of the live files instead of the snapshot's nextRowId. A table whose files from before the conversion already have overlapping row ids is refused. Tests:
DataEvolutionEnablerTest: copied rows only, copy 2 + append 2, and copy 2 + append 1 into an ordinary table. The copied rows keep row ids 0 and 1, the others continue after them, and so do rows written later. Files of two tables with the same row ids are refused.EnableDataEvolutionProcedureTest: your flow with the realsys.copyinto an ordinary table, followed by aMERGE INTOby row id.
The root cause is in sys.copy: it commits the copied files with their row ids without advancing nextRowId, which leaves duplicate row ids in row-tracking targets regardless of this PR. That is fixed separately in #10391, which also shifts copied row ids that collide with rows a partial overwrite keeps, and orders sequence numbers on data-evolution targets. The two PRs are independent: since 64aa7ce the tests here build their state without relying on how sys.copy behaves, and they pass with and without #10391.
| Long snapshotBefore = table.snapshotManager().latestSnapshotId(); | ||
| Assignment planned = plan(table); | ||
| if (enabled && !planned.hasChanges() && hasFence(table, planned.snapshot)) { | ||
| return Result.skipped( |
There was a problem hiding this comment.
[P1] Validate physical row-ID files before reporting conversion complete
A current-version row-tracking-only copy-on-write UPDATE/MERGE may commit after the initial check but before the schema switch. Its file stores physical _ROW_ID and has no firstRowId. The enabler switches DE on and commits its fence, then rejects this file, leaving the table on the DE schema. On retry, plan records it only in storingRowIds (then continues); hasChanges is false and this early Skipped return bypasses checkNoFileStoresRowIds. I reproduced this with the same real public BatchTableWrite/write-type and CommitMessage path as testRefusesFilesThatStoreTheirRowIds, injected between row-ID assignment and schema change: first run throws with DE=true/schema=1; second run returns Skipped; real DataEvolutionCompactCoordinator.plan still throws missing first row id; normal rollback to the pre-conversion snapshot is refused. The serial control correctly refuses before changing schema. This input is also produced by the normal Spark copy-on-write path, and docs only require stopping older-version writers. Coordinate these writes before the schema switch (or explicitly enforce a full writer pause), and always validate storingRowIds before the idempotent success check. Cover this interleaving and recovery, not just a file that exists at the initial scan.
There was a problem hiding this comment.
Fixed in bbadbc6. The check runs before the Skipped return, on every re-plan, and right before the schema switch. A copy-on-write commit that lands before the switch therefore fails the conversion and leaves the table unchanged (testCopyOnWriteDuringConversionIsRefusedBeforeTheSchemaSwitch).
The window that remains is a copy-on-write writer that passed its schema check before the switch and commits before the fence; closing it would take coordinating the writers. In that case every later run reports the files instead of Skipped until those rows are rewritten, after which the table is complete (testCopyOnWriteBeforeTheFenceIsReportedUntilTheRowsAreRewritten). The docs say not to run such statements during the conversion.
…w ids when enabling data evolution A copy-on-write UPDATE, DELETE or MERGE INTO on the row-tracking table that committed after the conversion planned could still be converted: the files it wrote, which store their row ids and have no first row id, were only noticed after the fence, with the schema already switched, and a later run reported the table as done. Check for them again on every replan and right before the schema switch, so that the table stays unchanged in the common case, and before reporting a table as done, so that a run after the remaining window keeps naming the files until those rows are rewritten. Files copied into the table with sys.copy keep their first row ids, but the next row id of the table does not cover them. The conversion continued from that next row id, so rows written afterwards got the row ids of the copied rows, and a data-evolution read let them hide the copied rows. Continue after the largest row id of the live files, and refuse a table whose files from before the conversion already have overlapping row ids.
JingsongLi
left a comment
There was a problem hiding this comment.
Requirement fit: SUPPORTED. Implementation: FINDINGS. Rechecked bbadbc6: the previous copied-row-ID allocation and false Skipped-after-COW findings are addressed. 152 formal core tests and all 11 Spark SQL procedure tests pass; the previous 6 Flink procedure tests passed. The actual public Spark copy/conversion/append controls retain all rows for both equal and unequal file sizes. However, an additional actual public-copy/column-update probe finds a sequence normalization gap for a row-tracking source copied into an existing ordinary append target. Production migration still requires stopping old-version writers and COW UPDATE/DELETE/MERGE, and discarding pre-conversion pending compaction state as documented.
| CoreOptions.fromMap( | ||
| table.schemaManager().schema(id).options()); | ||
| return fileOptions.rowTrackingEnabled() | ||
| && !fileOptions.dataEvolutionEnabled(); |
There was a problem hiding this comment.
[P1] Normalize copied tracked files rebound to a plain target schema
sys.copy_files permits a row-tracking source to be copied into an existing ordinary append target. It preserves firstRowId and the source sequence numbers, while rebinding file.schemaId to the target plain schema. This condition therefore excludes those files from resetSequences, and they already have a row ID so the missing-ID rewrite does not normalize them either. I reproduced this through the actual CopySchemaOperator/ListDataFilesOperator/CopyDataFilesOperator/CopyFilesCommitOperator with physical Parquet: after 10 real source overwrites, the one copied live row has firstRowId=9 and sequence=10. The plain target reads it correctly; conversion reports Success with nextRowId=10 at snapshot 3, but retains sequence=10. A column-v update to that row then commits successfully at snapshot 4 (sequence=4), yet actual reads still return old9 instead of NEW because the copied full-row file wins the sequence comparison. The same workflow into an existing row-tracking target is a passing control: conversion resets sequence to 1 and NEW is visible. Normalize retained complete-row files from every pre-DE schema, including these legal copied files, and keep the conversion-detection/fence helpers consistent; add a high-source-sequence plain-target copy followed by a real column-update regression test.
There was a problem hiding this comment.
Fixed in 8673cc1. Every file from a schema without data evolution now gets the baseline sequence number 1, whatever wrote it. That covers row-tracking-only writers, plain writers, and files copied in with the sequence numbers of another table. The same rule is used in all three places:
- the conversion plan;
DataEvolutionUtils.needsDataEvolutionConversion, which theSkippedcheck and the rollback guard use;- the stamping of restored pre-conversion appends in
FileStoreCommitImpl.
Tests:
testCopiedFilesWithHighSequenceNumbersGetTheBaselineSequence: a source overwritten ten times (first row id 9, sequence 10) is copied into an ordinary table and converted. A column update committed afterwards is visible, and a second run reportsSkipped.- The same flow in Spark with the real
sys.copy, followed byUPDATEandMERGE INTO.
…thout sys.copy The tests of copied row ids relied on sys.copy committing files into a row-tracking table without advancing its next row id, and on rows written afterwards getting the copied row ids again. That is a sys.copy bug, fixed separately, after which those states no longer arise. Build them the way they arise in any case: files with the row ids of another table in a table without row tracking, which keeps no next row id, and the files of two such tables with the same row ids in one table. Cover copied rows only and rows written after the copy, with equal and unequal row counts.
…ce number The conversion reset the sequence numbers only of files written with a row-tracking-only schema. A file copied into a table without row tracking, for example by sys.copy, keeps the first row id and the sequence numbers it had in its source, under the schema of the target. It needed no row id and matched no reset, so it kept a sequence number above the snapshot ids of the converted table, and a later column update, UPDATE or MERGE INTO of its rows was ignored. Give every file of a schema without data evolution the baseline sequence 1, older than every data-evolution update, whatever wrote it. Use the same rule where the conversion state is checked for Skipped and for rollbacks, and stamp appended files restored from a writer that predates the conversion with the same baseline, so that the rule holds for every such file.
|
@JingsongLi thanks for re-checking. All three points are answered in their threads:
Separate PR for
#10391 fixes these in This PR still has to handle copied files at conversion time, for two reasons:
The two PRs are independent. The tests here build their state without relying on how |
|
This is a reasonable requirement, but it is the wrong approach; the content of Data Evolution should not be modified, as doing so makes the code within it difficult to understand. |
|
Additionally, please split the PR. This is only a draft. And please share your production env, I need more information not just from AI. |
|
@JingsongLi thanks, agreed. The conversion should not make the Data Evolution code harder to understand, and this PR changes too much of the existing Data Evolution core logic for that. Our use case: in production we have a very large wide append table that needed to become a Data Evolution table. The launch was urgent and there was no way to convert it in place, so we had to rewrite all of its data into a new Data Evolution table. A metadata-only conversion would have avoided that rewrite. I'll redesign it so that the conversion stays outside the existing Data Evolution read/write/commit paths as much as possible, and split it into small PRs that are easy to review. I'll keep this PR as a draft until then. |
Purpose
Feature: convert an existing append table to a Data Evolution table in place, without rewriting any data file.
Why
row-tracking.enabledanddata-evolution.enabledare immutable (ALTER TABLE ... SETis 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 andINSERT OVERWRITEeverything 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.Global indexes are the second reason. A global index stores row ids, so
CreateGlobalIndexProcedurerefuses any table withoutrow-tracking.enabled, in Spark, Flink and PyPaimon alike — a plain append table cannot have one at all, whatever the index type:rewrite_file_indexbackfills themFile indexes live next to each data file and are rebuilt whenever a file is rewritten, so file-level pruning has always been available. Row-id addressing is what a global index needs, and it is what unlocks full-text and vector search,
LIKE '%x%'through FM, and index-driven top-k. Unlike partial-column updates, that costs nothing on the write side: it is pure read-side gain. Once again the tables people most want to index — a multi-TB log or document table — are the ones they can least afford to rewrite.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 nofirstRowId. Every code path that assumes the id (nonNullFirstRowId, the split generator, compaction,MERGE INTOon_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
The call returns one row:
Dry run. Enabling data evolution on table 'db.t' would do: ...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.OVERWRITEsnapshot that assigns the row ids, and an emptyAPPENDsnapshot that fences off writers which checked the previous schema (plus one moreOVERWRITEsnapshot if such a writer committed before the fence)Skipped. Table 'db.t' was not changed: data evolution is already enabled.What the table looks like afterwards
CREATE TABLE ... TBLPROPERTIES ('row-tracking.enabled'='true', 'data-evolution.enabled'='true'). Nothing else is switched on:bucket, deletion vectors anddata-evolution.nested-field.enabledstay as they were and can be changed afterwards withALTER TABLEas on any Data Evolution table.firstRowId, and its sequence numbers are set to the id of the snapshot that assigns it, as for any row-tracking commit, so that later column writes take precedence. On a table that already had row tracking, the files keep their row ids and their sequence numbers are reset to 1: a row-tracking-only writer numbers the rows of a file, which would hide later column writes. Row ids are contiguous per partition, in commit order, sosys.reassign_row_idhas nothing to do afterwards. New writes continue atnextRowId.ALTER TABLErefuses both options even before the first snapshot.Requirements: append table without primary key,
bucket = -1(the default for an append table),clustering.incrementaloff. Tables of a REST catalog are rejected in this version. A row-tracking table whose files store the row ids of their rows (written by a copy-on-writeUPDATE,DELETEorMERGE INTO) is refused before anything changes, because those ids need not be contiguous and no first row id can describe them; rewrite those rows first, for example withINSERT OVERWRITE.What to expect around the conversion (also in the docs):
_ROW_IDreads asNULL. Rolling back to them (rollback_to, also by tag, androllback_to_as_latest) is refused while data evolution is enabled, also for a table that already had row tracking, since that would restore the old sequence numbers. So is a rollback to a snapshot that still holds files the conversion had not repaired yet.What
One PR. Commits 1–4 below are the feature, followed by
[core] Report the table's row id when a resumed run assigns none(a run that finds the row ids already committed by an earlier, interrupted run reports the table'snextRowIdinstead of none). The commits after that address the review, see "Review follow-ups" below.1.
[core] Guard commits and rollbacks across a row-tracking switchTwo guards that make the switch safe against concurrent writers and against history:
FileStoreCommitImplrefuses 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) androllbackToAsLatestrefuse 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.replaceManifestListgets an overload that takes the schema id andnextRowIdexplicitly (used by commit 3).2.
[core] Read data-evolution files that have no first row id as plain filesA 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.splitByRowIdPresenceseparates them;DataEvolutionSplitGeneratorbin-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.DataEvolutionSplitReadProviderdoes not match such splits;AppendOnlyFileStoreTable.newReadalways registers the plain append raw-file provider after the Data Evolution one, so they are read as ordinary append files with_ROW_ID = NULL.DataEvolutionFileStoreScangroups them as singletons and applies file statistics to them individually (predicate pushdown keeps working on them).DataFileMeta.nonNullFirstRowIdandDataEvolutionCompactRangePlannername the file in their message and point atsys.enable_data_evolution, instead of a barefirstRowId must not be null.3.
[core] Add DataEvolutionEnabler to convert an append table in placeSchemaManagerUtilslets it through the immutable-option check;validateTableSchemastill enforces every other Data Evolution constraint (no primary key,bucket = -1, noclustering.incremental). After the review it becameDataEvolutionEnabler.EnableDataEvolutionwith a private constructor, so only the procedure can issue it; it is not part of the publicSchemaChangeAPI or the REST protocol.DataEvolutionEnabler(catalog, identifier).run(dryRun):validateTableSchema; REST catalog tables are rejected in this version (the manifest rewrite is a directreplaceManifestListon the file store, which the REST protocol does not expose yet).ADDentries without id are ordered per partition in manifest/commit order and given contiguous ids starting atnextRowId(or 0); only the manifests that hold converted files are rewritten, and the snapshot is replaced withreplaceManifestListcarrying the newnextRowId. On a conflict (another commit landed) it re-plans from the new latest snapshot, bounded by the table's commit retry settings.catalog.alterTablewith the change above, which keeps catalog locks and metastore sync. Skipped when a concurrent run switched it already.Skipped.reassign_row_idafterwards.4.
[spark][flink] Add the enable_data_evolution procedure and actionCALL sys.enable_data_evolution(table => 't' [, dry_run => true]), FlinkCALL sys.enable_data_evolution('db.t' [, dry_run])and theenable_data_evolutionaction jar.MERGE INTOon 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.BaseScanreads_ROW_ID/_SEQUENCE_NUMBERas nullable so such files read asNULL(a non-nullable read type turned it into0).copy(...)→copyWithoutTimeTravel(...)changes inMergeIntoPaimonDataEvolutionTable(Spark 3 and the Spark 4.0 copy) andDataEvolutionPaimonWriter. This fixes a pre-existing bug on native Data Evolution tables too: the merge pins its scan withscan.snapshot-id, andcopyre-applies that time travel and adopts the schema the snapshot was committed with. After a schema-only change (ALTER TABLEcreates no snapshot: enablingdata-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 inTableSchema.project.copyWithoutTimeTravelkeeps the current schema while still pinning the snapshot.Review follow-ups (each finding was first reproduced by a deterministic regression test):
[core] Check the latest schema in the rollback guards:rollbackTo(snapshot),rollbackTo(tag)androllbackToAsLatestdecide by the latest persisted schema, not by the options of the table object, so a table loaded before the conversion cannot roll back across it.[api][core] Keep enabling data evolution out of the REST protocol: removed from theSchemaChangeJSON subtypes and the OpenAPI spec (which is back to master).[core] Fence writers and verify row ids when enabling data evolution: the fence and repair above; every row id commit reads the latest schema right before the replacement (a concurrentADD COLUMNdoes not move the snapshot).[core] Stamp converted files with the sequence number of the row id commit: found while checking the tests. A plain append writer numbers its rows, so a file written in one commit of many rows has sequence numbers far above the snapshot ids of later commits; a data-evolution read takes each column from the file with the highest one, so after the conversion a partial-column write orMERGE INTOon such a file was silently ignored.[core] Test expiring snapshots and tags from before the data evolution conversion: the data files, which the conversion re-references through new manifests, are kept.[core] Only the enable_data_evolution procedure can switch data evolution on: the publicSchemaChange.enableDataEvolution()is removed; a directcatalog.alterTablecould switch the options without the fence. The schema-manager precondition and snapshot mark of the earlier follow-up are removed with it.[core] Refuse switching row tracking on with ALTER TABLE before the first snapshot: the same race existed on master, because the options are only immutable once a table has snapshots. Behaviour change:ALTER TABLE SET ('row-tracking.enabled' | 'data-evolution.enabled' = 'true')is now refused on an empty table too (set it atCREATE TABLEor call the procedure); switching them off is unaffected.[core] Fix row-tracking conversion sequences and resumed writer fences: a table that already had row tracking keeps its row ids, and its files get sequence number 1, below every data-evolution update, including one committed while the repair retries. A writer or compactor loaded before data evolution was enabled is refused even if row tracking was on already. An interrupted conversion reportsSkipped.only once a snapshot on the data-evolution schema has fenced old writers; otherwise the next run completes the fence and repair, also on an empty table.[core] Refuse rollbacks and restored compactions across the data evolution conversion(second review): the rollback guards also check the data-evolution boundary, in one helper shared byrollbackToandrollbackToAsLatest. A data-evolution commit refuses an added file that would be left without row ids (notAPPEND, no first row id, no stored row ids), which is what compaction output restored from a pre-conversion checkpoint is.[core] Do not depend on the listing order of schema files in DataEvolutionEnablerTest: test only.[core] Cover restored appends and unrepaired fence snapshots of the data evolution conversion: appended files restored from a pre-conversion writer are stamped with the commit's snapshot id instead of keeping row-count sequence numbers. A rollback of a converted table also checks the files the target makes live again, so a rollback to the fence snapshot before the repair finished is refused.[core] Refuse enabling data evolution on files that store their row ids: see Requirements; a rollback no longer counts such files as unconverted.Benefit
ALTER TABLEon existing Data Evolution tables now write files with the right schema id.CALL sys.create_global_index(...)— which is itself incremental, so the first call indexes the existing data and later calls only cover what has been appended since.Tests
paimon-core
RowTrackingEnableGuardTest(7): 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, tag and as-latest), also through a table object loaded before the switch, plain append table unaffected,replaceManifestListwith schema id.DataEvolutionFilesWithoutRowIdTest(5): full read / projection / predicate on files without id (_ROW_IDNULL, 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,nonNullFirstRowIdmessage.DataEvolutionEnablerTest(42): empty table switches the schema and commits the fence; 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; dry run changes nothing; partial-column write and Data Evolution compaction on a converted table; a column write after converting a 20-row commit wins over the converted file; expiring pre-conversion snapshots and deleting a pre-conversion tag keep the data; converts one branch only; a row-tracking-only table keeps its row ids; rejects primary-key / bucketed / incremental-clustering / REST catalog tables; concurrent append and compaction before the row id commit are assigned on retry; a stale write between the row id commit and the switch is repaired; the row id commit references a schema changed concurrently; two concurrent runs switch the schema once; a writer paused right after its schema check (aFileIOthat blocks its first manifest write) is refused when released after the procedure, repaired when released before the fence, refused on an empty table; the schema change cannot be created outside the enabler;ALTER TABLEcannot enable row tracking, with or without snapshots, also with a writer in flight (the reviewer's probe). Every conversion test asserts that no live file lacks a row id or has a sequence number above the latest snapshot id. Added in the later follow-ups: a column write after converting a row-tracking-only table wins over its old sequence numbers, also when committed while the repair retries; row-tracking-only writers and compactors loaded before the conversion are refused; row-tracking-only compaction committed before the fence is fully repaired; restored compaction output of a row-tracking-only table is refused after the conversion (a realAppendCompactTaskmessage serialized and restored throughfilterAndCommitMultiple); a restored append of a pre-conversion writer is stamped with its snapshot; rollback to a snapshot, a tag or as-latest from before converting a row-tracking-only table is refused, and so is a rollback to the fence before the repair; a retry after an interrupted schema switch still fences a paused writer, also on an empty table; files that store their row ids are refused.SchemaManagerTest#testSetOptionCannotEnableRowTrackingWithoutSnapshots;RESTApiJsonTestasserts that a server cannot parse anenableDataEvolutionaction.DataEvolution*suites,AppendOnlySimpleTableTest,BlobTableTestpass.paimon-spark (Spark 3.5)
EnableDataEvolutionProcedureTest(9):UPDATEandMERGE INTOafter converting a 100-row commit are applied; basic conversion + dry run + idempotence +MERGE INTO/UPDATE/INSERTafterwards (partial-column files, row ids continue after the assigned range); partition-contiguous ids andreassign_row_idreportsSkipped.; deletion vectors kept andDELETE/ merge-delete keep working; sub-field data evolution on a converted table; pre-conversion snapshot and tag readable with_ROW_IDNULL; refusals (primary key, bucketed,ALTER TABLEwith and without data);MERGE INTOrefuses files without id, then the procedure repairs the table.RowTrackingTestBase: newmerge 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,DataEvolutionUpdateSnapshotTestpass.paimon-flink
EnableDataEvolutionProcedureITCase(6): procedure (ids read throught$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 snapshots without replaying the table (default andstreaming-read-append-overwrite),ALTER TABLEcannot enable row tracking.spotless:checkandcheckstyle:checkpass on paimon-api, paimon-core, paimon-spark-common, paimon-spark-4.0, paimon-spark-ut and paimon-flink-common.API and Format
SchemaChange, no REST protocol change (the OpenAPI spec is unchanged). The REST catalog does not support the conversion yet; the procedure reports this.ALTER TABLEcan no longer switchrow-tracking.enabledordata-evolution.enabledon for a table without snapshots (it already could not once the table had snapshots).firstRowId = null; they are read as complete-row files with_ROW_ID = NULL. The file format itself is unchanged.sys.enable_data_evolution, Flink proceduresys.enable_data_evolutionand actionenable_data_evolution.Documentation
multimodal-table/data-evolution.mdx: new section "Enable on an Existing Append Table" (what the procedure does, that it is the only way to switch the options on, what to expect around the conversion: older-version writers must be stopped, stale writers, jobs restored from a checkpoint, tables whose files store row ids, 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.