Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -165,6 +166,27 @@ public void markDonePartitions(Identifier identifier, List<Map<String, String>>
wrapped.markDonePartitions(identifier, partitions);
}

@Override
public void createPartitions(Identifier identifier, List<Map<String, String>> partitions)
throws TableNotExistException {
privilegeManager.getPrivilegeChecker().assertCanInsert(identifier);
wrapped.createPartitions(identifier, partitions);
}

@Override
public void dropPartitions(Identifier identifier, List<Map<String, String>> partitions)
throws TableNotExistException {
privilegeManager.getPrivilegeChecker().assertCanInsert(identifier);
wrapped.dropPartitions(identifier, partitions);
}

@Override
public void alterPartitions(Identifier identifier, List<PartitionStatistics> partitions)
throws TableNotExistException {
privilegeManager.getPrivilegeChecker().assertCanInsert(identifier);
wrapped.alterPartitions(identifier, partitions);
}

public void createPrivilegedUser(String user, String password) {
privilegeManager.createUser(user, password);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Snapshot> latestSnapshot() {
privilegeChecker.assertCanSelectOrInsert(identifier);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"};
Expand All @@ -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"};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading