zhuxiangyi commented on code in PR #10098: URL: https://github.com/apache/paimon/pull/10098#discussion_r4183373255
########## 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: Fixed in bbadbc65f. 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 real `sys.copy` into an ordinary table, followed by a `MERGE INTO` by 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 64aa7ce8c the tests here build their state without relying on how `sys.copy` behaves, and they pass with and without #10391. ########## 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: Fixed in bbadbc65f. 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. ########## paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionEnabler.java: ########## @@ -0,0 +1,764 @@ +/* + * 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.io.DataFileMeta; +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.Comparator; +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.format.blob.BlobFileFormat.isBlobFile; +import static org.apache.paimon.operation.commit.RowTrackingCommitUtils.storesRowIds; +import static org.apache.paimon.types.VectorType.isVectorStoreFile; +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); + // Also before reporting the table as done: an earlier run may have switched the schema and + // then found files it cannot convert. + checkConvertible(planned, enabled); + if (enabled && !planned.hasChanges() && hasFence(table, planned.snapshot)) { + return Result.skipped( + schemaBefore, snapshotBefore, "data evolution is already enabled"); + } + 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); + checkConvertible(remaining, true); + 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; + } + // A writer on the current schema may have committed files that cannot be converted since + // the plan. Refuse while the table is still unchanged: once the schema is switched, the + // table stays on it. Only a writer that passed its schema check before the switch and + // commits before the fence can still slip in, see run. + checkConvertible(plan(table), false); + catalog.alterTable(identifier, new EnableDataEvolution(), false); + } + + /** + * Refuses files that the conversion cannot give correct row ids. {@code enabled} tells whether + * the schema is switched already, so that the message says the table is left on it. + * + * <ul> + * <li>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. + * <li>Files written before the conversion are complete-row files, so their row id ranges must + * not overlap: a data-evolution read merges files of overlapping ranges, and rows of one + * would hide the rows of the other. Ranges overlap when files that carry row ids are + * copied into a table, for example by {@code sys.copy}, which does not advance the next + * row id of the table, and rows are written afterwards. + * </ul> + * + * <p>Rewriting those rows, for example with INSERT OVERWRITE, gives them new row ids. + */ + private void checkConvertible(Assignment assignment, boolean enabled) { + String problem; + if (!assignment.storingRowIds.isEmpty()) { + problem = + String.format( + "%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. " + + "Files: %s", + assignment.storingRowIds.size(), describe(assignment.storingRowIds)); + } else if (!assignment.overlappingRowIds.isEmpty()) { + problem = + String.format( + "data files written before the conversion were assigned " + + "overlapping row ids, for example by copying files into " + + "the table with sys.copy and writing rows afterwards, so a " + + "data-evolution read would let the rows of one hide the " + + "rows of another. Files: %s", + describe(assignment.overlappingRowIds)); + } else { + return; + } + throw new IllegalArgumentException( + String.format( + "%s: %s. Rewrite those rows first, for example with INSERT OVERWRITE, " + + "and run the procedure again.", + enabled + ? String.format( + "Table %s has data evolution enabled, but cannot be " + + "fully converted", + identifier.getFullName()) + : String.format( + "Cannot enable data evolution on table %s", + identifier.getFullName()), + problem)); + } + + private static String describe(List<String> files) { + int shown = Math.min(5, files.size()); + return files.subList(0, shown) + (shown < files.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<>(); + Map<Long, Boolean> dataEvolutionSchemas = new HashMap<>(); + List<ManifestEntry> completeRowFiles = new ArrayList<>(); + long maxRowIdEnd = 0L; + for (ManifestEntry entry : live.values()) { + if (entry.kind() != FileKind.ADD) { + continue; + } + if (entry.file().firstRowId() == null && storesRowIds(entry.file())) { + // see checkConvertible + storingRowIds.add(entry.file().fileName()); + continue; + } + if (entry.file().firstRowId() != null) { + maxRowIdEnd = + Math.max(maxRowIdEnd, entry.file().firstRowId() + entry.file().rowCount()); + if (!isBlobFile(entry.file().fileName()) + && !isVectorStoreFile(entry.file().fileName()) + && !dataEvolutionSchemas.computeIfAbsent( + entry.file().schemaId(), + id -> + CoreOptions.fromMap( + table.schemaManager().schema(id).options()) + .dataEvolutionEnabled())) { + completeRowFiles.add(entry); + } + } + 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(); Review Comment: Fixed in 8673cc17a. 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 the `Skipped` check 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 reports `Skipped`. - The same flow in Spark with the real `sys.copy`, followed by `UPDATE` and `MERGE INTO`. -- 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]
