diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerConfig.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerConfig.java index 4af66ebb88385..24dbaabd7307a 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerConfig.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerConfig.java @@ -56,8 +56,6 @@ public class ManagedLedgerConfig { private long retentionTimeMs = 0; private long retentionSizeInMB = 0; private boolean autoSkipNonRecoverableData; - private long offloadLedgerDeletionLagMs = TimeUnit.HOURS.toMillis(4); - private long offloadAutoTriggerSizeThresholdBytes = -1; private long metadataOperationsTimeoutSeconds = 60; private long readEntryTimeoutSeconds = 120; private long addEntryTimeoutSeconds = 120; @@ -408,48 +406,6 @@ public long getRetentionSizeInMB() { return retentionSizeInMB; } - /** - * When a ledger is offloaded from bookkeeper storage to longterm storage, the bookkeeper ledger - * is not deleted immediately. Instead we wait for a grace period before deleting from bookkeeper. - * The offloadLedgerDeleteLag sets this grace period. - * - * @param lagTime period to wait before deleting offloaded ledgers from bookkeeper - * @param unit timeunit for lagTime - */ - public ManagedLedgerConfig setOffloadLedgerDeletionLag(long lagTime, TimeUnit unit) { - this.offloadLedgerDeletionLagMs = unit.toMillis(lagTime); - return this; - } - - /** - * Number of milliseconds before an offloaded ledger will be deleted from bookkeeper. - * - * @return the offload ledger deletion lag time in milliseconds - */ - public long getOffloadLedgerDeletionLagMillis() { - return offloadLedgerDeletionLagMs; - } - - /** - * Size, in bytes, at which the managed ledger will start to automatically offload ledgers to longterm storage. - * A negative value disables autotriggering. A threshold of 0 offloads data as soon as possible. - * Offloading will not occur if no offloader has been set {@link #setLedgerOffloader(LedgerOffloader)}. - * Automatical offloading occurs when the ledger is rolled, and the ledgers up to that point exceed the threshold. - * - * @param threshold Threshold in bytes at which offload is automatically triggered - */ - public ManagedLedgerConfig setOffloadAutoTriggerSizeThresholdBytes(long threshold) { - this.offloadAutoTriggerSizeThresholdBytes = threshold; - return this; - } - - /** - * Size, in bytes, at which offloading will automatically be triggered for this managed ledger. - * @return the trigger threshold, in bytes - */ - public long getOffloadAutoTriggerSizeThresholdBytes() { - return this.offloadAutoTriggerSizeThresholdBytes; - } /** * Skip reading non-recoverable/unreadable data-ledger under managed-ledger's list. It helps when data-ledgers gets diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index 38ed3ea12995b..e539a257b091a 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java @@ -1856,8 +1856,11 @@ private void scheduleDeferredTrimming(CompletableFuture promise) { } private void maybeOffloadInBackground(CompletableFuture promise) { - if (config.getOffloadAutoTriggerSizeThresholdBytes() >= 0) { - executor.executeOrdered(name, safeRun(() -> maybeOffload(promise))); + if (config.getLedgerOffloader() != null && config.getLedgerOffloader() != NullLedgerOffloader.INSTANCE + && config.getLedgerOffloader().getOffloadPolicies() != null) { + if (config.getLedgerOffloader().getOffloadPolicies().getManagedLedgerOffloadThresholdInBytes() >= 0) { + executor.executeOrdered(name, safeRun(() -> maybeOffload(promise))); + } } } @@ -1876,39 +1879,43 @@ private void maybeOffload(CompletableFuture finalPromise) { } }); - long threshold = config.getOffloadAutoTriggerSizeThresholdBytes(); - long sizeSummed = 0; - long alreadyOffloadedSize = 0; - long toOffloadSize = 0; - - ConcurrentLinkedDeque toOffload = new ConcurrentLinkedDeque(); - - // go through ledger list from newest to oldest and build a list to offload in oldest to newest order - for (Map.Entry e : ledgers.descendingMap().entrySet()) { - long size = e.getValue().getSize(); - sizeSummed += size; - boolean alreadyOffloaded = e.getValue().hasOffloadContext() - && e.getValue().getOffloadContext().getComplete(); - if (alreadyOffloaded) { - alreadyOffloadedSize += size; - } else if (sizeSummed > threshold) { - toOffloadSize += size; - toOffload.addFirst(e.getValue()); + if (config.getLedgerOffloader() != null && config.getLedgerOffloader() != NullLedgerOffloader.INSTANCE + && config.getLedgerOffloader().getOffloadPolicies() != null) { + long threshold = config.getLedgerOffloader().getOffloadPolicies().getManagedLedgerOffloadThresholdInBytes(); + + long sizeSummed = 0; + long alreadyOffloadedSize = 0; + long toOffloadSize = 0; + + ConcurrentLinkedDeque toOffload = new ConcurrentLinkedDeque(); + + // go through ledger list from newest to oldest and build a list to offload in oldest to newest order + for (Map.Entry e : ledgers.descendingMap().entrySet()) { + long size = e.getValue().getSize(); + sizeSummed += size; + boolean alreadyOffloaded = e.getValue().hasOffloadContext() + && e.getValue().getOffloadContext().getComplete(); + if (alreadyOffloaded) { + alreadyOffloadedSize += size; + } else if (sizeSummed > threshold) { + toOffloadSize += size; + toOffload.addFirst(e.getValue()); + } } - } - if (toOffload.size() > 0) { - log.info("[{}] Going to automatically offload ledgers {}" - + ", total size = {}, already offloaded = {}, to offload = {}", - name, toOffload.stream().map(l -> l.getLedgerId()).collect(Collectors.toList()), - sizeSummed, alreadyOffloadedSize, toOffloadSize); - } else { - // offloadLoop will complete immediately with an empty list to offload - log.debug("[{}] Nothing to offload, total size = {}, already offloaded = {}, threshold = {}", - name, sizeSummed, alreadyOffloadedSize, threshold); - } + if (toOffload.size() > 0) { + log.info("[{}] Going to automatically offload ledgers {}" + + ", total size = {}, already offloaded = {}, to offload = {}", + name, toOffload.stream().map(l -> l.getLedgerId()).collect(Collectors.toList()), + sizeSummed, alreadyOffloadedSize, toOffloadSize); + } else { + // offloadLoop will complete immediately with an empty list to offload + log.debug("[{}] Nothing to offload, total size = {}, already offloaded = {}, threshold = {}", + name, sizeSummed, alreadyOffloadedSize, threshold); + } - offloadLoop(unlockingPromise, toOffload, PositionImpl.latest, Optional.empty()); + offloadLoop(unlockingPromise, toOffload, PositionImpl.latest, Optional.empty()); + } } } @@ -1930,8 +1937,15 @@ private boolean isLedgerRetentionOverSizeQuota() { private boolean isOffloadedNeedsDelete(OffloadContext offload) { long elapsedMs = clock.millis() - offload.getTimestamp(); - return offload.getComplete() && !offload.getBookkeeperDeleted() - && elapsedMs > config.getOffloadLedgerDeletionLagMillis(); + + if (config.getLedgerOffloader() != null && config.getLedgerOffloader() != NullLedgerOffloader.INSTANCE + && config.getLedgerOffloader().getOffloadPolicies() != null) { + return offload.getComplete() && !offload.getBookkeeperDeleted() + && elapsedMs > config.getLedgerOffloader() + .getOffloadPolicies().getManagedLedgerOffloadDeletionLagInMillis(); + } else { + return false; + } } /** diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java index 02fff694f185a..8fbc588b56086 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java @@ -48,7 +48,7 @@ public void testLaggedDelete() throws Exception { config.setMinimumRolloverTime(0, TimeUnit.SECONDS); config.setRetentionTime(10, TimeUnit.MINUTES); config.setRetentionSizeInMB(10); - config.setOffloadLedgerDeletionLag(5, TimeUnit.MINUTES); + offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(new Long(300000)); config.setLedgerOffloader(offloader); config.setClock(clock); @@ -109,8 +109,8 @@ public void testLaggedDeleteRetentionSetLower() throws Exception { config.setMaxEntriesPerLedger(10); config.setMinimumRolloverTime(0, TimeUnit.SECONDS); config.setRetentionTime(5, TimeUnit.MINUTES); - config.setOffloadLedgerDeletionLag(10, TimeUnit.MINUTES); config.setRetentionSizeInMB(10); + offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(new Long(600000)); config.setLedgerOffloader(offloader); config.setClock(clock); @@ -157,7 +157,7 @@ public void testLaggedDeleteSlowConsumer() throws Exception { config.setMaxEntriesPerLedger(10); config.setMinimumRolloverTime(0, TimeUnit.SECONDS); config.setRetentionTime(10, TimeUnit.MINUTES); - config.setOffloadLedgerDeletionLag(5, TimeUnit.MINUTES); + offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(new Long(300000)); config.setLedgerOffloader(offloader); config.setClock(clock); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java index 997f3f66f3310..97bab568e4cc7 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java @@ -608,9 +608,12 @@ public void testOffloadDelete() throws Exception { config.setMaxEntriesPerLedger(10); config.setMinimumRolloverTime(0, TimeUnit.SECONDS); config.setRetentionTime(0, TimeUnit.MINUTES); + offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(new Long(100)); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(100); config.setLedgerOffloader(offloader); ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("my_test_ledger", config); ManagedCursor cursor = ledger.openCursor("foobar"); + for (int i = 0; i < 15; i++) { String content = "entry-" + i; ledger.addEntry(content.getBytes()); @@ -746,9 +749,9 @@ public void testAutoTriggerOffload() throws Exception { MockLedgerOffloader offloader = new MockLedgerOffloader(); ManagedLedgerConfig config = new ManagedLedgerConfig(); config.setMaxEntriesPerLedger(10); - config.setOffloadAutoTriggerSizeThresholdBytes(100); config.setRetentionTime(10, TimeUnit.MINUTES); config.setRetentionSizeInMB(10); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(100); config.setLedgerOffloader(offloader); ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("my_test_ledger", config); @@ -782,9 +785,9 @@ public CompletableFuture offload(ReadHandle ledger, ManagedLedgerConfig config = new ManagedLedgerConfig(); config.setMaxEntriesPerLedger(10); - config.setOffloadAutoTriggerSizeThresholdBytes(100); config.setRetentionTime(10, TimeUnit.MINUTES); config.setRetentionSizeInMB(10); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(100); config.setLedgerOffloader(offloader); ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("my_test_ledger", config); @@ -843,9 +846,9 @@ public CompletableFuture offload(ReadHandle ledger, ManagedLedgerConfig config = new ManagedLedgerConfig(); config.setMaxEntriesPerLedger(10); - config.setOffloadAutoTriggerSizeThresholdBytes(100); config.setRetentionTime(10, TimeUnit.MINUTES); config.setRetentionSizeInMB(10); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(100); config.setLedgerOffloader(offloader); ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("my_test_ledger", config); @@ -894,9 +897,9 @@ public CompletableFuture offload(ReadHandle ledger, ManagedLedgerConfig config = new ManagedLedgerConfig(); config.setMaxEntriesPerLedger(10); - config.setOffloadAutoTriggerSizeThresholdBytes(100); config.setRetentionTime(10, TimeUnit.MINUTES); config.setRetentionSizeInMB(10); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(100); config.setLedgerOffloader(offloader); ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("my_test_ledger", config); @@ -926,13 +929,12 @@ public CompletableFuture offload(ReadHandle ledger, @Test public void offloadAsSoonAsClosed() throws Exception { - MockLedgerOffloader offloader = new MockLedgerOffloader(); ManagedLedgerConfig config = new ManagedLedgerConfig(); config.setMaxEntriesPerLedger(10); - config.setOffloadAutoTriggerSizeThresholdBytes(0); config.setRetentionTime(10, TimeUnit.MINUTES); config.setRetentionSizeInMB(10); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(0); config.setLedgerOffloader(offloader); ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("my_test_ledger", config); @@ -988,6 +990,12 @@ Set deletedOffloads() { return deletes.keySet(); } + OffloadPolicies offloadPolicies = OffloadPolicies.create("S3", "", "", "", + OffloadPolicies.DEFAULT_MAX_BLOCK_SIZE_IN_BYTES, + OffloadPolicies.DEFAULT_READ_BUFFER_SIZE_IN_BYTES, + OffloadPolicies.DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES, + OffloadPolicies.DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS); + @Override public String getOffloadDriverName() { return "mock"; @@ -1029,7 +1037,7 @@ public CompletableFuture deleteOffloaded(long ledgerId, UUID uuid, @Override public OffloadPolicies getOffloadPolicies() { - return null; + return offloadPolicies; } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index c0ce1e090ab40..1790e2f6a75d2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -2122,7 +2122,12 @@ protected void internalSetCompactionThreshold(long newThreshold) { protected long internalGetOffloadThreshold() { validateAdminAccessForTenant(namespaceName.getTenant()); - return getNamespacePolicies(namespaceName).offload_threshold; + Policies policies = getNamespacePolicies(namespaceName); + if (policies.offload_policies == null) { + return policies.offload_threshold; + } else { + return policies.offload_policies.getManagedLedgerOffloadThresholdInBytes(); + } } protected void internalSetOffloadThreshold(long newThreshold) { @@ -2133,8 +2138,13 @@ protected void internalSetOffloadThreshold(long newThreshold) { Stat nodeStat = new Stat(); final String path = path(POLICIES, namespaceName.toString()); byte[] content = globalZk().getData(path, null, nodeStat); + Policies policies = jsonMapper().readValue(content, Policies.class); + if (policies.offload_policies != null) { + policies.offload_policies.setManagedLedgerOffloadThresholdInBytes(newThreshold); + } policies.offload_threshold = newThreshold; + globalZk().setData(path, jsonMapper().writeValueAsBytes(policies), nodeStat.getVersion()); policiesCache().invalidate(path(POLICIES, namespaceName.toString())); log.info("[{}] Successfully updated offloadThreshold configuration: namespace={}, value={}", @@ -2159,7 +2169,12 @@ protected void internalSetOffloadThreshold(long newThreshold) { protected Long internalGetOffloadDeletionLag() { validateAdminAccessForTenant(namespaceName.getTenant()); - return getNamespacePolicies(namespaceName).offload_deletion_lag_ms; + Policies policies = getNamespacePolicies(namespaceName); + if (policies.offload_policies == null) { + return policies.offload_deletion_lag_ms; + } else { + return policies.offload_policies.getManagedLedgerOffloadDeletionLagInMillis(); + } } protected void internalSetOffloadDeletionLag(Long newDeletionLagMs) { @@ -2170,8 +2185,13 @@ protected void internalSetOffloadDeletionLag(Long newDeletionLagMs) { Stat nodeStat = new Stat(); final String path = path(POLICIES, namespaceName.toString()); byte[] content = globalZk().getData(path, null, nodeStat); + Policies policies = jsonMapper().readValue(content, Policies.class); + if (policies.offload_policies != null) { + policies.offload_policies.setManagedLedgerOffloadDeletionLagInMillis(newDeletionLagMs); + } policies.offload_deletion_lag_ms = newDeletionLagMs; + globalZk().setData(path, jsonMapper().writeValueAsBytes(policies), nodeStat.getVersion()); policiesCache().invalidate(path(POLICIES, namespaceName.toString())); log.info("[{}] Successfully updated offloadDeletionLagMs configuration: namespace={}, value={}", @@ -2310,6 +2330,20 @@ protected void internalSetOffloadPolicies(AsyncResponse asyncResponse, OffloadPo final String path = path(POLICIES, namespaceName.toString()); byte[] content = globalZk().getData(path, null, nodeStat); Policies policies = jsonMapper().readValue(content, Policies.class); + + if (offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis() + .equals(OffloadPolicies.DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS)) { + offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(policies.offload_deletion_lag_ms); + } else { + policies.offload_deletion_lag_ms = offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis(); + } + if (offloadPolicies.getManagedLedgerOffloadThresholdInBytes() == + OffloadPolicies.DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES) { + offloadPolicies.setManagedLedgerOffloadThresholdInBytes(policies.offload_threshold); + } else { + policies.offload_threshold = offloadPolicies.getManagedLedgerOffloadThresholdInBytes(); + } + policies.offload_policies = offloadPolicies; globalZk().setData(path, jsonMapper().writeValueAsBytes(policies), nodeStat.getVersion(), (rc, path1, ctx, stat) -> { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index a050b47d350b9..3fba996212945 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -1005,19 +1005,18 @@ public CompletableFuture getManagedLedgerConfig(TopicName t managedLedgerConfig.setRetentionSizeInMB(retentionPolicies.getRetentionSizeInMB()); OffloadPolicies offloadPolicies = policies.map(p -> p.offload_policies).orElse(null); + + if (offloadPolicies == null) { + offloadPolicies = new OffloadPolicies(); + offloadPolicies.setManagedLedgerOffloadDriver(pulsar.getConfiguration().getManagedLedgerOffloadDriver()); + offloadPolicies.setManagedLedgerOffloadThresholdInBytes( + pulsar.getConfiguration().getManagedLedgerOffloadAutoTriggerSizeThresholdBytes() + ); + offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis( + pulsar.getConfiguration().getManagedLedgerOffloadDeletionLagMs() + ); + } managedLedgerConfig.setLedgerOffloader(pulsar.getManagedLedgerOffloader(namespace, offloadPolicies)); - policies.ifPresent(p -> { - long lag = serviceConfig.getManagedLedgerOffloadDeletionLagMs(); - if (p.offload_deletion_lag_ms != null) { - lag = p.offload_deletion_lag_ms; - } - long bytes = serviceConfig.getManagedLedgerOffloadAutoTriggerSizeThresholdBytes(); - if (p.offload_threshold != -1L) { - bytes = p.offload_threshold; - } - managedLedgerConfig.setOffloadLedgerDeletionLag(lag, TimeUnit.MILLISECONDS); - managedLedgerConfig.setOffloadAutoTriggerSizeThresholdBytes(bytes); - }); future.complete(managedLedgerConfig); }, (exception) -> future.completeExceptionally(exception))); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiOffloadTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiOffloadTest.java index 25c45dfd29a7d..76272ef8611ce 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiOffloadTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiOffloadTest.java @@ -149,7 +149,7 @@ public void testOffloadPolicies() throws Exception { String endpoint = "test-endpoint"; OffloadPolicies offload1 = OffloadPolicies.create( - driver, region, bucket, endpoint, 100, 100); + driver, region, bucket, endpoint, 100, 100, -1, null); admin.namespaces().setOffloadPolicies(namespaceName, offload1); OffloadPolicies offload2 = admin.namespaces().getOffloadPolicies(namespaceName); Assert.assertEquals(offload1, offload2); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java index 0f7d100e4d122..e5d535cacb6b9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java @@ -43,8 +43,12 @@ import java.net.URL; import java.util.EnumSet; import java.util.List; +import java.util.Map; import java.util.Optional; +import java.util.Set; +import java.util.UUID; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; import javax.ws.rs.ClientErrorException; import javax.ws.rs.WebApplicationException; @@ -54,6 +58,8 @@ import javax.ws.rs.core.UriBuilder; import javax.ws.rs.core.UriInfo; +import org.apache.bookkeeper.client.api.ReadHandle; +import org.apache.bookkeeper.mledger.LedgerOffloader; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.util.ZkUtils; import org.apache.pulsar.broker.admin.v1.Namespaces; @@ -75,6 +81,7 @@ import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.BundlesData; import org.apache.pulsar.common.policies.data.ClusterData; +import org.apache.pulsar.common.policies.data.OffloadPolicies; import org.apache.pulsar.common.policies.data.PersistencePolicies; import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.RetentionPolicies; @@ -1075,6 +1082,74 @@ public void testSubscribeRate() throws Exception { admin.tenants().deleteTenant("my-tenants"); } + class MockLedgerOffloader implements LedgerOffloader { + ConcurrentHashMap offloads = new ConcurrentHashMap(); + ConcurrentHashMap deletes = new ConcurrentHashMap(); + + Set offloadedLedgers() { + return offloads.keySet(); + } + + Set deletedOffloads() { + return deletes.keySet(); + } + + OffloadPolicies offloadPolicies; + + public MockLedgerOffloader(OffloadPolicies offloadPolicies) { + this.offloadPolicies = offloadPolicies; + } + + @Override + public String getOffloadDriverName() { + return "mock"; + } + + @Override + public CompletableFuture offload(ReadHandle ledger, + UUID uuid, + Map extraMetadata) { + CompletableFuture promise = new CompletableFuture<>(); + if (offloads.putIfAbsent(ledger.getId(), uuid) == null) { + promise.complete(null); + } else { + promise.completeExceptionally(new Exception("Already exists exception")); + } + return promise; + } + + @Override + public CompletableFuture readOffloaded(long ledgerId, UUID uuid, + Map offloadDriverMetadata) { + CompletableFuture promise = new CompletableFuture<>(); + promise.completeExceptionally(new UnsupportedOperationException()); + return promise; + } + + @Override + public CompletableFuture deleteOffloaded(long ledgerId, UUID uuid, + Map offloadDriverMetadata) { + CompletableFuture promise = new CompletableFuture<>(); + if (offloads.remove(ledgerId, uuid)) { + deletes.put(ledgerId, uuid); + promise.complete(null); + } else { + promise.completeExceptionally(new Exception("Not found")); + } + return promise; + }; + + @Override + public OffloadPolicies getOffloadPolicies() { + return offloadPolicies; + } + + @Override + public void close() { + + } + } + @Test public void testSetOffloadThreshold() throws Exception { TopicName topicName = TopicName.get("persistent", this.testTenant, "offload", "offload-topic"); @@ -1090,25 +1165,54 @@ public void testSetOffloadThreshold() throws Exception { assertEquals(-1, admin.namespaces().getOffloadThreshold(namespace)); // the ledger config should have the expected value ManagedLedgerConfig ledgerConf = pulsar.getBrokerService().getManagedLedgerConfig(topicName).get(); - assertEquals(ledgerConf.getOffloadAutoTriggerSizeThresholdBytes(), 1); + MockLedgerOffloader offloader = new MockLedgerOffloader(OffloadPolicies.create("S3", "", "", "", + OffloadPolicies.DEFAULT_MAX_BLOCK_SIZE_IN_BYTES, + OffloadPolicies.DEFAULT_READ_BUFFER_SIZE_IN_BYTES, + admin.namespaces().getOffloadThreshold(namespace), + pulsar.getConfiguration().getManagedLedgerOffloadDeletionLagMs())); + ledgerConf.setLedgerOffloader(offloader); + assertEquals(ledgerConf.getLedgerOffloader().getOffloadPolicies().getManagedLedgerOffloadThresholdInBytes(), + -1); // set an override for the namespace admin.namespaces().setOffloadThreshold(namespace, 100); assertEquals(100, admin.namespaces().getOffloadThreshold(namespace)); ledgerConf = pulsar.getBrokerService().getManagedLedgerConfig(topicName).get(); - assertEquals(ledgerConf.getOffloadAutoTriggerSizeThresholdBytes(), 100); + admin.namespaces().getOffloadPolicies(namespace); + offloader = new MockLedgerOffloader(OffloadPolicies.create("S3", "", "", "", + OffloadPolicies.DEFAULT_MAX_BLOCK_SIZE_IN_BYTES, + OffloadPolicies.DEFAULT_READ_BUFFER_SIZE_IN_BYTES, + admin.namespaces().getOffloadThreshold(namespace), + pulsar.getConfiguration().getManagedLedgerOffloadDeletionLagMs())); + ledgerConf.setLedgerOffloader(offloader); + assertEquals(ledgerConf.getLedgerOffloader().getOffloadPolicies().getManagedLedgerOffloadThresholdInBytes(), + 100); // set another negative value to disable admin.namespaces().setOffloadThreshold(namespace, -2); assertEquals(-2, admin.namespaces().getOffloadThreshold(namespace)); ledgerConf = pulsar.getBrokerService().getManagedLedgerConfig(topicName).get(); - assertEquals(ledgerConf.getOffloadAutoTriggerSizeThresholdBytes(), -2); + offloader = new MockLedgerOffloader(OffloadPolicies.create("S3", "", "", "", + OffloadPolicies.DEFAULT_MAX_BLOCK_SIZE_IN_BYTES, + OffloadPolicies.DEFAULT_READ_BUFFER_SIZE_IN_BYTES, + admin.namespaces().getOffloadThreshold(namespace), + pulsar.getConfiguration().getManagedLedgerOffloadDeletionLagMs())); + ledgerConf.setLedgerOffloader(offloader); + assertEquals(ledgerConf.getLedgerOffloader().getOffloadPolicies().getManagedLedgerOffloadThresholdInBytes(), + -2); // set back to -1 and fall back to default admin.namespaces().setOffloadThreshold(namespace, -1); assertEquals(-1, admin.namespaces().getOffloadThreshold(namespace)); ledgerConf = pulsar.getBrokerService().getManagedLedgerConfig(topicName).get(); - assertEquals(ledgerConf.getOffloadAutoTriggerSizeThresholdBytes(), 1); + offloader = new MockLedgerOffloader(OffloadPolicies.create("S3", "", "", "", + OffloadPolicies.DEFAULT_MAX_BLOCK_SIZE_IN_BYTES, + OffloadPolicies.DEFAULT_READ_BUFFER_SIZE_IN_BYTES, + admin.namespaces().getOffloadThreshold(namespace), + pulsar.getConfiguration().getManagedLedgerOffloadDeletionLagMs())); + ledgerConf.setLedgerOffloader(offloader); + assertEquals(ledgerConf.getLedgerOffloader().getOffloadPolicies().getManagedLedgerOffloadThresholdInBytes(), + -1); // cleanup admin.topics().delete(topicName.toString(), true); diff --git a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java index 283c5fe7c12ac..d8f04388c9c63 100644 --- a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java +++ b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java @@ -501,10 +501,11 @@ void namespaces() throws Exception { namespaces.run(split("clear-offload-deletion-lag myprop/clust/ns1")); verify(mockNamespaces).clearOffloadDeleteLag("myprop/clust/ns1"); - namespaces.run(split("set-offload-policies myprop/clust/ns1 -r test-region -d aws-s3 -b test-bucket -e http://test.endpoint -mbs 32M -rbs 5M")); + namespaces.run(split("set-offload-policies myprop/clust/ns1 -r test-region -d aws-s3 -b test-bucket -e http://test.endpoint -mbs 32M -rbs 5M -oat 10M -oae 10s")); verify(mockNamespaces).setOffloadPolicies("myprop/clust/ns1", OffloadPolicies.create("aws-s3", "test-region", "test-bucket", - "http://test.endpoint", 32 * 1024 * 1024, 5 * 1024 * 1024)); + "http://test.endpoint", 32 * 1024 * 1024, 5 * 1024 * 1024, + 10 * 1024 * 1024, 10000L)); namespaces.run(split("get-offload-policies myprop/clust/ns1")); verify(mockNamespaces).getOffloadPolicies("myprop/clust/ns1"); diff --git a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaces.java b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaces.java index d77b8ec187311..e057cf0fe9a2c 100644 --- a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaces.java +++ b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaces.java @@ -1392,6 +1392,18 @@ private class SetOffloadPolicies extends CliCommand { required = false) private String readBufferSizeStr; + @Parameter( + names = {"--offloadAfterElapsed", "-oae"}, + description = "Offload after elapsed in minutes (or minutes, hours,days,weeks eg: 100m, 3h, 2d, 5w).", + required = false) + private String offloadAfterElapsedStr; + + @Parameter( + names = {"--offloadAfterThreshold", "-oat"}, + description = "Offload after threshold size (eg: 1M, 5M)", + required = false) + private String offloadAfterThresholdStr; + private final String[] DRIVER_NAMES = {"S3", "aws-s3", "google-cloud-storage"}; public boolean driverSupported(String driver) { @@ -1453,8 +1465,27 @@ && maxValueCheck("ReadBufferSize", readBufferSize, Integer.MAX_VALUE)) { } } + Long offloadAfterElapsedInMillis = OffloadPolicies.DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS; + if (StringUtils.isNotEmpty(offloadAfterElapsedStr)) { + Long offloadAfterElapsed = TimeUnit.SECONDS.toMillis(RelativeTimeUtil.parseRelativeTimeInSeconds(offloadAfterElapsedStr)); + if (positiveCheck("OffloadAfterElapsed", offloadAfterElapsed) + && maxValueCheck("OffloadAfterElapsed", offloadAfterElapsed, Long.MAX_VALUE)) { + offloadAfterElapsedInMillis = new Long(offloadAfterElapsed); + } + } + + long offloadAfterThresholdInBytes = OffloadPolicies.DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES; + if (StringUtils.isNotEmpty(offloadAfterThresholdStr)) { + long offloadAfterThreshold = validateSizeString(offloadAfterThresholdStr); + if (positiveCheck("OffloadAfterThreshold", offloadAfterThreshold) + && maxValueCheck("OffloadAfterThreshold", offloadAfterThreshold, Long.MAX_VALUE)) { + offloadAfterThresholdInBytes = new Long(offloadAfterThreshold); + } + } + OffloadPolicies offloadPolicies = OffloadPolicies.create(driver, region, bucket, endpoint, - maxBlockSizeInBytes, readBufferSizeInBytes); + maxBlockSizeInBytes, readBufferSizeInBytes, offloadAfterThresholdInBytes, + offloadAfterElapsedInMillis); admin.namespaces().setOffloadPolicies(namespace, offloadPolicies); } } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java index f46b44f031fbe..5ccb75c88957a 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java @@ -39,11 +39,15 @@ public class OffloadPolicies { public final static int DEFAULT_OFFLOAD_MAX_THREADS = 2; public final static String[] DRIVER_NAMES = {"S3", "aws-s3", "google-cloud-storage", "filesystem"}; public final static String DEFAULT_OFFLOADER_DIRECTORY = "./offloaders"; + public final static long DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES = -1; + public final static Long DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS = null; // common config private String offloadersDirectory = DEFAULT_OFFLOADER_DIRECTORY; private String managedLedgerOffloadDriver = null; private int managedLedgerOffloadMaxThreads = DEFAULT_OFFLOAD_MAX_THREADS; + private long managedLedgerOffloadThresholdInBytes = DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES; + private Long managedLedgerOffloadDeletionLagInMillis = DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS; // s3 config, set by service configuration or cli private String s3ManagedLedgerOffloadRegion = null; @@ -68,9 +72,13 @@ public class OffloadPolicies { private String fileSystemURI = null; public static OffloadPolicies create(String driver, String region, String bucket, String endpoint, - int maxBlockSizeInBytes, int readBufferSizeInBytes) { + int maxBlockSizeInBytes, int readBufferSizeInBytes, + long offloadThresholdInBytes, Long offloadDeletionLagInMillis) { OffloadPolicies offloadPolicies = new OffloadPolicies(); offloadPolicies.setManagedLedgerOffloadDriver(driver); + offloadPolicies.setManagedLedgerOffloadThresholdInBytes(offloadThresholdInBytes); + offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(offloadDeletionLagInMillis); + if (driver.equalsIgnoreCase(DRIVER_NAMES[0]) || driver.equalsIgnoreCase(DRIVER_NAMES[1])) { offloadPolicies.setS3ManagedLedgerOffloadRegion(region); offloadPolicies.setS3ManagedLedgerOffloadBucket(bucket); @@ -153,6 +161,8 @@ public int hashCode() { return Objects.hash( managedLedgerOffloadDriver, managedLedgerOffloadMaxThreads, + managedLedgerOffloadThresholdInBytes, + managedLedgerOffloadDeletionLagInMillis, s3ManagedLedgerOffloadRegion, s3ManagedLedgerOffloadBucket, s3ManagedLedgerOffloadServiceEndpoint, @@ -180,6 +190,10 @@ public boolean equals(Object obj) { OffloadPolicies other = (OffloadPolicies) obj; return Objects.equals(managedLedgerOffloadDriver, other.getManagedLedgerOffloadDriver()) && Objects.equals(managedLedgerOffloadMaxThreads, other.getManagedLedgerOffloadMaxThreads()) + && Objects.equals(managedLedgerOffloadThresholdInBytes, + other.getManagedLedgerOffloadThresholdInBytes()) + && Objects.equals(managedLedgerOffloadDeletionLagInMillis, + other.getManagedLedgerOffloadDeletionLagInMillis()) && Objects.equals(s3ManagedLedgerOffloadRegion, other.getS3ManagedLedgerOffloadRegion()) && Objects.equals(s3ManagedLedgerOffloadBucket, other.getS3ManagedLedgerOffloadBucket()) && Objects.equals(s3ManagedLedgerOffloadServiceEndpoint, @@ -208,6 +222,8 @@ public String toString() { return MoreObjects.toStringHelper(this) .add("managedLedgerOffloadDriver", managedLedgerOffloadDriver) .add("managedLedgerOffloadMaxThreads", managedLedgerOffloadMaxThreads) + .add("managedLedgerOffloadThresholdInBytes", managedLedgerOffloadThresholdInBytes) + .add("managedLedgerOffloadDeletionLagInMillis", managedLedgerOffloadDeletionLagInMillis) .add("s3ManagedLedgerOffloadRegion", s3ManagedLedgerOffloadRegion) .add("s3ManagedLedgerOffloadBucket", s3ManagedLedgerOffloadBucket) .add("s3ManagedLedgerOffloadServiceEndpoint", s3ManagedLedgerOffloadServiceEndpoint) diff --git a/site2/docs/reference-pulsar-admin.md b/site2/docs/reference-pulsar-admin.md index c52f87bad31f3..966ae388ab850 100644 --- a/site2/docs/reference-pulsar-admin.md +++ b/site2/docs/reference-pulsar-admin.md @@ -2319,3 +2319,5 @@ Options |`-e`, `--endpoint`|Alternative endpoint to connect to|| |`-mbs`, `--maxBlockSize`|Max block size|64MB| |`-rbs`, `--readBufferSize`|Read buffer size|1MB| +|`-oat`, `--offloadAfterThreshold`|Offload after threshold size (eg: 1M, 5M)|| +|`-oae`, `--offloadAfterElapsed`|Offload after elapsed in millis (or minutes, hours,days,weeks eg: 100m, 3h, 2d, 5w).||