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 17f700cd87 [core] Tag diff query at postpone table should check bucket
number (#9048)
17f700cd87 is described below
commit 17f700cd8728730d5cd6bab9d25a5c74220400a0
Author: yuzelin <[email protected]>
AuthorDate: Fri Aug 7 13:16:15 2026 +0800
[core] Tag diff query at postpone table should check bucket number (#9048)
---
.../snapshot/IncrementalDiffStartingScanner.java | 9 --
.../table/source/snapshot/SnapshotReaderImpl.java | 2 +
.../table/source/snapshot/TimeTravelUtil.java | 75 ++++++++--
.../apache/paimon/table/IncrementalTableTest.java | 156 +++++++++++++++++++++
.../spark/sql/TableValuedFunctionsTest.scala | 40 ++++++
5 files changed, 258 insertions(+), 24 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/IncrementalDiffStartingScanner.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/IncrementalDiffStartingScanner.java
index 3a42f16e2c..f7bb09e112 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/IncrementalDiffStartingScanner.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/IncrementalDiffStartingScanner.java
@@ -21,7 +21,6 @@ package org.apache.paimon.table.source.snapshot;
import org.apache.paimon.CoreOptions;
import org.apache.paimon.Snapshot;
import org.apache.paimon.manifest.PartitionEntry;
-import org.apache.paimon.schema.SchemaManager;
import org.apache.paimon.tag.Tag;
import org.apache.paimon.tag.TagPeriodHandler;
import org.apache.paimon.utils.Pair;
@@ -53,14 +52,6 @@ public class IncrementalDiffStartingScanner extends
AbstractStartingScanner {
this.start = start;
this.end = end;
this.startingSnapshotId = start.id();
-
- TimeTravelUtil.checkRescaleBucketForIncrementalDiffQuery(
- new SchemaManager(
- snapshotManager.fileIO(),
- snapshotManager.tablePath(),
- snapshotManager.branch()),
- start,
- end);
}
@Override
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/SnapshotReaderImpl.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/SnapshotReaderImpl.java
index 9826965039..cb338c8570 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/SnapshotReaderImpl.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/SnapshotReaderImpl.java
@@ -608,6 +608,8 @@ public class SnapshotReaderImpl implements SnapshotReader {
groupByPartFiles(plan.files(FileKind.ADD));
Map<BinaryRow, Map<Integer, List<ManifestEntry>>> beforeFiles =
groupByPartFiles(scan.withSnapshot(before).plan().files(FileKind.ADD));
+ TimeTravelUtil.checkRescaleBucketForIncrementalDiffQuery(
+ tableSchema, before, beforeFiles, plan.snapshot(), afterFiles);
return toIncrementalPlan(
false,
new LazyField<>(() -> before),
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java
index 786fecd0cb..47d7bc6d44 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java
@@ -20,8 +20,9 @@ package org.apache.paimon.table.source.snapshot;
import org.apache.paimon.CoreOptions;
import org.apache.paimon.Snapshot;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.manifest.ManifestEntry;
import org.apache.paimon.options.Options;
-import org.apache.paimon.schema.SchemaManager;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.utils.ChangelogManager;
@@ -38,8 +39,10 @@ import java.io.FileNotFoundException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
+import java.util.Map;
import java.util.Optional;
import java.util.TimeZone;
+import java.util.stream.Collectors;
import static org.apache.paimon.CoreOptions.SCAN_SNAPSHOT_ID;
import static org.apache.paimon.CoreOptions.SCAN_TAG_NAME;
@@ -251,24 +254,66 @@ public class TimeTravelUtil {
}
public static void checkRescaleBucketForIncrementalDiffQuery(
- SchemaManager schemaManager, Snapshot start, Snapshot end) {
- if (start.schemaId() != end.schemaId()) {
- int startBucketNumber = bucketNumber(schemaManager,
start.schemaId());
- int endBucketNumber = bucketNumber(schemaManager, end.schemaId());
- if (startBucketNumber != endBucketNumber) {
- throw new InconsistentTagBucketException(
- start.id(),
- end.id(),
- String.format(
- "The bucket number of two snapshots are
different (%s, %s), which is not supported in incremental diff query.",
- startBucketNumber, endBucketNumber));
+ TableSchema schema,
+ Snapshot start,
+ Map<BinaryRow, Map<Integer, List<ManifestEntry>>> startFiles,
+ Snapshot end,
+ Map<BinaryRow, Map<Integer, List<ManifestEntry>>> endFiles) {
+ if (schema.numBuckets() == -1) {
+ return;
+ }
+
+ for (Map.Entry<BinaryRow, Map<Integer, List<ManifestEntry>>> entry :
+ startFiles.entrySet()) {
+ Map<Integer, List<ManifestEntry>> endPartitionFiles =
endFiles.get(entry.getKey());
+ if (endPartitionFiles == null) {
+ continue;
+ }
+
+ Integer startPartitionBucketNumber =
+ realBucketNumbers(
+ entry.getValue().values().stream()
+ .flatMap(List::stream)
+ .collect(Collectors.toList()));
+ Integer endPartitionBucketNumber =
+ realBucketNumbers(
+ endPartitionFiles.values().stream()
+ .flatMap(List::stream)
+ .collect(Collectors.toList()));
+
+ if (startPartitionBucketNumber != null
+ && endPartitionBucketNumber != null
+ &&
startPartitionBucketNumber.equals(endPartitionBucketNumber)) {
+ continue;
}
+
+ throw new InconsistentTagBucketException(
+ start.id(),
+ end.id(),
+ String.format(
+ "The bucket number of two snapshots are different
(%s, %s), which is not supported in incremental diff query.",
+ startPartitionBucketNumber,
endPartitionBucketNumber));
}
}
- private static int bucketNumber(SchemaManager schemaManager, long
schemaId) {
- TableSchema schema = schemaManager.schema(schemaId);
- return CoreOptions.fromMap(schema.options()).bucket();
+ @Nullable
+ private static Integer realBucketNumbers(List<ManifestEntry> entries) {
+ Integer totalBuckets = null;
+ for (ManifestEntry entry : entries) {
+ if (entry.totalBuckets() >= 0) {
+ if (totalBuckets != null && totalBuckets !=
entry.totalBuckets()) {
+ throw new IllegalStateException(
+ "Partition "
+ + entry.partition()
+ + " has different totalBuckets "
+ + totalBuckets
+ + " and "
+ + entry.totalBuckets());
+ }
+ totalBuckets = entry.totalBuckets();
+ }
+ }
+ return totalBuckets;
}
/**
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/IncrementalTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/IncrementalTableTest.java
index 3a85110bec..51f245982e 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/IncrementalTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/IncrementalTableTest.java
@@ -31,8 +31,10 @@ 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.PostponeFixedBucketWriteBuilder;
import org.apache.paimon.table.sink.TableCommitImpl;
import org.apache.paimon.table.sink.TableWriteImpl;
+import
org.apache.paimon.table.source.snapshot.TimeTravelUtil.InconsistentTagBucketException;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowKind;
import org.apache.paimon.utils.Pair;
@@ -42,6 +44,7 @@ import org.apache.paimon.utils.TagManager;
import org.junit.jupiter.api.Test;
import java.time.LocalDateTime;
+import java.util.Collections;
import java.util.List;
import static org.apache.paimon.CoreOptions.INCREMENTAL_BETWEEN;
@@ -581,6 +584,159 @@ public class IncrementalTableTest extends TableTestBase {
.isEmpty();
}
+ @Test
+ public void testPostponeSameBucketNumberWithDifferentActiveBuckets()
throws Exception {
+ Identifier identifier = identifier("T");
+ Schema schema =
+ Schema.newBuilder()
+ .column("pk", DataTypes.INT())
+ .column("col1", DataTypes.INT())
+ .primaryKey("pk")
+ .option("bucket",
String.valueOf(BucketMode.POSTPONE_BUCKET))
+ .build();
+ catalog.createTable(identifier, schema, true);
+ FileStoreTable table = (FileStoreTable) catalog.getTable(identifier);
+
+ PostponeFixedBucketWriteBuilder builder =
table.newPostponeFixedBucketWriteBuilder();
+ try (TableWriteImpl<?> write = builder.newWrite();
+ BatchTableCommit commit = builder.newCommit()) {
+ write.writeAndReturn(GenericRow.of(1, 1), 0, 2);
+ commit.commit(write.prepareCommit());
+ }
+ table.createTag("TAG1", 1);
+
+ try (TableWriteImpl<?> write = builder.newWrite();
+ BatchTableCommit commit = builder.newCommit()) {
+ write.writeAndReturn(GenericRow.of(2, 2), 1, 2);
+ commit.commit(write.prepareCommit());
+ }
+ table.createTag("TAG2", 2);
+
+ assertThat(read(table, Pair.of(INCREMENTAL_BETWEEN, "TAG1,TAG2")))
+ .containsExactly(GenericRow.of(2, 2));
+ }
+
+ @Test
+ public void testPostponeBucketNumberChangedInIncrementalDiff() throws
Exception {
+ Identifier identifier = identifier("T");
+ Schema schema =
+ Schema.newBuilder()
+ .column("pk", DataTypes.INT())
+ .column("col1", DataTypes.INT())
+ .primaryKey("pk")
+ .option("bucket",
String.valueOf(BucketMode.POSTPONE_BUCKET))
+ .build();
+ catalog.createTable(identifier, schema, true);
+ FileStoreTable table = (FileStoreTable) catalog.getTable(identifier);
+
+ PostponeFixedBucketWriteBuilder builder =
table.newPostponeFixedBucketWriteBuilder();
+ try (TableWriteImpl<?> write = builder.newWrite();
+ BatchTableCommit commit = builder.newCommit()) {
+ write.writeAndReturn(GenericRow.of(1, 1), 0, 1);
+ commit.commit(write.prepareCommit());
+ }
+ table.createTag("TAG1", 1);
+
+ builder =
table.newPostponeFixedBucketWriteBuilder().withOverwrite(Collections.emptyMap());
+ try (TableWriteImpl<?> write = builder.newWrite();
+ BatchTableCommit commit = builder.newCommit()) {
+ write.writeAndReturn(GenericRow.of(1, 2), 0, 2);
+ commit.commit(write.prepareCommit());
+ }
+ table.createTag("TAG2", 2);
+
+ assertThatThrownBy(() -> read(table, Pair.of(INCREMENTAL_BETWEEN,
"TAG1,TAG2")))
+ .isInstanceOf(InconsistentTagBucketException.class)
+ .hasMessageContaining(
+ "The bucket number of two snapshots are different (1,
2), "
+ + "which is not supported in incremental diff
query.");
+ }
+
+ @Test
+ public void testPostponeBucketNumberChangedInLaterPartition() throws
Exception {
+ Identifier identifier = identifier("T");
+ Schema schema =
+ Schema.newBuilder()
+ .column("pt", DataTypes.INT())
+ .column("pk", DataTypes.INT())
+ .column("col1", DataTypes.INT())
+ .partitionKeys("pt")
+ .primaryKey("pk", "pt")
+ .option("bucket",
String.valueOf(BucketMode.POSTPONE_BUCKET))
+ .build();
+ catalog.createTable(identifier, schema, true);
+ FileStoreTable table = (FileStoreTable) catalog.getTable(identifier);
+
+ PostponeFixedBucketWriteBuilder builder =
table.newPostponeFixedBucketWriteBuilder();
+ try (TableWriteImpl<?> write = builder.newWrite();
+ BatchTableCommit commit = builder.newCommit()) {
+ write.writeAndReturn(GenericRow.of(2, 2, 2), 0, 1);
+ write.writeAndReturn(GenericRow.of(1, 1, 1), 0, 1);
+ commit.commit(write.prepareCommit());
+ }
+ table.createTag("TAG1", 1);
+
+ builder =
table.newPostponeFixedBucketWriteBuilder().withOverwrite(Collections.emptyMap());
+ try (TableWriteImpl<?> write = builder.newWrite();
+ BatchTableCommit commit = builder.newCommit()) {
+ write.writeAndReturn(GenericRow.of(2, 2, 3), 0, 1);
+ write.writeAndReturn(GenericRow.of(1, 1, 2), 0, 2);
+ commit.commit(write.prepareCommit());
+ }
+ table.createTag("TAG2", 2);
+
+ assertThatThrownBy(() -> read(table, Pair.of(INCREMENTAL_BETWEEN,
"TAG1,TAG2")))
+ .isInstanceOf(InconsistentTagBucketException.class)
+ .hasMessageContaining(
+ "The bucket number of two snapshots are different (1,
2), "
+ + "which is not supported in incremental diff
query.");
+ }
+
+ @Test
+ public void testPostponeDifferentBucketNumbersForDifferentPartitions()
throws Exception {
+ Identifier identifier = identifier("T");
+ Schema schema =
+ Schema.newBuilder()
+ .column("pt", DataTypes.INT())
+ .column("pk", DataTypes.INT())
+ .column("col1", DataTypes.INT())
+ .partitionKeys("pt")
+ .primaryKey("pk", "pt")
+ .option("bucket",
String.valueOf(BucketMode.POSTPONE_BUCKET))
+ .build();
+ catalog.createTable(identifier, schema, true);
+ FileStoreTable table = (FileStoreTable) catalog.getTable(identifier);
+
+ PostponeFixedBucketWriteBuilder builder =
table.newPostponeFixedBucketWriteBuilder();
+ try (TableWriteImpl<?> write = builder.newWrite();
+ BatchTableCommit commit = builder.newCommit()) {
+ write.writeAndReturn(GenericRow.of(1, 1, 1), 0, 1);
+ commit.commit(write.prepareCommit());
+ }
+ table.createTag("TAG1", 1);
+
+ try (TableWriteImpl<?> write = builder.newWrite();
+ BatchTableCommit commit = builder.newCommit()) {
+ write.writeAndReturn(GenericRow.of(2, 1, 1), 1, 2);
+ commit.commit(write.prepareCommit());
+ }
+ table.createTag("TAG2", 2);
+
+ // test snapshot expiration won't affect tag diff query
+ try (TableWriteImpl<?> write = builder.newWrite();
+ BatchTableCommit commit = builder.newCommit()) {
+ write.writeAndReturn(GenericRow.of(3, 1, 1), 0, 1);
+ commit.commit(write.prepareCommit());
+ }
+ table.newExpireSnapshots()
+
.config(ExpireConfig.builder().snapshotRetainMax(1).snapshotRetainMin(1).build())
+ .expire();
+ assertThat(table.snapshotManager().snapshotCount()).isEqualTo(1);
+
+ assertThat(read(table, Pair.of(INCREMENTAL_BETWEEN, "TAG1,TAG2")))
+ .containsExactly(GenericRow.of(2, 1, 1));
+ }
+
private static long utcMills(String timestamp) {
return
Timestamp.fromLocalDateTime(LocalDateTime.parse(timestamp)).getMillisecond();
}
diff --git
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/TableValuedFunctionsTest.scala
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/TableValuedFunctionsTest.scala
index e46a00358d..301a175611 100644
---
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/TableValuedFunctionsTest.scala
+++
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/TableValuedFunctionsTest.scala
@@ -754,6 +754,46 @@ class TableValuedFunctionsTest extends PaimonHiveTestBase
with AdaptiveSparkPlan
}
}
+ test("Table Valued Functions: incremental query with inconsistent postpone
bucket") {
+ withTable("t") {
+ sql("""
+ |CREATE TABLE t (a INT, b INT) USING paimon
+ |TBLPROPERTIES ('primary-key'='a', 'bucket' = '-2')
+ |""".stripMargin)
+
+ val table = loadTable("t")
+ var builder = table.newPostponeFixedBucketWriteBuilder()
+ var write = builder.newWrite()
+ var commit = builder.newCommit()
+ try {
+ write.writeAndReturn(GenericRow.of(1, 11), 0, 1)
+ commit.commit(write.prepareCommit())
+ } finally {
+ write.close()
+ commit.close()
+ }
+ table.createTag("2024-01-01", 1)
+
+ builder =
table.newPostponeFixedBucketWriteBuilder().withOverwrite(Collections.emptyMap())
+ write = builder.newWrite()
+ commit = builder.newCommit()
+ try {
+ write.writeAndReturn(GenericRow.of(1, 22), 0, 2)
+ write.writeAndReturn(GenericRow.of(2, 22), 1, 2)
+ commit.commit(write.prepareCommit())
+ } finally {
+ write.close()
+ commit.close()
+ }
+ table.createTag("2024-01-02", 2)
+
+ checkAnswer(
+ sql(
+ "SELECT * FROM paimon_incremental_query('t', '2024-01-01',
'2024-01-02') ORDER BY a, b"),
+ Seq(Row(1, 22), Row(2, 22)))
+ }
+ }
+
test("Table Valued Functions: incremental query with delete after minor
compact") {
withTable("t") {
sql("""