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 f52dde235a [core] Fix query auth for fallback reads and index
bootstrap (#8770)
f52dde235a is described below
commit f52dde235aa573e3701e6461fd6d12806be5dd17
Author: umi <[email protected]>
AuthorDate: Thu Jul 23 18:41:39 2026 +0800
[core] Fix query auth for fallback reads and index bootstrap (#8770)
---
.../paimon/crosspartition/IndexBootstrap.java | 21 ++-
.../paimon/table/FallbackReadFileStoreTable.java | 22 ++-
.../paimon/crosspartition/IndexBootstrapTest.java | 183 ++++++++++++++++++++-
.../table/FallbackReadFileStoreTableTest.java | 82 +++++++++
4 files changed, 293 insertions(+), 15 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/crosspartition/IndexBootstrap.java
b/paimon-core/src/main/java/org/apache/paimon/crosspartition/IndexBootstrap.java
index 24166243b2..8f9d99edd6 100644
---
a/paimon-core/src/main/java/org/apache/paimon/crosspartition/IndexBootstrap.java
+++
b/paimon-core/src/main/java/org/apache/paimon/crosspartition/IndexBootstrap.java
@@ -24,6 +24,7 @@ import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.JoinedRow;
import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.options.Options;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.FileStoreTable;
@@ -41,12 +42,12 @@ import java.io.IOException;
import java.io.Serializable;
import java.time.Duration;
import java.util.ArrayList;
-import java.util.Collections;
import java.util.List;
import java.util.function.Consumer;
import java.util.stream.Collectors;
import java.util.stream.Stream;
+import static org.apache.paimon.CoreOptions.QUERY_AUTH_ENABLED;
import static org.apache.paimon.CoreOptions.SCAN_MODE;
import static org.apache.paimon.CoreOptions.StartupMode.LATEST;
import static org.apache.paimon.io.SplitsParallelReadUtil.parallelExecute;
@@ -80,11 +81,12 @@ public class IndexBootstrap implements Serializable {
.mapToInt(Integer::intValue)
.toArray();
- // force using the latest scan mode
+ // Force using the latest scan mode and bypass query auth for this
internal index read.
+ Options bootstrapOptions = new Options();
+ bootstrapOptions.set(SCAN_MODE, LATEST);
+ bootstrapOptions.set(QUERY_AUTH_ENABLED, false);
ReadBuilder readBuilder =
- table.copy(Collections.singletonMap(SCAN_MODE.key(),
LATEST.toString()))
- .newReadBuilder()
- .withProjection(keyProjection);
+
table.copy(bootstrapOptions.toMap()).newReadBuilder().withProjection(keyProjection);
DataTableScan tableScan = (DataTableScan) readBuilder.newScan();
List<Split> splits =
@@ -119,7 +121,7 @@ public class IndexBootstrap implements Serializable {
options.pageSize(),
options.crossPartitionUpsertBootstrapParallelism(),
split -> {
- DataSplit dataSplit = ((DataSplit) split);
+ DataSplit dataSplit = unwrapDataSplit(split);
int bucket = dataSplit.bucket();
return partBucketConverter.toGenericRow(
new JoinedRow(dataSplit.partition(),
GenericRow.of(bucket)));
@@ -129,7 +131,7 @@ public class IndexBootstrap implements Serializable {
@VisibleForTesting
static boolean filterSplit(Split split, long indexTtl, long currentTime) {
- List<DataFileMeta> files = ((DataSplit) split).dataFiles();
+ List<DataFileMeta> files = unwrapDataSplit(split).dataFiles();
for (DataFileMeta file : files) {
long fileTime = file.creationTimeEpochMillis();
if (currentTime <= fileTime + indexTtl) {
@@ -139,6 +141,11 @@ public class IndexBootstrap implements Serializable {
return false;
}
+ @VisibleForTesting
+ static DataSplit unwrapDataSplit(Split split) {
+ return (DataSplit) split;
+ }
+
public static RowType bootstrapType(TableSchema schema) {
List<String> primaryKeys = schema.trimmedPrimaryKeys();
List<String> partitionKeys = schema.partitionKeys();
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
index bc84bca88a..0d970f2b99 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
@@ -56,6 +56,8 @@ import org.apache.paimon.utils.SegmentsCache;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import javax.annotation.Nullable;
+
import java.io.IOException;
import java.io.ObjectInputStream;
import java.io.ObjectOutputStream;
@@ -438,6 +440,13 @@ public class FallbackReadFileStoreTable extends
DelegatedFileStoreTable {
return this;
}
+ @Override
+ public FallbackReadScan withReadType(@Nullable RowType readType) {
+ mainScan.withReadType(readType);
+ fallbackScan.withReadType(readType);
+ return this;
+ }
+
@Override
public FallbackReadScan withLimit(int limit) {
mainScan.withLimit(limit);
@@ -564,8 +573,7 @@ public class FallbackReadFileStoreTable extends
DelegatedFileStoreTable {
Set<BinaryRow> completePartitions =
new
HashSet<>(newPartitionListingScan(true).listPartitions());
for (Split split : mainScan.plan().splits()) {
- DataSplit dataSplit = (DataSplit) split;
- splits.add(toFallbackSplit(dataSplit, false));
+ splits.add(toFallbackSplit(split, false));
}
List<BinaryRow> remainingPartitions =
@@ -691,9 +699,10 @@ public class FallbackReadFileStoreTable extends
DelegatedFileStoreTable {
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
if (split instanceof FallbackSplit) {
FallbackSplit fallbackSplit = (FallbackSplit) split;
+ Split wrappedSplit = fallbackSplit.wrapped();
if (fallbackSplit.isFallback()) {
try {
- return
fallbackRead.createReader(fallbackSplit.wrapped());
+ return fallbackRead.createReader(wrappedSplit);
} catch (Exception e) {
if (fallbackReadFailFast) {
if (e instanceof IOException) {
@@ -703,16 +712,15 @@ public class FallbackReadFileStoreTable extends
DelegatedFileStoreTable {
throw (RuntimeException) e;
}
throw new IOException(
- "Failed to read fallback branch split: "
- + fallbackSplit.wrapped(),
- e);
+ "Failed to read fallback branch split: " +
wrappedSplit, e);
}
LOG.error(
"Reading from supplemental branch has
problems: {}",
- fallbackSplit.wrapped(),
+ wrappedSplit,
e);
}
}
+ return mainRead.createReader(wrappedSplit);
}
return mainRead.createReader(split);
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/crosspartition/IndexBootstrapTest.java
b/paimon-core/src/test/java/org/apache/paimon/crosspartition/IndexBootstrapTest.java
index 526b6bd29f..0921ddf565 100644
---
a/paimon-core/src/test/java/org/apache/paimon/crosspartition/IndexBootstrapTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/crosspartition/IndexBootstrapTest.java
@@ -20,22 +20,37 @@ package org.apache.paimon.crosspartition;
import org.apache.paimon.CoreOptions;
import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.catalog.TableQueryAuthResult;
import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.Timestamp;
import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.manifest.FileSource;
import org.apache.paimon.options.Options;
+import org.apache.paimon.predicate.FieldRef;
+import org.apache.paimon.predicate.FieldTransform;
+import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.TableSchema;
+import org.apache.paimon.table.DelegatedFileStoreTable;
+import org.apache.paimon.table.FallbackReadFileStoreTable;
+import org.apache.paimon.table.FallbackReadFileStoreTable.FallbackSplit;
import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.table.Table;
import org.apache.paimon.table.TableTestBase;
import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.DataTableScan;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.Filter;
+import org.apache.paimon.utils.JsonSerdeUtil;
import org.apache.paimon.utils.Pair;
import org.junit.jupiter.api.Test;
+import org.mockito.AdditionalAnswers;
+import org.mockito.Mockito;
import java.time.Instant;
import java.time.ZoneId;
@@ -43,6 +58,7 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
+import java.util.Map;
import java.util.function.Consumer;
import static org.apache.paimon.crosspartition.IndexBootstrap.BUCKET_FIELD;
@@ -97,7 +113,11 @@ public class IndexBootstrapTest extends TableTestBase {
}
private Table createTable() throws Exception {
- Identifier identifier = identifier("T");
+ return createTable("T");
+ }
+
+ private Table createTable(String tableName) throws Exception {
+ Identifier identifier = identifier(tableName);
Options options = new Options();
options.set(CoreOptions.BUCKET, -1);
Schema schema =
@@ -129,6 +149,167 @@ public class IndexBootstrapTest extends TableTestBase {
assertThat(filterSplit(newSplit(newFile(100), newFile(200)), 200,
230)).isTrue();
}
+ @Test
+ public void testFallbackDataSplit() {
+ DataSplit dataSplit = newSplit(newFile(100), newFile(200));
+ FallbackSplit fallbackSplit =
FallbackReadFileStoreTable.toFallbackSplit(dataSplit, true);
+
+ assertThat(fallbackSplit).isInstanceOf(DataSplit.class);
+ assertThat(filterSplit(fallbackSplit, 50, 230)).isTrue();
+ assertThat(filterSplit(fallbackSplit, 50, 300)).isFalse();
+ assertThat(fallbackSplit.isFallback()).isTrue();
+ }
+
+ @Test
+ public void testBootstrapIgnoresQueryAuth() throws Exception {
+ FileStoreTable table = (FileStoreTable) createTable();
+ write(
+ table,
+ row(1, 1, 1, 2),
+ row(1, 2, 2, 3),
+ row(1, 3, 3, 4),
+ row(2, 4, 4, 5),
+ row(2, 5, 5, 6),
+ row(3, 6, 6, 7),
+ row(3, 7, 7, 8));
+
+ TableQueryAuthResult authResult = queryAuthResult(table);
+ IndexBootstrap indexBootstrap =
+ new IndexBootstrap(new QueryAuthFileStoreTable(table,
authResult));
+
+ List<GenericRow> result = new ArrayList<>();
+ try (RecordReader<InternalRow> reader = indexBootstrap.bootstrap(1,
0)) {
+ reader.forEachRemaining(
+ row -> result.add(GenericRow.of(row.getInt(0),
row.getInt(1), row.getInt(2))));
+ }
+
+ assertThat(result)
+ .containsExactlyInAnyOrder(
+ GenericRow.of(1, 1, 2),
+ GenericRow.of(2, 1, 3),
+ GenericRow.of(3, 1, 4),
+ GenericRow.of(4, 2, 5),
+ GenericRow.of(5, 2, 6),
+ GenericRow.of(6, 3, 7),
+ GenericRow.of(7, 3, 8));
+ }
+
+ @Test
+ public void testBootstrapIgnoresQueryAuthWithFallback() throws Exception {
+ FileStoreTable mainTable = (FileStoreTable) createTable("MAIN");
+ FileStoreTable fallbackTable = (FileStoreTable)
createTable("FALLBACK");
+ write(mainTable, row(1, 10, 10, 2));
+ write(fallbackTable, row(2, 20, 20, 3));
+
+ TableQueryAuthResult authResult = queryAuthResult(mainTable);
+ FallbackReadFileStoreTable table =
+ new FallbackReadFileStoreTable(
+ new QueryAuthFileStoreTable(mainTable, authResult),
+ new QueryAuthFileStoreTable(fallbackTable, authResult),
+ true);
+
+ List<GenericRow> result = new ArrayList<>();
+ try (RecordReader<InternalRow> reader = new
IndexBootstrap(table).bootstrap(1, 0)) {
+ reader.forEachRemaining(
+ row -> result.add(GenericRow.of(row.getInt(0),
row.getInt(1), row.getInt(2))));
+ }
+
+ assertThat(result)
+ .containsExactlyInAnyOrder(GenericRow.of(10, 1, 2),
GenericRow.of(20, 2, 3));
+ }
+
+ private TableQueryAuthResult queryAuthResult(FileStoreTable table) {
+ Predicate filter = new PredicateBuilder(table.rowType()).equal(0, 1);
+ return new TableQueryAuthResult(
+ Collections.singletonList(JsonSerdeUtil.toFlatJson(filter)),
+ Collections.singletonMap(
+ "pk",
+ JsonSerdeUtil.toFlatJson(
+ new FieldTransform(new FieldRef(0, "pt",
DataTypes.INT())))));
+ }
+
+ private static class QueryAuthFileStoreTable extends
DelegatedFileStoreTable {
+
+ private final TableQueryAuthResult authResult;
+ private final boolean queryAuthEnabled;
+
+ private QueryAuthFileStoreTable(FileStoreTable wrapped,
TableQueryAuthResult authResult) {
+ this(wrapped, authResult, true);
+ }
+
+ private QueryAuthFileStoreTable(
+ FileStoreTable wrapped, TableQueryAuthResult authResult,
boolean queryAuthEnabled) {
+ super(wrapped);
+ this.authResult = authResult;
+ this.queryAuthEnabled = queryAuthEnabled;
+ }
+
+ @Override
+ public DataTableScan newScan() {
+ DataTableScan delegate = wrapped.newScan();
+ if (!queryAuthEnabled) {
+ return delegate;
+ }
+ DataTableScan scan =
+ Mockito.mock(DataTableScan.class,
AdditionalAnswers.delegatesTo(delegate));
+ Mockito.doAnswer(
+ invocation -> {
+ Filter<Integer> filter =
invocation.getArgument(0);
+ delegate.withBucketFilter(filter);
+ return scan;
+ })
+ .when(scan)
+ .withBucketFilter(Mockito.any());
+ Mockito.doAnswer(
+ invocation -> {
+ Filter<Integer> filter =
invocation.getArgument(0);
+ delegate.withLevelFilter(filter);
+ return scan;
+ })
+ .when(scan)
+ .withLevelFilter(Mockito.any());
+ Mockito.doAnswer(ignored ->
authResult.convertPlan(delegate.plan())).when(scan).plan();
+ return scan;
+ }
+
+ @Override
+ public FileStoreTable copy(Map<String, String> dynamicOptions) {
+ return new QueryAuthFileStoreTable(
+ wrapped.copy(dynamicOptions), authResult,
queryAuthEnabled(dynamicOptions));
+ }
+
+ @Override
+ public FileStoreTable copy(TableSchema newTableSchema) {
+ return new QueryAuthFileStoreTable(
+ wrapped.copy(newTableSchema), authResult,
queryAuthEnabled);
+ }
+
+ @Override
+ public FileStoreTable copyWithoutTimeTravel(Map<String, String>
dynamicOptions) {
+ return new QueryAuthFileStoreTable(
+ wrapped.copyWithoutTimeTravel(dynamicOptions),
+ authResult,
+ queryAuthEnabled(dynamicOptions));
+ }
+
+ @Override
+ public FileStoreTable copyWithLatestSchema() {
+ return new QueryAuthFileStoreTable(
+ wrapped.copyWithLatestSchema(), authResult,
queryAuthEnabled);
+ }
+
+ @Override
+ public FileStoreTable switchToBranch(String branchName) {
+ return new QueryAuthFileStoreTable(
+ wrapped.switchToBranch(branchName), authResult,
queryAuthEnabled);
+ }
+
+ private boolean queryAuthEnabled(Map<String, String> dynamicOptions) {
+ String value =
dynamicOptions.get(CoreOptions.QUERY_AUTH_ENABLED.key());
+ return value == null ? queryAuthEnabled :
Boolean.parseBoolean(value);
+ }
+ }
+
private DataSplit newSplit(DataFileMeta... files) {
return DataSplit.builder()
.withSnapshot(1)
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/FallbackReadFileStoreTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/FallbackReadFileStoreTableTest.java
index e0841df4d0..6716dc56a7 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/FallbackReadFileStoreTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/FallbackReadFileStoreTableTest.java
@@ -20,6 +20,7 @@ package org.apache.paimon.table;
import org.apache.paimon.CoreOptions;
import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.catalog.TableQueryAuthResult;
import org.apache.paimon.data.BinaryRow;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.fs.FileIO;
@@ -29,6 +30,8 @@ import org.apache.paimon.fs.local.LocalFileIO;
import org.apache.paimon.manifest.PartitionEntry;
import org.apache.paimon.options.Options;
import org.apache.paimon.partition.PartitionPredicate;
+import org.apache.paimon.predicate.FieldRef;
+import org.apache.paimon.predicate.FieldTransform;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
import org.apache.paimon.reader.RecordReader;
@@ -40,11 +43,14 @@ import org.apache.paimon.table.sink.StreamTableCommit;
import org.apache.paimon.table.sink.StreamTableWrite;
import org.apache.paimon.table.source.DataTableScan;
import org.apache.paimon.table.source.InnerTableRead;
+import org.apache.paimon.table.source.QueryAuthSplit;
import org.apache.paimon.table.source.Split;
import org.apache.paimon.table.source.TableRead;
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.InstantiationUtil;
+import org.apache.paimon.utils.JsonSerdeUtil;
import org.apache.paimon.utils.Pair;
import org.apache.paimon.utils.TraceableFileIO;
@@ -53,6 +59,7 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.AdditionalAnswers;
import org.mockito.Mockito;
import java.io.IOException;
@@ -91,6 +98,81 @@ public class FallbackReadFileStoreTableTest {
fileIO = FileIOFinder.find(tablePath);
}
+ @Test
+ public void testScanForwardsReadType() {
+ FileStoreTable mainTable = Mockito.mock(FileStoreTable.class);
+ FileStoreTable fallbackTable = Mockito.mock(FileStoreTable.class);
+ DataTableScan mainScan = Mockito.mock(DataTableScan.class);
+ DataTableScan fallbackScan = Mockito.mock(DataTableScan.class);
+ RowType readType = ROW_TYPE.project("a");
+
+ FallbackReadFileStoreTable.FallbackReadScan scan =
+ new FallbackReadFileStoreTable.FallbackReadScan(
+ mainTable,
+ fallbackTable,
+ Mockito.mock(TableSchema.class),
+ table -> table == mainTable ? mainScan : fallbackScan);
+
+ assertThat(scan.withReadType(readType)).isSameAs(scan);
+ Mockito.verify(mainScan).withReadType(readType);
+ Mockito.verify(fallbackScan).withReadType(readType);
+ }
+
+ @Test
+ public void testPlanAndReadWithQueryAuthSplit() throws Exception {
+ FileStoreTable mainTable = createTable();
+ writeDataIntoTable(mainTable, 0, rowData(1, 10));
+
+ mainTable.createBranch("bc");
+ FileStoreTable branchTable = createTableFromBranch(mainTable, "bc");
+ writeDataIntoTable(branchTable, 0, rowData(2, 20));
+
+ FallbackReadFileStoreTable table =
+ new FallbackReadFileStoreTable(mainTable, branchTable, true);
+ TableQueryAuthResult authResult =
+ new TableQueryAuthResult(
+ null,
+ Collections.singletonMap(
+ "a",
+ JsonSerdeUtil.toFlatJson(
+ new FieldTransform(
+ new FieldRef(0, "pt",
DataTypes.INT())))));
+ DataTableScan scan =
+ table.newFallbackScan(
+ fileStoreTable ->
queryAuthScan(fileStoreTable.newScan(), authResult));
+
+ List<Split> splits = scan.plan().splits();
+ assertThat(splits).hasSize(2);
+ assertThat(splits)
+ .allSatisfy(
+ split -> {
+ assertThat(split)
+
.isInstanceOf(FallbackReadFileStoreTable.FallbackSplit.class);
+ Split wrapped =
+
((FallbackReadFileStoreTable.FallbackSplit) split).wrapped();
+
assertThat(wrapped).isInstanceOf(QueryAuthSplit.class);
+ });
+
+ List<Pair<Integer, Integer>> result = new ArrayList<>();
+ for (Split split : splits) {
+ Split deserialized =
+ InstantiationUtil.deserializeObject(
+ InstantiationUtil.serializeObject(split),
getClass().getClassLoader());
+ RecordReader<InternalRow> reader =
table.newRead().createReader(deserialized);
+ reader.forEachRemaining(r -> result.add(Pair.of(r.getInt(0),
r.getInt(1))));
+ reader.close();
+ }
+
+ assertThat(result).containsExactlyInAnyOrder(Pair.of(1, 1), Pair.of(2,
2));
+ }
+
+ private DataTableScan queryAuthScan(DataTableScan delegate,
TableQueryAuthResult authResult) {
+ DataTableScan scan =
+ Mockito.mock(DataTableScan.class,
AdditionalAnswers.delegatesTo(delegate));
+ Mockito.doAnswer(ignored ->
authResult.convertPlan(delegate.plan())).when(scan).plan();
+ return scan;
+ }
+
@ParameterizedTest
@ValueSource(booleans = {false, true})
public void testListPartitions(boolean wrappedFirst) throws Exception {