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]

Reply via email to