JingsongLi commented on code in PR #10098:
URL: https://github.com/apache/paimon/pull/10098#discussion_r4183041057


##########
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:
   [P1] Normalize copied tracked files rebound to a plain target schema
   
   sys.copy_files permits a row-tracking source to be copied into an existing 
ordinary append target. It preserves firstRowId and the source sequence 
numbers, while rebinding file.schemaId to the target plain schema. This 
condition therefore excludes those files from resetSequences, and they already 
have a row ID so the missing-ID rewrite does not normalize them either. I 
reproduced this through the actual 
CopySchemaOperator/ListDataFilesOperator/CopyDataFilesOperator/CopyFilesCommitOperator
 with physical Parquet: after 10 real source overwrites, the one copied live 
row has firstRowId=9 and sequence=10. The plain target reads it correctly; 
conversion reports Success with nextRowId=10 at snapshot 3, but retains 
sequence=10. A column-v update to that row then commits successfully at 
snapshot 4 (sequence=4), yet actual reads still return old9 instead of NEW 
because the copied full-row file wins the sequence comparison. The same 
workflow into an existing row-tracking target is a passing 
 control: conversion resets sequence to 1 and NEW is visible. Normalize 
retained complete-row files from every pre-DE schema, including these legal 
copied files, and keep the conversion-detection/fence helpers consistent; add a 
high-source-sequence plain-target copy followed by a real column-update 
regression test.



-- 
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