This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new f7e9614a98 [flink] Fix batch database compaction for Data Evolution 
tables (#10110)
f7e9614a98 is described below

commit f7e9614a981d061c6774a29697ed6880957d4d2a
Author: tzphh <[email protected]>
AuthorDate: Thu Sep 24 13:10:02 2026 +0800

    [flink] Fix batch database compaction for Data Evolution tables (#10110)
---
 .../paimon/flink/action/CompactDatabaseAction.java |  34 +++-
 .../flink/action/CompactDatabaseActionITCase.java  | 214 +++++++++++++++++++++
 2 files changed, 246 insertions(+), 2 deletions(-)

diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/CompactDatabaseAction.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/CompactDatabaseAction.java
index 3000c319b5..47323c68dc 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/CompactDatabaseAction.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/CompactDatabaseAction.java
@@ -25,6 +25,7 @@ import org.apache.paimon.catalog.Identifier;
 import org.apache.paimon.flink.FlinkConnectorOptions;
 import 
org.apache.paimon.flink.FlinkConnectorOptions.CompactionBucketDistributionStrategy;
 import org.apache.paimon.flink.compact.AppendTableCompact;
+import org.apache.paimon.flink.compact.DataEvolutionTableCompact;
 import org.apache.paimon.flink.sink.BucketsRowChannelComputer;
 import org.apache.paimon.flink.sink.CombinedTableCompactorSink;
 import org.apache.paimon.flink.sink.CompactorSinkBuilder;
@@ -48,6 +49,7 @@ import org.slf4j.LoggerFactory;
 import javax.annotation.Nullable;
 
 import java.time.Duration;
+import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -147,7 +149,7 @@ public class CompactDatabaseAction extends ActionBase {
         }
     }
 
-    private void buildForDividedMode() {
+    private void collectTables() {
         try {
             List<String> databases = catalog.listDatabases();
             for (String databaseName : databases) {
@@ -181,7 +183,10 @@ public class CompactDatabaseAction extends ActionBase {
         } catch (Catalog.DatabaseNotExistException | 
Catalog.TableNotExistException e) {
             throw new RuntimeException(e);
         }
+    }
 
+    private void buildForDividedMode() {
+        collectTables();
         Preconditions.checkState(
                 !tableMap.isEmpty(),
                 "no tables to be compacted. possible cause is that there are 
no tables detected after pattern matching");
@@ -205,12 +210,30 @@ public class CompactDatabaseAction extends ActionBase {
     }
 
     private void buildForCombinedMode() {
+        Pattern combinedIncludingPattern = includingPattern;
+        if (!isStreaming) {
+            collectTables();
+            List<String> combinedTables = new ArrayList<>();
+            for (Map.Entry<String, FileStoreTable> entry : 
tableMap.entrySet()) {
+                if (entry.getValue().coreOptions().dataEvolutionEnabled()) {
+                    buildForUnawareBucketCompaction(env, entry.getKey(), 
entry.getValue());
+                } else {
+                    combinedTables.add(Pattern.quote(entry.getKey()));
+                }
+            }
+            if (combinedTables.isEmpty() && !tableMap.isEmpty()) {
+                return;
+            }
+            // Batch table selection is fixed when building the job. 
Data-evolution tables
+            // already have dedicated compactors and must not enter the shared 
append path.
+            combinedIncludingPattern = Pattern.compile(String.join("|", 
combinedTables));
+        }
 
         CombinedTableCompactorSourceBuilder sourceBuilder =
                 new CombinedTableCompactorSourceBuilder(
                                 catalogLoader(),
                                 databasePattern,
-                                includingPattern,
+                                combinedIncludingPattern,
                                 excludingPattern,
                                 tableOptions
                                         
.get(CoreOptions.CONTINUOUS_DISCOVERY_INTERVAL)
@@ -285,6 +308,13 @@ public class CompactDatabaseAction extends ActionBase {
 
     private void buildForUnawareBucketCompaction(
             StreamExecutionEnvironment env, String fullName, FileStoreTable 
table) {
+        if (table.coreOptions().dataEvolutionEnabled()) {
+            Preconditions.checkArgument(
+                    !isStreaming, "Data evolution table compact only supports 
batch mode yet.");
+            new DataEvolutionTableCompact(env, fullName, table).build();
+            return;
+        }
+
         AppendTableCompact unawareBucketCompactionTopoBuilder =
                 new AppendTableCompact(env, fullName, table);
 
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/CompactDatabaseActionITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/CompactDatabaseActionITCase.java
index acb0267440..b1534c9fae 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/CompactDatabaseActionITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/CompactDatabaseActionITCase.java
@@ -24,8 +24,14 @@ import org.apache.paimon.catalog.Identifier;
 import org.apache.paimon.data.BinaryString;
 import org.apache.paimon.data.GenericRow;
 import org.apache.paimon.flink.FlinkConnectorOptions;
+import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.schema.Schema;
 import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.sink.BatchTableCommit;
+import org.apache.paimon.table.sink.BatchTableWrite;
+import org.apache.paimon.table.sink.BatchWriteBuilder;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.CommitMessageImpl;
 import org.apache.paimon.table.sink.StreamTableCommit;
 import org.apache.paimon.table.sink.StreamTableWrite;
 import org.apache.paimon.table.sink.StreamWriteBuilder;
@@ -58,10 +64,12 @@ import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.ThreadLocalRandom;
+import java.util.regex.Pattern;
 import java.util.stream.Stream;
 
 import static org.apache.paimon.utils.CommonTestUtils.waitUtil;
 import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
 /** IT cases for {@link CompactDatabaseAction}. */
 public class CompactDatabaseActionITCase extends CompactActionITCaseBase {
@@ -297,6 +305,212 @@ public class CompactDatabaseActionITCase extends 
CompactActionITCaseBase {
         }
     }
 
+    @ParameterizedTest
+    @MethodSource("testData")
+    @Timeout(120)
+    public void testDataEvolutionTableCompact(String mode, String invoker) 
throws Exception {
+        List<FileStoreTable> evolutionTables =
+                Arrays.asList(
+                        createDataEvolutionTable(), 
createDataEvolutionTable("evolution_other"));
+        for (FileStoreTable table : evolutionTables) {
+            writeDataEvolutionRecords(table);
+        }
+        FileStoreTable table = evolutionTables.get(0);
+        // Compact an ordinary append table in the same job to exercise 
mixed-table routing.
+        Map<String, String> options = new HashMap<>();
+        options.put(CoreOptions.BUCKET.key(), "-1");
+        options.put(CoreOptions.WRITE_ONLY.key(), "true");
+        options.put(CoreOptions.COMPACTION_MIN_FILE_NUM.key(), "2");
+        FileStoreTable ordinaryTable =
+                createFileStoreTable(
+                        "ordinary_table",
+                        table.rowType(),
+                        Collections.emptyList(),
+                        Collections.emptyList(),
+                        Collections.emptyList(),
+                        options);
+        for (int i = 1; i <= 2; i++) {
+            BatchWriteBuilder builder = ordinaryTable.newBatchWriteBuilder();
+            try (BatchTableWrite write = builder.newWrite();
+                    BatchTableCommit commit = builder.newCommit()) {
+                write.write(rowData(i, i * 10));
+                commit.commit(write.prepareCommit());
+            }
+        }
+
+        assertThat(getResult(table.newRead(), table.newScan().plan().splits(), 
table.rowType()))
+                .containsExactlyInAnyOrder("+I[1, 30]", "+I[2, 60]");
+
+        switch (invoker) {
+            case "action":
+                StreamExecutionEnvironment env =
+                        
streamExecutionEnvironmentBuilder().batchMode().build();
+                createAction(
+                                CompactDatabaseAction.class,
+                                "compact_database",
+                                "--warehouse",
+                                warehouse,
+                                "--mode",
+                                mode)
+                        .withStreamExecutionEnvironment(env)
+                        .build();
+                env.execute();
+                break;
+            case "procedure_indexed":
+                executeSQL(String.format("CALL sys.compact_database('', 
'%s')", mode), false, true)
+                        .close();
+                break;
+            case "procedure_named":
+                executeSQL(
+                                String.format("CALL sys.compact_database(mode 
=> '%s')", mode),
+                                false,
+                                true)
+                        .close();
+                break;
+            default:
+                throw new UnsupportedOperationException(invoker);
+        }
+
+        for (FileStoreTable evolutionTable : evolutionTables) {
+            checkLatestSnapshot(evolutionTable, 4, 
Snapshot.CommitKind.COMPACT);
+            List<DataSplit> splits = 
evolutionTable.newSnapshotReader().read().dataSplits();
+            assertThat(splits).hasSize(1);
+            assertThat(splits.get(0).dataFiles()).hasSize(1);
+            DataFileMeta compacted = splits.get(0).dataFiles().get(0);
+            assertThat(compacted.nonNullFirstRowId()).isZero();
+            assertThat(compacted.rowCount()).isEqualTo(2);
+            assertThat(
+                            getResult(
+                                    evolutionTable.newRead(),
+                                    evolutionTable.newScan().plan().splits(),
+                                    evolutionTable.rowType()))
+                    .containsExactlyInAnyOrder("+I[1, 30]", "+I[2, 60]");
+        }
+
+        checkLatestSnapshot(ordinaryTable, 3, Snapshot.CommitKind.COMPACT);
+        
assertThat(ordinaryTable.newSnapshotReader().read().dataSplits().get(0).dataFiles())
+                .hasSize(1);
+        assertThat(
+                        getResult(
+                                ordinaryTable.newRead(),
+                                ordinaryTable.newScan().plan().splits(),
+                                ordinaryTable.rowType()))
+                .containsExactlyInAnyOrder("+I[1, 10]", "+I[2, 20]");
+    }
+
+    @ParameterizedTest
+    @ValueSource(booleans = {false, true})
+    public void testDataEvolutionOnlyCombinedCompact(boolean empty) throws 
Exception {
+        FileStoreTable table = createDataEvolutionTable();
+        if (!empty) {
+            writeDataEvolutionRecords(table);
+        }
+        StreamExecutionEnvironment env = 
streamExecutionEnvironmentBuilder().batchMode().build();
+        new CompactDatabaseAction(Collections.singletonMap("warehouse", 
warehouse))
+                .withDatabaseCompactMode("combined")
+                // Data Evolution uses the single-table compactor without 
idle-time filtering.
+                .withPartitionIdleTime(Duration.ofDays(1))
+                .withStreamExecutionEnvironment(env)
+                .build();
+        env.execute();
+        if (empty) {
+            assertThat(table.snapshotManager().latestSnapshotId()).isNull();
+        } else {
+            checkLatestSnapshot(table, 4, Snapshot.CommitKind.COMPACT);
+            assertThat(getResult(table.newRead(), 
table.newScan().plan().splits(), table.rowType()))
+                    .containsExactlyInAnyOrder("+I[1, 30]", "+I[2, 60]");
+        }
+    }
+
+    @ParameterizedTest
+    @ValueSource(booleans = {false, true})
+    public void testDataEvolutionCombinedTableFilter(boolean exclude) throws 
Exception {
+        FileStoreTable selected = createDataEvolutionTable();
+        FileStoreTable excluded = 
createDataEvolutionTable("excluded_evolution");
+        writeDataEvolutionRecords(selected);
+        writeDataEvolutionRecords(excluded);
+        StreamExecutionEnvironment env = 
streamExecutionEnvironmentBuilder().batchMode().build();
+        CompactDatabaseAction action =
+                new 
CompactDatabaseAction(Collections.singletonMap("warehouse", warehouse))
+                        .includingDatabases(Pattern.quote(database))
+                        .withDatabaseCompactMode("combined");
+        if (exclude) {
+            action.excludingTables(Pattern.quote(database + 
".excluded_evolution"));
+        } else {
+            action.includingTables(Pattern.quote(database + "." + tableName));
+        }
+        action.withStreamExecutionEnvironment(env).build();
+        env.execute();
+        checkLatestSnapshot(selected, 4, Snapshot.CommitKind.COMPACT);
+        checkLatestSnapshot(excluded, 3, Snapshot.CommitKind.APPEND);
+    }
+
+    @Test
+    public void testDataEvolutionStreamingCompactRejected() throws Exception {
+        createDataEvolutionTable();
+        CompactDatabaseAction action =
+                new 
CompactDatabaseAction(Collections.singletonMap("warehouse", warehouse));
+        action.withStreamExecutionEnvironment(
+                streamExecutionEnvironmentBuilder().streamingMode().build());
+        assertThatThrownBy(action::build)
+                .isInstanceOf(IllegalArgumentException.class)
+                .hasMessageContaining("Data evolution table compact only 
supports batch mode");
+    }
+
+    private void writeDataEvolutionRecords(FileStoreTable table) throws 
Exception {
+        BatchWriteBuilder builder = table.newBatchWriteBuilder();
+        try (BatchTableWrite write = builder.newWrite();
+                BatchTableCommit commit = builder.newCommit()) {
+            write.write(rowData(1, 10));
+            write.write(rowData(2, 20));
+            commit.commit(write.prepareCommit());
+        }
+
+        for (int multiplier : new int[] {2, 3}) {
+            builder = table.newBatchWriteBuilder();
+            try (BatchTableWrite write =
+                            builder.newWrite()
+                                    .withWriteType(
+                                            table.rowType()
+                                                    
.project(Collections.singletonList("v")));
+                    BatchTableCommit commit = builder.newCommit()) {
+                write.write(rowData(10 * multiplier));
+                write.write(rowData(20 * multiplier));
+                List<CommitMessage> messages = write.prepareCommit();
+                assertThat(messages).hasSize(1);
+                List<DataFileMeta> files =
+                        ((CommitMessageImpl) 
messages.get(0)).newFilesIncrement().newFiles();
+                assertThat(files).hasSize(1);
+                assertThat(files.get(0).writeCols()).containsExactly("v");
+                files.replaceAll(file -> file.assignFirstRowId(0L));
+                commit.commit(messages);
+            }
+        }
+    }
+
+    private FileStoreTable createDataEvolutionTable() throws Exception {
+        return createDataEvolutionTable(tableName);
+    }
+
+    private FileStoreTable createDataEvolutionTable(String name) throws 
Exception {
+        Map<String, String> options = new HashMap<>();
+        options.put(CoreOptions.BUCKET.key(), "-1");
+        options.put(CoreOptions.DATA_EVOLUTION_ENABLED.key(), "true");
+        options.put(CoreOptions.ROW_TRACKING_ENABLED.key(), "true");
+        options.put(CoreOptions.WRITE_ONLY.key(), "true");
+        options.put(CoreOptions.FILE_FORMAT.key(), "parquet");
+        options.put(CoreOptions.COMPACTION_MIN_FILE_NUM.key(), "2");
+        return createFileStoreTable(
+                name,
+                RowType.of(
+                        new DataType[] {DataTypes.INT().notNull(), 
DataTypes.INT()},
+                        new String[] {"id", "v"}),
+                Collections.emptyList(),
+                Collections.emptyList(),
+                Collections.emptyList(),
+                options);
+    }
+
     @ParameterizedTest(name = "mode = {0}, invoker = {1}")
     @MethodSource("testData")
     public void testStreamingCompact(String mode, String invoker) throws 
Exception {

Reply via email to