From 318dc0f50bd3d10f464b71bc3a4a71911e20f81c Mon Sep 17 00:00:00 2001 From: sychen Date: Thu, 30 Jul 2026 16:53:05 +0800 Subject: [PATCH] feat: add privilege checks for partition and consumer ops --- .../paimon/privilege/PrivilegedCatalog.java | 22 +++++++++++++++++++ .../privilege/PrivilegedFileStoreTable.java | 9 ++++++++ .../procedure/ClearConsumersProcedure.java | 6 +---- .../procedure/ResetConsumerProcedure.java | 12 ++-------- .../flink/action/ClearConsumerAction.java | 6 +---- .../flink/action/ResetConsumerAction.java | 6 +---- .../procedure/ClearConsumersProcedure.java | 6 +---- .../procedure/ResetConsumerProcedure.java | 6 +---- .../procedure/ClearConsumersProcedure.java | 6 +---- .../procedure/ResetConsumerProcedure.java | 6 +---- 10 files changed, 40 insertions(+), 45 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedCatalog.java b/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedCatalog.java index b408055e5112..95fa93d60f8b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedCatalog.java +++ b/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedCatalog.java @@ -27,6 +27,7 @@ import org.apache.paimon.options.ConfigOption; import org.apache.paimon.options.ConfigOptions; import org.apache.paimon.options.Options; +import org.apache.paimon.partition.PartitionStatistics; import org.apache.paimon.schema.Schema; import org.apache.paimon.schema.SchemaChange; import org.apache.paimon.table.FileStoreTable; @@ -165,6 +166,27 @@ public void markDonePartitions(Identifier identifier, List> wrapped.markDonePartitions(identifier, partitions); } + @Override + public void createPartitions(Identifier identifier, List> partitions) + throws TableNotExistException { + privilegeManager.getPrivilegeChecker().assertCanInsert(identifier); + wrapped.createPartitions(identifier, partitions); + } + + @Override + public void dropPartitions(Identifier identifier, List> partitions) + throws TableNotExistException { + privilegeManager.getPrivilegeChecker().assertCanInsert(identifier); + wrapped.dropPartitions(identifier, partitions); + } + + @Override + public void alterPartitions(Identifier identifier, List partitions) + throws TableNotExistException { + privilegeManager.getPrivilegeChecker().assertCanInsert(identifier); + wrapped.alterPartitions(identifier, partitions); + } + public void createPrivilegedUser(String user, String password) { privilegeManager.createUser(user, password); } diff --git a/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedFileStoreTable.java b/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedFileStoreTable.java index e1064e70458a..d7086261b70e 100644 --- a/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedFileStoreTable.java +++ b/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedFileStoreTable.java @@ -21,6 +21,7 @@ import org.apache.paimon.FileStore; import org.apache.paimon.Snapshot; import org.apache.paimon.catalog.Identifier; +import org.apache.paimon.consumer.ConsumerManager; import org.apache.paimon.schema.TableSchema; import org.apache.paimon.stats.Statistics; import org.apache.paimon.table.DelegatedFileStoreTable; @@ -72,6 +73,14 @@ public ChangelogManager changelogManager() { return wrapped.changelogManager(); } + @Override + public ConsumerManager consumerManager() { + // Resetting/deleting a consumer's progress affects what data downstream streaming + // readers will (re)consume, so treat it as a write/insert-level operation. + privilegeChecker.assertCanInsert(identifier); + return wrapped.consumerManager(); + } + @Override public Optional latestSnapshot() { privilegeChecker.assertCanSelectOrInsert(identifier); diff --git a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java index a3f6713f4cdb..cd2090e37573 100644 --- a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java +++ b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java @@ -56,11 +56,7 @@ public String[] call( throws Catalog.TableNotExistException { FileStoreTable fileStoreTable = (FileStoreTable) catalog.getTable(Identifier.fromString(tableId)); - ConsumerManager consumerManager = - new ConsumerManager( - fileStoreTable.fileIO(), - fileStoreTable.location(), - fileStoreTable.snapshotManager().branch()); + ConsumerManager consumerManager = fileStoreTable.consumerManager(); Pattern includingPattern = StringUtils.isNullOrWhitespaceOnly(includingConsumers) diff --git a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java index 7777ccda19de..85b0f39ed5d7 100644 --- a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java +++ b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java @@ -50,11 +50,7 @@ public String[] call( FileStoreTable fileStoreTable = (FileStoreTable) catalog.getTable(Identifier.fromString(tableId)); fileStoreTable.snapshotManager().snapshot(nextSnapshotId); - ConsumerManager consumerManager = - new ConsumerManager( - fileStoreTable.fileIO(), - fileStoreTable.location(), - fileStoreTable.snapshotManager().branch()); + ConsumerManager consumerManager = fileStoreTable.consumerManager(); consumerManager.resetConsumer(consumerId, new Consumer(nextSnapshotId)); return new String[] {"Success"}; @@ -64,11 +60,7 @@ public String[] call(ProcedureContext procedureContext, String tableId, String c throws Catalog.TableNotExistException { FileStoreTable fileStoreTable = (FileStoreTable) catalog.getTable(Identifier.fromString(tableId)); - ConsumerManager consumerManager = - new ConsumerManager( - fileStoreTable.fileIO(), - fileStoreTable.location(), - fileStoreTable.snapshotManager().branch()); + ConsumerManager consumerManager = fileStoreTable.consumerManager(); consumerManager.deleteConsumer(consumerId); return new String[] {"Success"}; diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ClearConsumerAction.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ClearConsumerAction.java index 124bf28fed41..485b021a61bc 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ClearConsumerAction.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ClearConsumerAction.java @@ -51,11 +51,7 @@ public ClearConsumerAction withExcludingConsumers(@Nullable String excludingCons @Override public void executeLocally() { FileStoreTable dataTable = (FileStoreTable) table; - ConsumerManager consumerManager = - new ConsumerManager( - dataTable.fileIO(), - dataTable.location(), - dataTable.snapshotManager().branch()); + ConsumerManager consumerManager = dataTable.consumerManager(); Pattern includingPattern = StringUtils.isNullOrWhitespaceOnly(includingConsumers) diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ResetConsumerAction.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ResetConsumerAction.java index cebc3b195986..418360f10847 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ResetConsumerAction.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ResetConsumerAction.java @@ -48,11 +48,7 @@ public ResetConsumerAction withNextSnapshotIds(Long nextSnapshotId) { @Override public void executeLocally() throws Exception { FileStoreTable dataTable = (FileStoreTable) table; - ConsumerManager consumerManager = - new ConsumerManager( - dataTable.fileIO(), - dataTable.location(), - dataTable.snapshotManager().branch()); + ConsumerManager consumerManager = dataTable.consumerManager(); if (Objects.isNull(nextSnapshotId)) { consumerManager.deleteConsumer(consumerId); } else { diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java index de4c371d30a3..3a510fc3cb6e 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java @@ -71,11 +71,7 @@ public String[] call( throws Catalog.TableNotExistException { FileStoreTable fileStoreTable = (FileStoreTable) catalog.getTable(Identifier.fromString(tableId)); - ConsumerManager consumerManager = - new ConsumerManager( - fileStoreTable.fileIO(), - fileStoreTable.location(), - fileStoreTable.snapshotManager().branch()); + ConsumerManager consumerManager = fileStoreTable.consumerManager(); Pattern includingPattern = StringUtils.isNullOrWhitespaceOnly(includingConsumers) diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java index 934ce182a09c..2b35e4b1e4d8 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java @@ -61,11 +61,7 @@ public String[] call( throws Catalog.TableNotExistException { FileStoreTable fileStoreTable = (FileStoreTable) catalog.getTable(Identifier.fromString(tableId)); - ConsumerManager consumerManager = - new ConsumerManager( - fileStoreTable.fileIO(), - fileStoreTable.location(), - fileStoreTable.snapshotManager().branch()); + ConsumerManager consumerManager = fileStoreTable.consumerManager(); if (nextSnapshotId != null) { fileStoreTable.snapshotManager().snapshot(nextSnapshotId); consumerManager.resetConsumer(consumerId, new Consumer(nextSnapshotId)); diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ClearConsumersProcedure.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ClearConsumersProcedure.java index cdde1c6f8350..1847b9cf3ee6 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ClearConsumersProcedure.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ClearConsumersProcedure.java @@ -97,11 +97,7 @@ public InternalRow[] call(InternalRow args) { tableIdent, table -> { FileStoreTable fileStoreTable = (FileStoreTable) table; - ConsumerManager consumerManager = - new ConsumerManager( - fileStoreTable.fileIO(), - fileStoreTable.location(), - fileStoreTable.snapshotManager().branch()); + ConsumerManager consumerManager = fileStoreTable.consumerManager(); consumerManager.clearConsumers(includingPattern, excludingPattern); InternalRow outputRow = newInternalRow(true); diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ResetConsumerProcedure.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ResetConsumerProcedure.java index 0f7fabd05d13..856693175db2 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ResetConsumerProcedure.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ResetConsumerProcedure.java @@ -82,11 +82,7 @@ public InternalRow[] call(InternalRow args) { tableIdent, table -> { FileStoreTable fileStoreTable = (FileStoreTable) table; - ConsumerManager consumerManager = - new ConsumerManager( - fileStoreTable.fileIO(), - fileStoreTable.location(), - fileStoreTable.snapshotManager().branch()); + ConsumerManager consumerManager = fileStoreTable.consumerManager(); if (nextSnapshotId == null) { consumerManager.deleteConsumer(consumerId); } else {