JingsongLi commented on code in PR #10098: URL: https://github.com/apache/paimon/pull/10098#discussion_r4180226998
########## paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionEnabler.java: ########## @@ -0,0 +1,664 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.append.dataevolution; + +import org.apache.paimon.CoreOptions; +import org.apache.paimon.Snapshot; +import org.apache.paimon.catalog.Catalog; +import org.apache.paimon.catalog.DelegateCatalog; +import org.apache.paimon.catalog.Identifier; +import org.apache.paimon.codegen.CodeGenUtils; +import org.apache.paimon.codegen.RecordComparator; +import org.apache.paimon.manifest.FileEntry; +import org.apache.paimon.manifest.FileKind; +import org.apache.paimon.manifest.ManifestCommittable; +import org.apache.paimon.manifest.ManifestEntry; +import org.apache.paimon.manifest.ManifestFile; +import org.apache.paimon.manifest.ManifestFileMeta; +import org.apache.paimon.manifest.ManifestList; +import org.apache.paimon.operation.FileStoreCommit; +import org.apache.paimon.operation.FileStoreCommitImpl; +import org.apache.paimon.rest.RESTCatalog; +import org.apache.paimon.schema.SchemaChange; +import org.apache.paimon.schema.SchemaValidation; +import org.apache.paimon.schema.TableSchema; +import org.apache.paimon.table.FileStoreTable; +import org.apache.paimon.table.Table; +import org.apache.paimon.table.sink.BatchWriteBuilder; +import org.apache.paimon.utils.Pair; +import org.apache.paimon.utils.RetryWaiter; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import javax.annotation.Nullable; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.UUID; + +import static org.apache.paimon.operation.commit.RowTrackingCommitUtils.storesRowIds; +import static org.apache.paimon.utils.Preconditions.checkArgument; +import static org.apache.paimon.utils.Preconditions.checkState; + +/** + * Enables data evolution on an existing append table without rewriting its data files. + * + * <p>{@code row-tracking.enabled} and {@code data-evolution.enabled} are immutable for {@code ALTER + * TABLE} because a data-evolution table derives every row id from the first row id of its file, and + * only the commit that adds a file assigns one. Files that are already in the table would therefore + * never get a row id. This class closes that gap in four steps: + * + * <ol> + * <li>Assign a first row id to every live data file that has none, by rewriting the manifests of + * the latest snapshot and committing them as a metadata-only snapshot. Ids are contiguous per + * partition, in the order {@code sys.reassign_row_id} would produce, so the converted table + * needs no reassignment afterwards. + * <li>Commit a schema with both options enabled through the catalog, so that catalog metadata + * stays in sync. Only this class can create that schema change, see {@link + * EnableDataEvolution}. + * <li>Commit a fence: an empty snapshot on the new schema. A writer that checked the previous + * schema read its base snapshot before that check, so it either committed before the fence, + * or its commit loses the race for the next snapshot id and is refused on retry, see {@code + * FileStoreCommitImpl}. + * <li>Assign ids to the files committed before the fence without one. + * </ol> + * + * <p>The procedure is idempotent: it reports {@code skipped} once the files have row ids and + * compatible sequence numbers, and a snapshot on a data-evolution schema has fenced old writers. + */ +public class DataEvolutionEnabler { + + private static final Logger LOG = LoggerFactory.getLogger(DataEvolutionEnabler.class); + private static final String COMMIT_USER_PREFIX = "enable-data-evolution"; + + private final Catalog catalog; + private final Identifier identifier; + private final Runnable beforeRowIdCommit; + private final Runnable beforeSchemaChange; + private final Runnable beforeFence; + + public DataEvolutionEnabler(Catalog catalog, Identifier identifier) { + this(catalog, identifier, () -> {}, () -> {}, () -> {}); + } + + DataEvolutionEnabler( + Catalog catalog, + Identifier identifier, + Runnable beforeRowIdCommit, + Runnable beforeSchemaChange) { + this(catalog, identifier, beforeRowIdCommit, beforeSchemaChange, () -> {}); + } + + /** Hooks for tests to inject concurrent activity between the steps. */ + DataEvolutionEnabler( + Catalog catalog, + Identifier identifier, + Runnable beforeRowIdCommit, + Runnable beforeSchemaChange, + Runnable beforeFence) { + this.catalog = catalog; + this.identifier = identifier; + this.beforeRowIdCommit = beforeRowIdCommit; + this.beforeSchemaChange = beforeSchemaChange; + this.beforeFence = beforeFence; + } + + /** Validates and, unless {@code dryRun}, converts the table. */ + public Result run(boolean dryRun) throws Exception { + FileStoreTable table = loadTable(); + CoreOptions options = table.coreOptions(); + boolean enabled = options.rowTrackingEnabled() && options.dataEvolutionEnabled(); + validate(table, enabled); + + long schemaBefore = table.schema().id(); + Long snapshotBefore = table.snapshotManager().latestSnapshotId(); + Assignment planned = plan(table); + if (enabled && !planned.hasChanges() && hasFence(table, planned.snapshot)) { + return Result.skipped( + schemaBefore, snapshotBefore, "data evolution is already enabled"); + } + checkNoFileStoresRowIds(planned); + if (dryRun) { + return Result.dryRun(schemaBefore, snapshotBefore, enabled, planned); + } + + Totals totals = new Totals(); + if (!enabled) { + assignRowIdsAndEnable(table, planned, totals); + table = loadTable(); + checkState( + table.coreOptions().rowTrackingEnabled() + && table.coreOptions().dataEvolutionEnabled(), + "Schema change did not enable data evolution on table %s.", + identifier.getFullName()); + } + + // A writer that checked the previous schema may still be on its way to commit. After the + // fence it can no longer succeed, so the files that need a row id are final. + beforeFence.run(); + commitFence(table); + Assignment remaining = plan(table); + checkNoFileStoresRowIds(remaining); + if (remaining.hasChanges()) { + LOG.info( + "Repairing row ids or sequence numbers of table {} after fencing old writers.", + identifier.getFullName()); + totals.add(assignRowIdsWithRetry(table, remaining)); + } + checkState( + !plan(table).hasChanges(), + "Table %s still has data files without a row id or with incompatible sequence " + + "numbers. A writer of an older Paimon " + + "version may still be writing to it; stop it and run the procedure " + + "again.", + identifier.getFullName()); + + Snapshot latest = table.snapshotManager().latestSnapshot(); + return new Result( + schemaBefore, + table.schema().id(), + snapshotBefore, + latest == null ? null : latest.id(), + totals.files, + totals.rows, + latest == null ? null : latest.nextRowId(), + false, + false, + null); + } + + /** + * Assigns row ids to the files of the latest snapshot and switches the schema. Files that a + * writer on the previous schema commits in between get their row ids after the fence. + */ + private void assignRowIdsAndEnable(FileStoreTable table, Assignment planned, Totals totals) + throws Exception { + if (planned.hasChanges()) { + totals.add(assignRowIdsWithRetry(table, planned)); + } + beforeSchemaChange.run(); + if (dataEvolutionEnabled(table)) { + // a concurrent run switched the schema already + return; + } + catalog.alterTable(identifier, new EnableDataEvolution(), false); + } + + /** + * 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. Such ids need not be contiguous, so no + * first row id describes them, and a new one would contradict the stored ids: a later column + * update by the row ids the rows read with would not reach them. A commit never assigns a first + * row id to such a file either. Refuse the conversion; rewriting those rows, for example with + * INSERT OVERWRITE, gives them new row ids that the conversion can assign. + */ + private void checkNoFileStoresRowIds(Assignment assignment) { + if (assignment.storingRowIds.isEmpty()) { + return; + } + int shown = Math.min(5, assignment.storingRowIds.size()); + throw new IllegalArgumentException( + String.format( + "Cannot enable data evolution on table %s: %d data file(s) store the row " + + "ids of their rows, written by a copy-on-write UPDATE, DELETE or " + + "MERGE INTO on the row-tracking table, and have no first row " + + "id, so their rows cannot be addressed by a data-evolution " + + "update. Rewrite those rows first, for example with INSERT " + + "OVERWRITE. Files: %s%s", + identifier.getFullName(), + assignment.storingRowIds.size(), + assignment.storingRowIds.subList(0, shown), + shown < assignment.storingRowIds.size() ? " ..." : "")); + } + + private boolean dataEvolutionEnabled(FileStoreTable table) { + CoreOptions latest = + CoreOptions.fromMap( + table.schemaManager() + .latestOrThrow( + "Cannot get latest schema for table " + + identifier.getFullName()) + .options()); + return latest.rowTrackingEnabled() && latest.dataEvolutionEnabled(); + } + + private boolean hasFence(FileStoreTable table, @Nullable Snapshot snapshot) { + // A writer that passed the old schema check read its base snapshot before that check. + // Any snapshot on a DE schema therefore forces it to retry (and fail the new check). + // Merely seeing the enabled schema is insufficient after a failure before commitFence. + return snapshot != null + && CoreOptions.fromMap(table.schemaManager().schema(snapshot.schemaId()).options()) + .dataEvolutionEnabled(); + } + + /** + * Commits an empty snapshot through the normal commit path, which refuses writers on a schema + * without row tracking. + */ + private void commitFence(FileStoreTable table) throws Exception { + String commitUser = COMMIT_USER_PREFIX + "-" + UUID.randomUUID(); + try (FileStoreCommit commit = table.store().newCommit(commitUser, table)) { + commit.ignoreEmptyCommit(false); + commit.commit(new ManifestCommittable(BatchWriteBuilder.COMMIT_IDENTIFIER), false); + } + } + + private FileStoreTable loadTable() throws Exception { + Table table = catalog.getTable(identifier); + checkArgument( + table instanceof FileStoreTable, + "Only a FileStoreTable can enable data evolution, but table %s is a %s.", + identifier.getFullName(), + table.getClass().getSimpleName()); + return (FileStoreTable) table; + } + + private void validate(FileStoreTable table, boolean enabled) { + checkArgument( + !(DelegateCatalog.rootCatalog(catalog) instanceof RESTCatalog), + "Enabling data evolution on table %s of a REST catalog is not supported yet.", + identifier.getFullName()); + if (enabled) { + return; + } + // The constraints of a row-tracking table are the ones of the schema this will create. + TableSchema current = table.schema(); + Map<String, String> options = new HashMap<>(current.options()); + options.put(CoreOptions.ROW_TRACKING_ENABLED.key(), "true"); + options.put(CoreOptions.DATA_EVOLUTION_ENABLED.key(), "true"); + try { + SchemaValidation.validateTableSchema( + new TableSchema( + current.id() + 1, + current.fields(), + current.highestFieldId(), + current.partitionKeys(), + current.primaryKeys(), + options, + current.comment())); + } catch (RuntimeException e) { + throw new IllegalArgumentException( + String.format( + "Cannot enable data evolution on table %s: %s", + identifier.getFullName(), e.getMessage()), + e); + } + } + + /** Plans missing row ids and normalizes the baseline of row-tracking-only files. */ + private Assignment plan(FileStoreTable table) { + Snapshot latest = table.snapshotManager().latestSnapshot(); + if (latest == null) { + return Assignment.empty(null, 0L); + } + ManifestFile manifestFile = table.store().manifestFileFactory().create(); + ManifestList manifestList = table.store().manifestListFactory().create(); + List<ManifestFileMeta> manifests = manifestList.readDataManifests(latest); + + Map<FileEntry.Identifier, ManifestEntry> live = new LinkedHashMap<>(); + FileEntry.mergeEntries( + manifestFile, manifests, live, table.coreOptions().scanManifestParallelism()); + + List<ManifestEntry> withoutRowId = new ArrayList<>(); + List<String> storingRowIds = new ArrayList<>(); + Set<FileEntry.Identifier> resetSequences = new HashSet<>(); + Map<Long, Boolean> rowTrackingOnlySchemas = new HashMap<>(); + for (ManifestEntry entry : live.values()) { + if (entry.kind() != FileKind.ADD) { + continue; + } + if (entry.file().firstRowId() == null && storesRowIds(entry.file())) { + // see checkNoFileStoresRowIds + storingRowIds.add(entry.file().fileName()); + continue; + } + if (entry.file().firstRowId() == null) { + withoutRowId.add(entry); + } + if ((entry.file().firstRowId() == null + || entry.file().minSequenceNumber() != Snapshot.FIRST_SNAPSHOT_ID + || entry.file().maxSequenceNumber() != Snapshot.FIRST_SNAPSHOT_ID) + && rowTrackingOnlySchemas.computeIfAbsent( + entry.file().schemaId(), + id -> { + CoreOptions fileOptions = + CoreOptions.fromMap( + table.schemaManager().schema(id).options()); + return fileOptions.rowTrackingEnabled() + && !fileOptions.dataEvolutionEnabled(); + })) { + resetSequences.add(entry.identifier()); + } + } + long start = latest.nextRowId() == null ? 0L : latest.nextRowId(); Review 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. ########## paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionEnabler.java: ########## @@ -0,0 +1,664 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.append.dataevolution; + +import org.apache.paimon.CoreOptions; +import org.apache.paimon.Snapshot; +import org.apache.paimon.catalog.Catalog; +import org.apache.paimon.catalog.DelegateCatalog; +import org.apache.paimon.catalog.Identifier; +import org.apache.paimon.codegen.CodeGenUtils; +import org.apache.paimon.codegen.RecordComparator; +import org.apache.paimon.manifest.FileEntry; +import org.apache.paimon.manifest.FileKind; +import org.apache.paimon.manifest.ManifestCommittable; +import org.apache.paimon.manifest.ManifestEntry; +import org.apache.paimon.manifest.ManifestFile; +import org.apache.paimon.manifest.ManifestFileMeta; +import org.apache.paimon.manifest.ManifestList; +import org.apache.paimon.operation.FileStoreCommit; +import org.apache.paimon.operation.FileStoreCommitImpl; +import org.apache.paimon.rest.RESTCatalog; +import org.apache.paimon.schema.SchemaChange; +import org.apache.paimon.schema.SchemaValidation; +import org.apache.paimon.schema.TableSchema; +import org.apache.paimon.table.FileStoreTable; +import org.apache.paimon.table.Table; +import org.apache.paimon.table.sink.BatchWriteBuilder; +import org.apache.paimon.utils.Pair; +import org.apache.paimon.utils.RetryWaiter; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import javax.annotation.Nullable; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.UUID; + +import static org.apache.paimon.operation.commit.RowTrackingCommitUtils.storesRowIds; +import static org.apache.paimon.utils.Preconditions.checkArgument; +import static org.apache.paimon.utils.Preconditions.checkState; + +/** + * Enables data evolution on an existing append table without rewriting its data files. + * + * <p>{@code row-tracking.enabled} and {@code data-evolution.enabled} are immutable for {@code ALTER + * TABLE} because a data-evolution table derives every row id from the first row id of its file, and + * only the commit that adds a file assigns one. Files that are already in the table would therefore + * never get a row id. This class closes that gap in four steps: + * + * <ol> + * <li>Assign a first row id to every live data file that has none, by rewriting the manifests of + * the latest snapshot and committing them as a metadata-only snapshot. Ids are contiguous per + * partition, in the order {@code sys.reassign_row_id} would produce, so the converted table + * needs no reassignment afterwards. + * <li>Commit a schema with both options enabled through the catalog, so that catalog metadata + * stays in sync. Only this class can create that schema change, see {@link + * EnableDataEvolution}. + * <li>Commit a fence: an empty snapshot on the new schema. A writer that checked the previous + * schema read its base snapshot before that check, so it either committed before the fence, + * or its commit loses the race for the next snapshot id and is refused on retry, see {@code + * FileStoreCommitImpl}. + * <li>Assign ids to the files committed before the fence without one. + * </ol> + * + * <p>The procedure is idempotent: it reports {@code skipped} once the files have row ids and + * compatible sequence numbers, and a snapshot on a data-evolution schema has fenced old writers. + */ +public class DataEvolutionEnabler { + + private static final Logger LOG = LoggerFactory.getLogger(DataEvolutionEnabler.class); + private static final String COMMIT_USER_PREFIX = "enable-data-evolution"; + + private final Catalog catalog; + private final Identifier identifier; + private final Runnable beforeRowIdCommit; + private final Runnable beforeSchemaChange; + private final Runnable beforeFence; + + public DataEvolutionEnabler(Catalog catalog, Identifier identifier) { + this(catalog, identifier, () -> {}, () -> {}, () -> {}); + } + + DataEvolutionEnabler( + Catalog catalog, + Identifier identifier, + Runnable beforeRowIdCommit, + Runnable beforeSchemaChange) { + this(catalog, identifier, beforeRowIdCommit, beforeSchemaChange, () -> {}); + } + + /** Hooks for tests to inject concurrent activity between the steps. */ + DataEvolutionEnabler( + Catalog catalog, + Identifier identifier, + Runnable beforeRowIdCommit, + Runnable beforeSchemaChange, + Runnable beforeFence) { + this.catalog = catalog; + this.identifier = identifier; + this.beforeRowIdCommit = beforeRowIdCommit; + this.beforeSchemaChange = beforeSchemaChange; + this.beforeFence = beforeFence; + } + + /** Validates and, unless {@code dryRun}, converts the table. */ + public Result run(boolean dryRun) throws Exception { + FileStoreTable table = loadTable(); + CoreOptions options = table.coreOptions(); + boolean enabled = options.rowTrackingEnabled() && options.dataEvolutionEnabled(); + validate(table, enabled); + + long schemaBefore = table.schema().id(); + Long snapshotBefore = table.snapshotManager().latestSnapshotId(); + Assignment planned = plan(table); + if (enabled && !planned.hasChanges() && hasFence(table, planned.snapshot)) { + return Result.skipped( Review 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 stoppin g 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. -- 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]
