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]
