zhuxiangyi commented on code in PR #10098:
URL: https://github.com/apache/paimon/pull/10098#discussion_r4095471478
##########
paimon-core/src/main/java/org/apache/paimon/table/AbstractFileStoreTable.java:
##########
@@ -648,6 +654,27 @@ public void rollbackTo(String tagName) {
}
}
+ /**
+ * A snapshot committed before {@code sys.enable_data_evolution} converted
the table holds files
+ * without a first row id. Making such a snapshot the latest again while
the schema still has
+ * row tracking enabled would leave the table unreadable as a
data-evolution table, so refuse
+ * it; roll the schema back first if the conversion really has to be
undone.
+ */
+ private void checkRollbackKeepsRowTracking(Snapshot target) {
Review Comment:
Fixed in 35110ec00. Both guards now go by the latest persisted schema
instead of the options of the table object:
- `AbstractFileStoreTable.checkRollbackKeepsRowTracking`, used by
`rollbackTo(long)` and `rollbackTo(String)`, reads `schemaManager().latest()`.
- `FileStoreCommitImpl.rollbackToAsLatest` reads
`schemaManager.latestOrThrow(...)`.
Regression test:
`RowTrackingEnableGuardTest#testRollbackThroughTableLoadedBeforeRowTrackingIsRefused`.
It keeps the table object loaded before row tracking was enabled, enables row
tracking through another handle, and checks that snapshot, tag and as-latest
rollback are all refused, that the latest snapshot and the tag are unchanged,
and that a target committed with row tracking still rolls back through the same
object. Reverting either guard makes the test fail.
##########
paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionEnabler.java:
##########
@@ -0,0 +1,485 @@
+/*
+ * 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.ManifestEntry;
+import org.apache.paimon.manifest.ManifestFile;
+import org.apache.paimon.manifest.ManifestFileMeta;
+import org.apache.paimon.manifest.ManifestList;
+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.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.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+
+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 three 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.
+ * <li>Assign ids to any file that a writer on the previous schema committed
in between, now that
+ * every later commit assigns ids on its own. A writer that loaded the
table before the switch
+ * is refused from then on, see {@code FileStoreCommitImpl}.
+ * </ol>
+ *
+ * <p>The procedure is idempotent: on a table that already has data evolution
enabled it only
+ * assigns ids to files that still lack one, and reports {@code skipped} when
there are none.
+ */
+public class DataEvolutionEnabler {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(DataEvolutionEnabler.class);
+ private static final String COMMIT_USER_PREFIX = "enable-data-evolution";
+ private static final int MAX_REPAIR_ROUNDS = 5;
+
+ private final Catalog catalog;
+ private final Identifier identifier;
+ private final Runnable beforeRowIdCommit;
+ private final Runnable beforeSchemaChange;
+
+ public DataEvolutionEnabler(Catalog catalog, Identifier identifier) {
+ this(catalog, identifier, () -> {}, () -> {});
+ }
+
+ /** Hooks for tests to inject concurrent activity between the steps. */
+ DataEvolutionEnabler(
+ Catalog catalog,
+ Identifier identifier,
+ Runnable beforeRowIdCommit,
+ Runnable beforeSchemaChange) {
+ this.catalog = catalog;
+ this.identifier = identifier;
+ this.beforeRowIdCommit = beforeRowIdCommit;
+ this.beforeSchemaChange = beforeSchemaChange;
+ }
+
+ /** 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.files.isEmpty()) {
+ return Result.skipped(
+ schemaBefore, snapshotBefore, "data evolution is already
enabled");
+ }
+ if (dryRun) {
+ return Result.dryRun(schemaBefore, snapshotBefore, enabled,
planned);
+ }
+
+ long assignedFiles = 0;
+ long assignedRows = 0;
+ // Even when this run assigns nothing, report the id the table is at:
a previous run may
+ // have committed the ids and failed before the schema change.
+ Long nextRowId = planned.snapshot == null ? null : planned.nextRowId;
+ if (!planned.files.isEmpty()) {
+ Committed committed = assignRowIdsWithRetry(table, planned);
+ assignedFiles += committed.assignment.files.size();
+ assignedRows += committed.assignment.rowCount;
+ nextRowId = committed.assignment.nextRowId;
+ }
+
+ if (!enabled) {
+ beforeSchemaChange.run();
+ catalog.alterTable(identifier, SchemaChange.enableDataEvolution(),
false);
+ table = loadTable();
+ checkState(
+ table.coreOptions().rowTrackingEnabled()
+ && table.coreOptions().dataEvolutionEnabled(),
+ "Schema change did not enable data evolution on table %s.",
+ identifier.getFullName());
+
+ // Repair what a writer on the previous schema committed between
the two steps.
+ for (int round = 0; round < MAX_REPAIR_ROUNDS; round++) {
+ Assignment remaining = plan(table);
+ if (remaining.files.isEmpty()) {
+ break;
+ }
+ LOG.info(
+ "Assigning row ids to {} file(s) committed to table {}
while data evolution was being enabled.",
+ remaining.files.size(),
+ identifier.getFullName());
+ Committed committed = assignRowIdsWithRetry(table, remaining);
+ assignedFiles += committed.assignment.files.size();
+ assignedRows += committed.assignment.rowCount;
+ nextRowId = committed.assignment.nextRowId;
+ }
+ checkState(
+ plan(table).files.isEmpty(),
+ "Table %s still has data files without a row id after %s
repair rounds; "
+ + "stop the writers that predate the schema change
and run the "
+ + "procedure again.",
+ identifier.getFullName(),
+ MAX_REPAIR_ROUNDS);
+ }
+
+ return new Result(
+ schemaBefore,
+ table.schema().id(),
+ snapshotBefore,
+ table.snapshotManager().latestSnapshotId(),
+ assignedFiles,
+ assignedRows,
+ nextRowId,
+ false,
+ false,
+ null);
+ }
+
+ 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 a first row id for every live data file of the latest snapshot
that has none. */
+ 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<>();
+ for (ManifestEntry entry : live.values()) {
+ if (entry.kind() == FileKind.ADD && entry.file().firstRowId() ==
null) {
+ withoutRowId.add(entry);
+ }
+ }
+ long start = latest.nextRowId() == null ? 0L : latest.nextRowId();
+ if (withoutRowId.isEmpty()) {
+ return Assignment.empty(latest, start);
+ }
+
+ // Contiguous per partition, partitions in order: the layout
reassign_row_id produces.
+ // Within a partition the files keep the order they have in the
manifests, which is the
+ // order they were committed in (the sort is stable).
+ RecordComparator partitionComparator =
+ CodeGenUtils.newRecordComparator(
+ table.schema().logicalPartitionType().getFieldTypes());
+ withoutRowId.sort(
+ (left, right) -> partitionComparator.compare(left.partition(),
right.partition()));
+
+ Map<FileEntry.Identifier, Long> firstRowIds = new HashMap<>();
+ long next = start;
+ long rowCount = 0;
+ for (ManifestEntry entry : withoutRowId) {
+ firstRowIds.put(entry.identifier(), next);
+ next += entry.file().rowCount();
+ rowCount += entry.file().rowCount();
+ }
+ return new Assignment(latest, manifests, withoutRowId, firstRowIds,
rowCount, next);
+ }
+
+ private Committed assignRowIdsWithRetry(FileStoreTable table, Assignment
initial)
+ throws Exception {
+ CoreOptions options = table.coreOptions();
+ RetryWaiter retryWaiter =
+ new RetryWaiter(options.commitMinRetryWait(),
options.commitMaxRetryWait());
+ long startMillis = System.currentTimeMillis();
+ Assignment assignment = initial;
+ int retryCount = 0;
+ while (true) {
+ if (commitAssignment(table, assignment)) {
+ return new Committed(assignment);
+ }
+ if (System.currentTimeMillis() - startMillis >
options.commitTimeout()
+ || retryCount >= options.commitMaxRetries()) {
+ throw new RuntimeException(
+ String.format(
+ "Failed to assign row ids to table %s after %s
millis and %s "
+ + "retries because newer snapshots
kept being committed.",
+ identifier.getFullName(),
+ System.currentTimeMillis() - startMillis,
+ retryCount));
+ }
+ retryWaiter.retryWait(retryCount);
+ retryCount++;
+ // Another commit landed: plan again from the new latest snapshot.
Files that already
+ // received an id in it (written by a writer on the new schema)
keep it.
+ assignment = plan(table);
+ if (assignment.files.isEmpty()) {
+ return new Committed(assignment);
+ }
+ LOG.info(
+ "Retrying row id assignment for table {} on snapshot {}
({}/{}).",
+ identifier.getFullName(),
+ assignment.snapshot.id(),
+ retryCount,
+ options.commitMaxRetries());
+ }
+ }
+
+ /**
+ * Rewrites the manifests holding the planned files and commits them,
referencing the table's
+ * current schema. Returns false when the snapshot moved on in the
meantime.
+ */
+ private boolean commitAssignment(FileStoreTable table, Assignment
assignment) {
+ ManifestFile manifestFile =
table.store().manifestFileFactory().create();
+ ManifestList manifestList =
table.store().manifestListFactory().create();
+
+ List<ManifestFileMeta> baseManifests = new ArrayList<>();
+ for (ManifestFileMeta manifest : assignment.manifests) {
+ List<ManifestEntry> entries =
+ manifestFile.read(manifest.fileName(),
manifest.fileSize());
+ List<ManifestEntry> rewritten = new ArrayList<>(entries.size());
+ boolean changed = false;
+ for (ManifestEntry entry : entries) {
+ Long firstRowId =
assignment.firstRowIds.get(entry.identifier());
+ if (firstRowId != null && entry.file().firstRowId() == null) {
+ rewritten.add(entry.assignFirstRowId(firstRowId));
+ changed = true;
+ } else {
+ rewritten.add(entry);
+ }
+ }
+ if (changed) {
+ baseManifests.addAll(manifestFile.write(rewritten));
+ } else {
+ baseManifests.add(manifest);
+ }
+ }
+
+ Pair<String, Long> baseManifestList =
manifestList.write(baseManifests);
+ Pair<String, Long> deltaManifestList =
manifestList.write(Collections.emptyList());
+ String commitUser = COMMIT_USER_PREFIX + "-" + UUID.randomUUID();
+ try (FileStoreCommitImpl commit =
+ (FileStoreCommitImpl) table.store().newCommit(commitUser,
table)) {
+ beforeRowIdCommit.run();
+ return commit.replaceManifestList(
+ assignment.snapshot,
+ table.schema().id(),
Review Comment:
Fixed in 37ff2bdfb. `commitAssignment` now reads the schema id with
`table.schemaManager().latestOrThrow(...)` inside every commit attempt, after
the `beforeRowIdCommit` hook and right before `replaceManifestList`, the same
way the normal commit path does. The cached `table.schema()` is no longer used
for the snapshot.
Regression test:
`DataEvolutionEnablerTest#testRowIdCommitReferencesTheLatestSchema`. It commits
`ADD COLUMN c` from the hook, between planning and the replacement, and asserts
that the row id snapshot references the new schema, that time travel to it
exposes `c`, and that the final schema has `c` with data evolution enabled.
Reverting the change makes the test fail.
##########
paimon-api/src/main/java/org/apache/paimon/schema/SchemaChange.java:
##########
@@ -171,6 +174,16 @@ static SchemaChange dropPrimaryKey() {
return new DropPrimaryKey();
}
+ /**
+ * Enables {@code row-tracking.enabled} and {@code data-evolution.enabled}
on a table that
+ * already has snapshots. Both options are immutable for {@code ALTER
TABLE}; this change is
+ * issued by the {@code sys.enable_data_evolution} procedure, which first
assigns a first row id
+ * to every existing data file.
+ */
+ static SchemaChange enableDataEvolution() {
Review Comment:
Fixed in 7eb1237cc (REST) and 37ff2bdfb (the other catalogs).
**REST:** the action is no longer part of the protocol.
- It is removed from the `SchemaChange` JSON subtypes and from
`rest-catalog-open-api.yaml`, which is back to master's version, so a REST
server cannot parse it.
- `RESTCatalog.alterTable` rejects it before sending, with the same message
the procedure gives.
- Tests: `RESTApiJsonTest` asserts that `{"action":"enableDataEvolution"}`
does not parse. `DataEvolutionEnablerTest#testRejectsRESTCatalog` now also
sends a direct `alterTable` on a REST table that has data: it is refused and
the options stay unchanged.
**Other catalogs:** `FileSystemSchemaManager.commitChanges` accepts
`EnableDataEvolution` on a table without row tracking only while every live
data file has a row id. That means the table has no snapshot yet, or its latest
snapshot is a row id commit of the procedure.
- Such a commit is marked with the snapshot property
`data-evolution.row-ids-assigned-snapshot-id`, whose value is the snapshot's
own id. A later snapshot that copies its base's properties therefore does not
carry the mark, the same pattern as the reassign plan in #10008.
- `SchemaManagerUtils` still skips the immutable-option check for this
change, but only after this precondition has been checked.
- Tests:
- `testSchemaChangeAloneIsRefusedWhileFilesHaveNoRowId`: a table with data
is refused and keeps its options.
- `testSchemaChangeAloneIsAcceptedOnTableWithoutSnapshot`
- `testSchemaChangeNeedsTheLatestSnapshotToBeAMarkedOne`: a snapshot that
only copied the mark is refused.
--
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]