From 8164838843b24ac19b3cfb2196e416d3d1955a59 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 17 Apr 2019 10:20:49 -0700 Subject: [PATCH 01/15] Allow to configure the managed ledger cache eviction frequency --- conf/broker.conf | 3 ++ .../mledger/ManagedLedgerConfig.java | 34 ++++++++++++---- .../mledger/impl/ManagedLedgerImpl.java | 40 +++++++++---------- .../pulsar/broker/ServiceConfiguration.java | 3 ++ .../pulsar/broker/service/BrokerService.java | 1 + 5 files changed, 53 insertions(+), 28 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 1bb171c7e7490..03b77d37190cf 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -442,6 +442,9 @@ managedLedgerCacheSizeMB= # Threshold to which bring down the cache level when eviction is triggered managedLedgerCacheEvictionWatermark=0.9 +# Configure the cache eviction frequency for the managed ledger cache +managedLedgerCacheEvictionFrequency=100.0 + # Rate limit the amount of writes per second generated by consumer acking the messages managedLedgerDefaultMarkDeleteRateLimit=1.0 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 8f890505167fe..27a1ec34b1209 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 @@ -63,6 +63,7 @@ public class ManagedLedgerConfig { private byte[] password = "".getBytes(Charsets.UTF_8); private LedgerOffloader ledgerOffloader = NullLedgerOffloader.INSTANCE; private Clock clock = Clock.systemUTC(); + private double cacheEvictionFrequency = 100; public boolean isCreateIfMissing() { return createIfMissing; @@ -515,9 +516,9 @@ public ManagedLedgerConfig setClock(Clock clock) { } /** - * + * * Ledger-Op (Create/Delete) timeout - * + * * @return */ public long getMetadataOperationsTimeoutSeconds() { @@ -526,17 +527,17 @@ public long getMetadataOperationsTimeoutSeconds() { /** * Ledger-Op (Create/Delete) timeout after which callback will be completed with failure - * + * * @param metadataOperationsTimeoutSeconds */ public ManagedLedgerConfig setMetadataOperationsTimeoutSeconds(long metadataOperationsTimeoutSeconds) { this.metadataOperationsTimeoutSeconds = metadataOperationsTimeoutSeconds; return this; } - + /** * Ledger read-entry timeout - * + * * @return */ public long getReadEntryTimeoutSeconds() { @@ -546,7 +547,7 @@ public long getReadEntryTimeoutSeconds() { /** * Ledger read entry timeout after which callback will be completed with failure. (disable timeout by setting * readTimeoutSeconds <= 0) - * + * * @param readTimeoutSeconds * @return */ @@ -554,18 +555,35 @@ public ManagedLedgerConfig setReadEntryTimeoutSeconds(long readEntryTimeoutSecon this.readEntryTimeoutSeconds = readEntryTimeoutSeconds; return this; } - + public long getAddEntryTimeoutSeconds() { return addEntryTimeoutSeconds; } /** * Add-entry timeout after which add-entry callback will be failed if add-entry is not succeeded. - * + * * @param addEntryTimeoutSeconds */ public ManagedLedgerConfig setAddEntryTimeoutSeconds(long addEntryTimeoutSeconds) { this.addEntryTimeoutSeconds = addEntryTimeoutSeconds; return this; } + + /** + * @return the configured cache eviction frequency for the managed ledger + */ + public double getCacheEvictionFrequency() { + return cacheEvictionFrequency; + } + + /** + * Configure the cache eviction frequency for the managed ledger + * + * @param cacheEvictionFrequency + * a frequence, in updates/s. Default is 100 + */ + public void setCacheEvictionFrequency(double cacheEvictionFrequency) { + this.cacheEvictionFrequency = cacheEvictionFrequency; + } } 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 eea6c1fb41b6a..9c2ffa988029b 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 @@ -20,10 +20,26 @@ import static com.google.common.base.Preconditions.checkArgument; import static java.lang.Math.min; +import static org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.FALSE; +import static org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.TRUE; import static org.apache.bookkeeper.mledger.util.SafeRun.safeRun; +import com.google.common.collect.BoundType; +import com.google.common.collect.ImmutableMap; +import com.google.common.collect.Lists; +import com.google.common.collect.Maps; +import com.google.common.collect.Queues; +import com.google.common.collect.Range; +import com.google.common.util.concurrent.RateLimiter; + +import io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; +import io.netty.util.Recycler; +import io.netty.util.Recycler.Handle; + import java.time.Clock; import java.util.Collections; +import java.util.HashMap; import java.util.Iterator; import java.util.List; import java.util.Map; @@ -61,13 +77,13 @@ import org.apache.bookkeeper.common.util.OrderedExecutor; import org.apache.bookkeeper.common.util.OrderedScheduler; import org.apache.bookkeeper.common.util.Retries; -import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.AddEntryCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.CloseCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteCursorCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteLedgerCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.OffloadCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenCursorCallback; +import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntryCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.TerminateCallback; import org.apache.bookkeeper.mledger.Entry; @@ -105,22 +121,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.google.common.collect.BoundType; -import com.google.common.collect.ImmutableMap; -import com.google.common.collect.Lists; -import com.google.common.collect.Maps; -import com.google.common.collect.Queues; -import com.google.common.collect.Range; -import com.google.common.util.concurrent.RateLimiter; - -import io.netty.buffer.ByteBuf; -import io.netty.buffer.Unpooled; -import io.netty.util.Recycler; -import io.netty.util.Recycler.Handle; -import java.util.HashMap; -import static org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.TRUE; -import static org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.FALSE; - public class ManagedLedgerImpl implements ManagedLedger, CreateCallback { private final static long MegaByte = 1024 * 1024; @@ -159,7 +159,7 @@ public class ManagedLedgerImpl implements ManagedLedger, CreateCallback { @SuppressWarnings("unused") private volatile long totalSize = 0; - private RateLimiter updateCursorRateLimit; + private RateLimiter cacheEvictionRateLimit; // Cursors that are waiting to be notified when new entries are persisted final ConcurrentLinkedQueue waitingCursors; @@ -263,7 +263,7 @@ public ManagedLedgerImpl(ManagedLedgerFactoryImpl factory, BookKeeper bookKeeper this.entryCache = factory.getEntryCacheManager().getEntryCache(this); this.waitingCursors = Queues.newConcurrentLinkedQueue(); this.uninitializedCursors = Maps.newHashMap(); - this.updateCursorRateLimit = RateLimiter.create(1); + this.cacheEvictionRateLimit = RateLimiter.create(config.getCacheEvictionFrequency()); this.clock = config.getClock(); // Get the next rollover time. Add a random value upto 5% to avoid rollover multiple ledgers at the same time @@ -1588,7 +1588,7 @@ private void internalReadFromLedger(ReadHandle ledger, OpReadEntry opReadEntry) } asyncReadEntry(ledger, firstEntry, lastEntry, false, opReadEntry, opReadEntry.ctx); - if (updateCursorRateLimit.tryAcquire()) { + if (cacheEvictionRateLimit.tryAcquire()) { if (isCursorActive(cursor)) { final PositionImpl lastReadPosition = PositionImpl.get(ledger.getId(), lastEntry); discardEntriesFromCache(cursor, lastReadPosition); diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index b3262c1e59659..5260f3b5cb335 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -771,6 +771,9 @@ public class ServiceConfiguration implements PulsarConfiguration { doc = "Threshold to which bring down the cache level when eviction is triggered" ) private double managedLedgerCacheEvictionWatermark = 0.9f; + @FieldContext(category = CATEGORY_STORAGE_ML, + doc = "Configure the cache eviction frequency for the managed ledger cache. Default is 100/s") + private double managedLedgerCacheEvictionFrequency = 100.0; @FieldContext( category = CATEGORY_STORAGE_ML, doc = "Rate limit the amount of writes per second generated by consumer acking the messages" 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 fddda96351f72..34c06cf3fe5b0 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 @@ -736,6 +736,7 @@ public CompletableFuture getManagedLedgerConfig(TopicName t managedLedgerConfig.setThrottleMarkDelete(persistencePolicies.getManagedLedgerMaxMarkDeleteRate()); managedLedgerConfig.setDigestType(serviceConfig.getManagedLedgerDigestType()); + managedLedgerConfig.setCacheEvictionFrequency(serviceConfig.getManagedLedgerCacheEvictionFrequency()); managedLedgerConfig.setMaxUnackedRangesToPersist(serviceConfig.getManagedLedgerMaxUnackedRangesToPersist()); managedLedgerConfig.setMaxUnackedRangesToPersistInZk(serviceConfig.getManagedLedgerMaxUnackedRangesToPersistInZooKeeper()); managedLedgerConfig.setMaxEntriesPerLedger(serviceConfig.getManagedLedgerMaxEntriesPerLedger()); From 1c32c929d290fb9054284655d295b85ff566a97e Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 17 Apr 2019 12:34:07 -0700 Subject: [PATCH 02/15] Fixed test --- .../bookkeeper/mledger/impl/EntryCacheManagerTest.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java index 47c515d13757a..c6fdf4f54702e 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java @@ -28,6 +28,7 @@ import org.apache.bookkeeper.common.util.OrderedScheduler; import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; import org.testng.annotations.BeforeClass; @@ -211,7 +212,9 @@ void verifyHitsMisses() throws Exception { factory = new ManagedLedgerFactoryImpl(bkc, bkc.getZkHandle(), config); EntryCacheManager cacheManager = factory.getEntryCacheManager(); - ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("ledger"); + ManagedLedgerConfig mlConf = new ManagedLedgerConfig(); + mlConf.setCacheEvictionFrequency(1); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("ledger", mlConf); ManagedCursorImpl c1 = (ManagedCursorImpl) ledger.openCursor("c1"); ManagedCursorImpl c2 = (ManagedCursorImpl) ledger.openCursor("c2"); From ca6e0bfa5c9940fa3cfb0caa57346123f8f0588f Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Tue, 23 Apr 2019 15:35:37 -0700 Subject: [PATCH 03/15] Simplified the cache eviction to make it predictable at the configured frequency --- conf/broker.conf | 4 +- .../bookkeeper/mledger/ManagedLedger.java | 6 - .../mledger/ManagedLedgerConfig.java | 18 --- .../mledger/ManagedLedgerFactoryConfig.java | 5 + .../bookkeeper/mledger/impl/EntryCache.java | 2 + .../mledger/impl/EntryCacheImpl.java | 19 ++- .../mledger/impl/EntryCacheManager.java | 4 + .../bookkeeper/mledger/impl/EntryImpl.java | 10 ++ .../impl/ManagedLedgerFactoryImpl.java | 30 ++++- .../mledger/impl/ManagedLedgerImpl.java | 76 ----------- .../bookkeeper/mledger/util/RangeCache.java | 45 ++++++- .../mledger/impl/EntryCacheManagerTest.java | 5 +- .../mledger/impl/ManagedLedgerTest.java | 123 ++++-------------- .../mledger/util/RangeCacheTest.java | 24 +++- .../broker/ManagedLedgerClientFactory.java | 1 + .../pulsar/broker/service/BrokerService.java | 1 - .../pulsar/broker/service/PulsarStats.java | 3 - .../api/SimpleProducerConsumerTest.java | 88 +------------ .../api/v1/V1_ProducerConsumerTest.java | 87 ------------- 19 files changed, 157 insertions(+), 394 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 03b77d37190cf..67c8cbcc1ce79 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -442,8 +442,8 @@ managedLedgerCacheSizeMB= # Threshold to which bring down the cache level when eviction is triggered managedLedgerCacheEvictionWatermark=0.9 -# Configure the cache eviction frequency for the managed ledger cache -managedLedgerCacheEvictionFrequency=100.0 +# Configure the cache eviction frequency for the managed ledger cache (evictions/sec) +managedLedgerCacheEvictionFrequency=10.0 # Rate limit the amount of writes per second generated by consumer acking the messages managedLedgerDefaultMarkDeleteRateLimit=1.0 diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java index d51b0d803a7d0..36e542903df25 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java @@ -322,12 +322,6 @@ public interface ManagedLedger { */ long getEstimatedBacklogSize(); - /** - * Activate cursors those caught up backlog-threshold entries and deactivate slow cursors which are creating - * backlog. - */ - void checkBackloggedCursors(); - void asyncTerminate(TerminateCallback callback, Object ctx); /** 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 27a1ec34b1209..4a57055393c67 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 @@ -63,7 +63,6 @@ public class ManagedLedgerConfig { private byte[] password = "".getBytes(Charsets.UTF_8); private LedgerOffloader ledgerOffloader = NullLedgerOffloader.INSTANCE; private Clock clock = Clock.systemUTC(); - private double cacheEvictionFrequency = 100; public boolean isCreateIfMissing() { return createIfMissing; @@ -569,21 +568,4 @@ public ManagedLedgerConfig setAddEntryTimeoutSeconds(long addEntryTimeoutSeconds this.addEntryTimeoutSeconds = addEntryTimeoutSeconds; return this; } - - /** - * @return the configured cache eviction frequency for the managed ledger - */ - public double getCacheEvictionFrequency() { - return cacheEvictionFrequency; - } - - /** - * Configure the cache eviction frequency for the managed ledger - * - * @param cacheEvictionFrequency - * a frequence, in updates/s. Default is 100 - */ - public void setCacheEvictionFrequency(double cacheEvictionFrequency) { - this.cacheEvictionFrequency = cacheEvictionFrequency; - } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java index 81673f4351ba9..d7894792e7de0 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java @@ -37,6 +37,11 @@ public class ManagedLedgerFactoryConfig { private int numManagedLedgerWorkerThreads = Runtime.getRuntime().availableProcessors(); private int numManagedLedgerSchedulerThreads = Runtime.getRuntime().availableProcessors(); + /** + * Frequency of cache eviction triggering. Default is 10 times per second. + */ + private double cacheEvictionFrequency = 10; + public long getMaxCacheSize() { return maxCacheSize; } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCache.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCache.java index 0c99650cecf18..32ad3a0392e09 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCache.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCache.java @@ -54,6 +54,8 @@ public interface EntryCache extends Comparable { */ void invalidateEntries(PositionImpl lastPosition); + void invalidateEntriesBeforeTimestamp(long timestamp); + /** * Remove from the cache all the entries belonging to a specific ledger. * diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java index a435c1442b90b..66c91c6d052b9 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java @@ -21,15 +21,17 @@ import static com.google.common.base.Preconditions.checkArgument; import static com.google.common.base.Preconditions.checkNotNull; import static org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.createManagedLedgerException; -import static org.apache.bookkeeper.mledger.util.SafeRun.safeRun; import com.google.common.collect.Lists; import com.google.common.primitives.Longs; + import io.netty.buffer.ByteBuf; import io.netty.buffer.PooledByteBufAllocator; + import java.util.Collection; import java.util.Iterator; import java.util.List; + import org.apache.bookkeeper.client.api.BKException; import org.apache.bookkeeper.client.api.LedgerEntry; import org.apache.bookkeeper.client.api.ReadHandle; @@ -37,7 +39,6 @@ import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntryCallback; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.util.RangeCache; -import org.apache.bookkeeper.mledger.util.RangeCache.Weighter; import org.apache.commons.lang3.tuple.Pair; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -53,12 +54,10 @@ public class EntryCacheImpl implements EntryCache { private static final double MB = 1024 * 1024; - private static final Weighter entryWeighter = EntryImpl::getLength; - public EntryCacheImpl(EntryCacheManager manager, ManagedLedgerImpl ml) { this.manager = manager; this.ml = ml; - this.entries = new RangeCache<>(entryWeighter); + this.entries = new RangeCache<>(EntryImpl::getLength, EntryImpl::getTimestamp); if (log.isDebugEnabled()) { log.debug("[{}] Initialized managed-ledger entry cache", ml.getName()); @@ -173,7 +172,7 @@ public void asyncReadEntry(ReadHandle lh, PositionImpl position, final ReadEntry callback.readEntryFailed(createManagedLedgerException(t), ctx); } } - + private void asyncReadEntry0(ReadHandle lh, PositionImpl position, final ReadEntryCallback callback, final Object ctx) { if (log.isDebugEnabled()) { @@ -229,7 +228,7 @@ public void asyncReadEntry(ReadHandle lh, long firstEntry, long lastEntry, boole callback.readEntriesFailed(createManagedLedgerException(t), ctx); } } - + @SuppressWarnings({ "unchecked", "rawtypes" }) private void asyncReadEntry0(ReadHandle lh, long firstEntry, long lastEntry, boolean isSlowestReader, final ReadEntriesCallback callback, Object ctx) { @@ -341,5 +340,11 @@ public Pair evictEntries(long sizeToFree) { return evicted; } + @Override + public void invalidateEntriesBeforeTimestamp(long timestamp) { + Pair evicted = entries.evictLEntriesBeforeTimestamp(timestamp); + manager.entriesRemoved(evicted.getRight()); + } + private static final Logger log = LoggerFactory.getLogger(EntryCacheImpl.class); } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheManager.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheManager.java index c551002b474f9..526796e3c26bd 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheManager.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheManager.java @@ -189,6 +189,10 @@ public Pair evictEntries(long sizeToFree) { return Pair.of(0, (long) 0); } + @Override + public void invalidateEntriesBeforeTimestamp(long timestamp) { + } + @Override public void asyncReadEntry(ReadHandle lh, long firstEntry, long lastEntry, boolean isSlowestReader, final ReadEntriesCallback callback, Object ctx) { diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java index eeccfe7857d2c..2a04f0e538569 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java @@ -40,12 +40,14 @@ protected EntryImpl newObject(Handle handle) { }; private final Handle recyclerHandle; + private long timestamp; private long ledgerId; private long entryId; ByteBuf data; public static EntryImpl create(LedgerEntry ledgerEntry) { EntryImpl entry = RECYCLER.get(); + entry.timestamp = System.nanoTime(); entry.ledgerId = ledgerEntry.getLedgerId(); entry.entryId = ledgerEntry.getEntryId(); entry.data = ledgerEntry.getEntryBuffer(); @@ -57,6 +59,7 @@ public static EntryImpl create(LedgerEntry ledgerEntry) { // Used just for tests public static EntryImpl create(long ledgerId, long entryId, byte[] data) { EntryImpl entry = RECYCLER.get(); + entry.timestamp = System.nanoTime(); entry.ledgerId = ledgerId; entry.entryId = entryId; entry.data = Unpooled.wrappedBuffer(data); @@ -66,6 +69,7 @@ public static EntryImpl create(long ledgerId, long entryId, byte[] data) { public static EntryImpl create(long ledgerId, long entryId, ByteBuf data) { EntryImpl entry = RECYCLER.get(); + entry.timestamp = System.nanoTime(); entry.ledgerId = ledgerId; entry.entryId = entryId; entry.data = data; @@ -76,6 +80,7 @@ public static EntryImpl create(long ledgerId, long entryId, ByteBuf data) { public static EntryImpl create(PositionImpl position, ByteBuf data) { EntryImpl entry = RECYCLER.get(); + entry.timestamp = System.nanoTime(); entry.ledgerId = position.getLedgerId(); entry.entryId = position.getEntryId(); entry.data = data; @@ -86,6 +91,7 @@ public static EntryImpl create(PositionImpl position, ByteBuf data) { public static EntryImpl create(EntryImpl other) { EntryImpl entry = RECYCLER.get(); + entry.timestamp = System.nanoTime(); entry.ledgerId = other.ledgerId; entry.entryId = other.entryId; entry.data = other.data.retainedDuplicate(); @@ -97,6 +103,10 @@ private EntryImpl(Recycler.Handle recyclerHandle) { this.recyclerHandle = recyclerHandle; } + public long getTimestamp() { + return timestamp; + } + @Override public ByteBuf getDataBuffer() { return data; diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index b9f63a850b9e9..73a3be7d12a31 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -88,7 +88,9 @@ public class ManagedLedgerFactoryImpl implements ManagedLedgerFactory { private final EntryCacheManager entryCacheManager; private long lastStatTimestamp = System.nanoTime(); + private long lastCacheEvictionTimestamp = System.nanoTime(); private final ScheduledFuture statsTask; + private final ScheduledFuture cacheEvictionTask; private static final int StatsPeriodSeconds = 60; public ManagedLedgerFactoryImpl(ClientConfiguration bkClientConfiguration) throws Exception { @@ -139,6 +141,8 @@ private ManagedLedgerFactoryImpl(BookKeeper bookKeeper, boolean isBookkeeperMana this.mbean = new ManagedLedgerFactoryMBeanImpl(this); this.entryCacheManager = new EntryCacheManager(this); this.statsTask = scheduledExecutor.scheduleAtFixedRate(() -> refreshStats(), 0, StatsPeriodSeconds, TimeUnit.SECONDS); + this.cacheEvictionTask = scheduledExecutor.scheduleAtFixedRate(() -> cacheEviction(), 0, + (long) (1000 / config.getCacheEvictionFrequency()), TimeUnit.MILLISECONDS); } private synchronized void refreshStats() { @@ -147,15 +151,34 @@ private synchronized void refreshStats() { mbean.refreshStats(period, TimeUnit.NANOSECONDS); ledgers.values().forEach(mlfuture -> { - ManagedLedgerImpl ml = mlfuture.getNow(null); - if (ml != null) { - ml.mbean.refreshStats(period, TimeUnit.NANOSECONDS); + if (mlfuture.isDone() && !mlfuture.isCompletedExceptionally()) { + ManagedLedgerImpl ml = mlfuture.getNow(null); + if (ml != null) { + ml.mbean.refreshStats(period, TimeUnit.NANOSECONDS); + } } }); lastStatTimestamp = now; } + private synchronized void cacheEviction() { + long now = System.nanoTime(); + long period = now - lastCacheEvictionTimestamp; + long maxTimestamp = now - period; + + ledgers.values().forEach(mlfuture -> { + if (mlfuture.isDone() && !mlfuture.isCompletedExceptionally()) { + ManagedLedgerImpl ml = mlfuture.getNow(null); + if (ml != null) { + ml.entryCache.invalidateEntriesBeforeTimestamp(maxTimestamp); + } + } + }); + + lastCacheEvictionTimestamp = now; + } + /** * Helper for getting stats. * @@ -323,6 +346,7 @@ void close(ManagedLedger ledger) { @Override public void shutdown() throws InterruptedException, ManagedLedgerException { statsTask.cancel(true); + cacheEvictionTask.cancel(true); int numLedgers = ledgers.size(); final CountDownLatch latch = new CountDownLatch(numLedgers); 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 9c2ffa988029b..d883c740000d9 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 @@ -30,7 +30,6 @@ import com.google.common.collect.Maps; import com.google.common.collect.Queues; import com.google.common.collect.Range; -import com.google.common.util.concurrent.RateLimiter; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; @@ -114,9 +113,7 @@ import org.apache.bookkeeper.mledger.util.CallbackMutex; import org.apache.bookkeeper.mledger.util.Futures; import org.apache.commons.lang3.tuple.Pair; -import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.InitialPosition; -import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata; import org.apache.pulsar.common.util.collections.ConcurrentLongHashMap; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -125,8 +122,6 @@ public class ManagedLedgerImpl implements ManagedLedger, CreateCallback { private final static long MegaByte = 1024 * 1024; protected final static int AsyncOperationTimeoutSeconds = 30; - private final static long maxActiveCursorBacklogEntries = 100; - private static long maxMessageCacheRetentionTimeMillis = 10 * 1000; protected final BookKeeper bookKeeper; protected final String name; @@ -159,8 +154,6 @@ public class ManagedLedgerImpl implements ManagedLedger, CreateCallback { @SuppressWarnings("unused") private volatile long totalSize = 0; - private RateLimiter cacheEvictionRateLimit; - // Cursors that are waiting to be notified when new entries are persisted final ConcurrentLinkedQueue waitingCursors; @@ -263,7 +256,6 @@ public ManagedLedgerImpl(ManagedLedgerFactoryImpl factory, BookKeeper bookKeeper this.entryCache = factory.getEntryCacheManager().getEntryCache(this); this.waitingCursors = Queues.newConcurrentLinkedQueue(); this.uninitializedCursors = Maps.newHashMap(); - this.cacheEvictionRateLimit = RateLimiter.create(config.getCacheEvictionFrequency()); this.clock = config.getClock(); // Get the next rollover time. Add a random value upto 5% to avoid rollover multiple ledgers at the same time @@ -912,67 +904,6 @@ public long getTotalSize() { return TOTAL_SIZE_UPDATER.get(this); } - @Override - public void checkBackloggedCursors() { - - // activate caught up cursors - cursors.forEach(cursor -> { - if (cursor.getNumberOfEntries() < maxActiveCursorBacklogEntries) { - cursor.setActive(); - } - }); - - // deactivate backlog cursors - Iterator cursors = activeCursors.iterator(); - while (cursors.hasNext()) { - ManagedCursor cursor = cursors.next(); - long backlogEntries = cursor.getNumberOfEntries(); - if (backlogEntries > maxActiveCursorBacklogEntries) { - PositionImpl readPosition = (PositionImpl) cursor.getReadPosition(); - readPosition = isValidPosition(readPosition) ? readPosition : getNextValidPosition(readPosition); - if (readPosition == null) { - if (log.isDebugEnabled()) { - log.debug("[{}] Couldn't find valid read position [{}] {}", name, cursor.getName(), - cursor.getReadPosition()); - } - continue; - } - try { - asyncReadEntry(readPosition, new ReadEntryCallback() { - - @Override - public void readEntryFailed(ManagedLedgerException e, Object ctx) { - log.warn("[{}] Failed while reading entries on [{}] {}", name, cursor.getName(), - e.getMessage(), e); - - } - - @Override - public void readEntryComplete(Entry entry, Object ctx) { - MessageMetadata msgMetadata = null; - try { - msgMetadata = Commands.parseMessageMetadata(entry.getDataBuffer()); - long msgTimeSincePublish = (clock.millis() - msgMetadata.getPublishTime()); - if (msgTimeSincePublish > maxMessageCacheRetentionTimeMillis) { - cursor.setInactive(); - } - } finally { - if (msgMetadata != null) { - msgMetadata.recycle(); - } - entry.release(); - } - - } - }, null); - } catch (Exception e) { - log.warn("[{}] Failed while reading entries from cache on [{}] {}", name, cursor.getName(), - e.getMessage(), e); - } - } - } - } - @Override public long getEstimatedBacklogSize() { @@ -1587,13 +1518,6 @@ private void internalReadFromLedger(ReadHandle ledger, OpReadEntry opReadEntry) lastEntry); } asyncReadEntry(ledger, firstEntry, lastEntry, false, opReadEntry, opReadEntry.ctx); - - if (cacheEvictionRateLimit.tryAcquire()) { - if (isCursorActive(cursor)) { - final PositionImpl lastReadPosition = PositionImpl.get(ledger.getId(), lastEntry); - discardEntriesFromCache(cursor, lastReadPosition); - } - } } protected void asyncReadEntry(ReadHandle ledger, PositionImpl position, ReadEntryCallback callback, Object ctx) { diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java index b9b3aceb0df9b..dcf26b50048ed 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java @@ -43,12 +43,13 @@ public class RangeCache, Value extends ReferenceCoun private final ConcurrentNavigableMap entries; private AtomicLong size; // Total size of values stored in cache private final Weighter weighter; // Weighter object used to extract the size from values + private final TimestampExtractor timestampExtractor; // Extract the timestamp associated with a value /** * Construct a new RangeLruCache with default Weighter. */ public RangeCache() { - this(new DefaultWeighter()); + this(new DefaultWeighter(), (x) -> System.nanoTime()); } /** @@ -57,10 +58,11 @@ public RangeCache() { * @param weighter * a custom weighter to compute the size of each stored value */ - public RangeCache(Weighter weighter) { + public RangeCache(Weighter weighter, TimestampExtractor timestampExtractor) { this.size = new AtomicLong(0); this.entries = new ConcurrentSkipListMap<>(); this.weighter = weighter; + this.timestampExtractor = timestampExtractor; } /** @@ -175,6 +177,36 @@ public Pair evictLeastAccessedEntries(long minSize) { return Pair.of(removedEntries, removedSize); } + /** + * + * @param minSize + * @return a pair containing the number of entries evicted and their total size + */ + public Pair evictLEntriesBeforeTimestamp(long maxTimestamp) { + long removedSize = 0; + int removedEntries = 0; + + while (true) { + Map.Entry entry = entries.firstEntry(); + if (entry == null || timestampExtractor.getTimestamp(entry.getValue()) > maxTimestamp) { + break; + } + + entry = entries.pollFirstEntry(); + if (entry == null) { + break; + } + + Value value = entry.getValue(); + ++removedEntries; + removedSize += weighter.getSize(value); + value.release(); + } + + size.addAndGet(-removedSize); + return Pair.of(removedEntries, removedSize); + } + /** * Just for testing. Getting the number of entries is very expensive on the conncurrent map */ @@ -217,6 +249,15 @@ public interface Weighter { long getSize(ValueT value); } + /** + * Interface of a object that is able to the extract the "timestamp" of the cached values. + * + * @param + */ + public interface TimestampExtractor { + long getTimestamp(ValueT value); + } + /** * Default cache weighter, every value is assumed the same cost. * diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java index c6fdf4f54702e..d3f5733fff74c 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java @@ -208,13 +208,12 @@ void verifyHitsMisses() throws Exception { ManagedLedgerFactoryConfig config = new ManagedLedgerFactoryConfig(); config.setMaxCacheSize(7 * 10); config.setCacheEvictionWatermark(0.8); + config.setCacheEvictionFrequency(1); factory = new ManagedLedgerFactoryImpl(bkc, bkc.getZkHandle(), config); EntryCacheManager cacheManager = factory.getEntryCacheManager(); - ManagedLedgerConfig mlConf = new ManagedLedgerConfig(); - mlConf.setCacheEvictionFrequency(1); - ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("ledger", mlConf); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("ledger"); ManagedCursorImpl c1 = (ManagedCursorImpl) ledger.openCursor("c1"); ManagedCursorImpl c2 = (ManagedCursorImpl) ledger.openCursor("c2"); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index 497ae9e72224e..eefb0afb6639b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -22,7 +22,6 @@ import static org.mockito.Matchers.anyInt; import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.doReturn; -import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; @@ -31,14 +30,19 @@ import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; +import com.google.common.base.Charsets; +import com.google.common.collect.Sets; + +import io.netty.buffer.ByteBuf; +import io.netty.buffer.PooledByteBufAllocator; +import io.netty.buffer.Unpooled; + import java.io.IOException; import java.lang.reflect.Field; -import java.lang.reflect.Modifier; import java.nio.charset.Charset; import java.security.GeneralSecurityException; import java.util.ArrayList; import java.util.Collections; -import java.util.EnumSet; import java.util.Iterator; import java.util.List; import java.util.Set; @@ -60,13 +64,10 @@ import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.client.BookKeeper.DigestType; import org.apache.bookkeeper.client.LedgerHandle; -import org.apache.bookkeeper.client.api.LedgerEntries; -import org.apache.bookkeeper.client.api.ReadHandle; import org.apache.bookkeeper.client.PulsarMockBookKeeper; import org.apache.bookkeeper.client.PulsarMockLedgerHandle; import org.apache.bookkeeper.client.api.LedgerEntries; import org.apache.bookkeeper.client.api.ReadHandle; -import org.apache.bookkeeper.client.api.WriteFlag; import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.mledger.AsyncCallbacks.AddEntryCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.CloseCallback; @@ -86,6 +87,7 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerNotFoundException; import org.apache.bookkeeper.mledger.ManagedLedgerException.MetaStoreException; import org.apache.bookkeeper.mledger.ManagedLedgerFactory; +import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; import org.apache.bookkeeper.mledger.impl.MetaStore.Stat; @@ -105,13 +107,6 @@ import org.slf4j.LoggerFactory; import org.testng.annotations.Test; -import com.google.common.base.Charsets; -import com.google.common.collect.Sets; - -import io.netty.buffer.ByteBuf; -import io.netty.buffer.PooledByteBufAllocator; -import io.netty.buffer.Unpooled; - public class ManagedLedgerTest extends MockedBookKeeperTestCase { private static final Logger log = LoggerFactory.getLogger(ManagedLedgerTest.class); @@ -1903,6 +1898,9 @@ public void testGetNextValidPosition() throws Exception { */ @Test public void testActiveDeactiveCursorWithDiscardEntriesFromCache() throws Exception { + ManagedLedgerFactoryConfig conf = new ManagedLedgerFactoryConfig(); + conf.setCacheEvictionFrequency(0.1); + ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(bkc, zkc, conf); ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("cache_eviction_ledger"); // Open Cursor also adds cursor into activeCursor-container @@ -1954,8 +1952,8 @@ public void testActiveDeactiveCursorWithDiscardEntriesFromCache() throws Excepti } // (3) Validate: cache should remove all entries read by both active cursors - log.info("expected, found : {}, {}", (5 * (totalInsertedEntries - readEntries)), entryCache.getSize()); - assertEquals((5 * (totalInsertedEntries - readEntries)), entryCache.getSize()); + log.info("expected, found : {}, {}", (5 * (totalInsertedEntries)), entryCache.getSize()); + assertEquals((5 * totalInsertedEntries), entryCache.getSize()); final int remainingEntries = totalInsertedEntries - readEntries; entries1 = cursor1.readEntries(remainingEntries); @@ -1968,7 +1966,7 @@ public void testActiveDeactiveCursorWithDiscardEntriesFromCache() throws Excepti // (4) Validate: cursor2 is active cursor and has not read these entries yet: so, cache should not remove these // entries - assertEquals((5 * (remainingEntries)), entryCache.getSize()); + assertEquals((5 * totalInsertedEntries), entryCache.getSize()); ledger.deactivateCursor(cursor2); @@ -1978,6 +1976,7 @@ public void testActiveDeactiveCursorWithDiscardEntriesFromCache() throws Excepti log.info("Finished reading entries"); ledger.close(); + factory.shutdown(); } @Test @@ -2017,7 +2016,8 @@ public void testActiveDeactiveCursor() throws Exception { entry.release(); } - // (3) Validate: cache discards all entries as read by active cursor + // (3) Validate: cache discards all entries after all cursors are deactivated + ledger.deactivateCursor(cursor1); assertEquals(0, entryCache.getSize()); ledger.close(); @@ -2041,81 +2041,6 @@ public void testCursorRecoveryForEmptyLedgers() throws Exception { assertEquals(c1.getMarkDeletedPosition(), ledger.lastConfirmedEntry); } - @Test - public void testBacklogCursor() throws Exception { - ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("cache_backlog_ledger"); - - final long maxMessageCacheRetentionTimeMillis = 100; - Field field = ManagedLedgerImpl.class.getDeclaredField("maxMessageCacheRetentionTimeMillis"); - field.setAccessible(true); - Field modifiersField = Field.class.getDeclaredField("modifiers"); - modifiersField.setAccessible(true); - modifiersField.setInt(field, field.getModifiers() & ~Modifier.FINAL); - field.set(ledger, maxMessageCacheRetentionTimeMillis); - Field backlogThresholdField = ManagedLedgerImpl.class.getDeclaredField("maxActiveCursorBacklogEntries"); - backlogThresholdField.setAccessible(true); - final long maxActiveCursorBacklogEntries = (long) backlogThresholdField.get(ledger); - - // Open Cursor also adds cursor into activeCursor-container - ManagedCursor cursor1 = ledger.openCursor("c1"); - ManagedCursor cursor2 = ledger.openCursor("c2"); - - final int totalBacklogSizeEntries = (int) maxActiveCursorBacklogEntries; - CountDownLatch latch = new CountDownLatch(totalBacklogSizeEntries); - for (int i = 0; i < totalBacklogSizeEntries + 1; i++) { - String content = "entry"; // 5 bytes - ByteBuf entry = getMessageWithMetadata(content.getBytes()); - ledger.asyncAddEntry(entry, new AddEntryCallback() { - @Override - public void addComplete(Position position, Object ctx) { - latch.countDown(); - entry.release(); - } - - @Override - public void addFailed(ManagedLedgerException exception, Object ctx) { - latch.countDown(); - entry.release(); - } - - }, null); - } - latch.await(); - - // Verify: cursors are active as :haven't started deactivateBacklogCursor scan - assertTrue(cursor1.isActive()); - assertTrue(cursor2.isActive()); - - // it allows message to be older enough to be considered in backlog - Thread.sleep(maxMessageCacheRetentionTimeMillis * 2); - - // deactivate backlog cursors - ledger.checkBackloggedCursors(); - Thread.sleep(100); - - // both cursors have to be inactive - assertFalse(cursor1.isActive()); - assertFalse(cursor2.isActive()); - - // read entries so, cursor1 reaches maxBacklog threshold again to be active again - List entries1 = cursor1.readEntries(50); - for (Entry entry : entries1) { - log.info("Read entry. Position={} Content='{}'", entry.getPosition(), new String(entry.getData())); - entry.release(); - } - - // activate cursors which caught up maxbacklog threshold - ledger.checkBackloggedCursors(); - - // verify: cursor1 has consumed messages so, under maxBacklog threshold => active - assertTrue(cursor1.isActive()); - - // verify: cursor2 has not consumed messages so, above maxBacklog threshold => inactive - assertFalse(cursor2.isActive()); - - ledger.close(); - } - @Test public void testConcurrentOpenCursor() throws Exception { ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testConcurrentOpenCursor"); @@ -2240,12 +2165,12 @@ public void testManagedLedgerWithoutAutoCreate() throws Exception { assertFalse(factory.getManagedLedgers().containsKey("testManagedLedgerWithoutAutoCreate")); } - + @Test public void testManagedLedgerWithCreateLedgerTimeOut() throws Exception { ManagedLedgerConfig config = new ManagedLedgerConfig().setMetadataOperationsTimeoutSeconds(3); ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("timeout_ledger_test", config); - + BookKeeper bk = mock(BookKeeper.class); doNothing().when(bk).asyncCreateLedger(anyInt(), anyInt(), anyInt(), any(), any(), any(), any(), any()); AtomicInteger response = new AtomicInteger(0); @@ -2260,13 +2185,13 @@ public void createComplete(int rc, LedgerHandle lh, Object ctx) { latch.await(config.getMetadataOperationsTimeoutSeconds() + 2, TimeUnit.SECONDS); assertEquals(response.get(), BKException.Code.TimeoutException); - + ledger.close(); } - + /** * It verifies that asyncRead timesout if it doesn't receive response from bk-client in configured timeout - * + * * @throws Exception */ @Test @@ -2339,8 +2264,8 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { /** * It verifies that if bk-client doesn't complete the add-entry in given time out then broker is resilient enought * to create new ledger and add entry successfully. - * - * + * + * * @throws Exception */ @Test(timeOut = 20000) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java index d1d2e5d2c6305..2cbeb6bced2fc 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java @@ -111,7 +111,7 @@ void simple() { @Test void customWeighter() { - RangeCache cache = new RangeCache<>(value -> value.s.length()); + RangeCache cache = new RangeCache<>(value -> value.s.length(), x -> 0); cache.put(0, new RefString("zero")); cache.put(1, new RefString("one")); @@ -120,6 +120,26 @@ void customWeighter() { assertEquals(cache.getNumberOfEntries(), 2); } + @Test + void customTimeExtraction() { + RangeCache cache = new RangeCache<>(value -> value.s.length(), x -> x.s.length()); + + cache.put(1, new RefString("1")); + cache.put(2, new RefString("22")); + cache.put(3, new RefString("333")); + cache.put(4, new RefString("4444")); + + assertEquals(cache.getSize(), 10); + assertEquals(cache.getNumberOfEntries(), 4); + + Pair p = cache.evictLEntriesBeforeTimestamp(3); + assertEquals(p.getLeft().intValue(), 3); + assertEquals(p.getRight().intValue(), 6); + + assertEquals(cache.getSize(), 4); + assertEquals(cache.getNumberOfEntries(), 1); + } + @Test void doubleInsert() { RangeCache cache = new RangeCache<>(); @@ -172,7 +192,7 @@ void getRange() { @Test void eviction() { - RangeCache cache = new RangeCache<>(value -> value.s.length()); + RangeCache cache = new RangeCache<>(value -> value.s.length(), x -> 0); cache.put(0, new RefString("zero")); cache.put(1, new RefString("one")); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java index 2b82204750ccb..2c1e78fc5b2be 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java @@ -47,6 +47,7 @@ public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient, managedLedgerFactoryConfig.setCacheEvictionWatermark(conf.getManagedLedgerCacheEvictionWatermark()); managedLedgerFactoryConfig.setNumManagedLedgerWorkerThreads(conf.getManagedLedgerNumWorkerThreads()); managedLedgerFactoryConfig.setNumManagedLedgerSchedulerThreads(conf.getManagedLedgerNumSchedulerThreads()); + managedLedgerFactoryConfig.setCacheEvictionFrequency(conf.getManagedLedgerCacheEvictionFrequency()); this.managedLedgerFactory = new ManagedLedgerFactoryImpl(bkClient, zkClient, managedLedgerFactoryConfig); } 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 34c06cf3fe5b0..fddda96351f72 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 @@ -736,7 +736,6 @@ public CompletableFuture getManagedLedgerConfig(TopicName t managedLedgerConfig.setThrottleMarkDelete(persistencePolicies.getManagedLedgerMaxMarkDeleteRate()); managedLedgerConfig.setDigestType(serviceConfig.getManagedLedgerDigestType()); - managedLedgerConfig.setCacheEvictionFrequency(serviceConfig.getManagedLedgerCacheEvictionFrequency()); managedLedgerConfig.setMaxUnackedRangesToPersist(serviceConfig.getManagedLedgerMaxUnackedRangesToPersist()); managedLedgerConfig.setMaxUnackedRangesToPersistInZk(serviceConfig.getManagedLedgerMaxUnackedRangesToPersistInZooKeeper()); managedLedgerConfig.setMaxEntriesPerLedger(serviceConfig.getManagedLedgerMaxEntriesPerLedger()); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarStats.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarStats.java index 678abf3bd5a67..d99dcfcb8c5cb 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarStats.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarStats.java @@ -136,9 +136,6 @@ public synchronized void updateStats( } catch (Exception e) { log.error("Failed to generate topic stats for topic {}: {}", name, e.getMessage(), e); } - // this task: helps to activate inactive-backlog-cursors which have caught up and - // connected, also deactivate active-backlog-cursors which has backlog - ((PersistentTopic) topic).getManagedLedger().checkBackloggedCursors(); }else if (topic instanceof NonPersistentTopic) { tempNonPersistentTopics.add((NonPersistentTopic) topic); } else { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index af9be1046787c..dfd64a8c8cb71 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -838,88 +838,6 @@ public void testActiveAndInActiveConsumerEntryCacheBehavior() throws Exception { log.info("-- Exiting {} test --", methodName); } - @Test - public void testDeactivatingBacklogConsumer() throws Exception { - log.info("-- Starting {} test --", methodName); - - final long batchMessageDelayMs = 100; - final int receiverSize = 10; - final String topicName = "cache-topic"; - final String topic = "persistent://my-property/my-ns/" + topicName; - final String sub1 = "faster-sub1"; - final String sub2 = "slower-sub2"; - - // 1. Subscriber Faster subscriber: let it consume all messages immediately - Consumer subscriber1 = pulsarClient.newConsumer() - .topic("persistent://my-property/my-ns/" + topicName).subscriptionName(sub1) - .subscriptionType(SubscriptionType.Shared).receiverQueueSize(receiverSize).subscribe(); - // 1.b. Subscriber Slow subscriber: - Consumer subscriber2 = pulsarClient.newConsumer() - .topic("persistent://my-property/my-ns/" + topicName).subscriptionName(sub2) - .subscriptionType(SubscriptionType.Shared).receiverQueueSize(receiverSize).subscribe(); - - ProducerBuilder producerBuilder = pulsarClient.newProducer().topic(topic); - if (batchMessageDelayMs != 0) { - producerBuilder.enableBatching(true).batchingMaxPublishDelay(batchMessageDelayMs, TimeUnit.MILLISECONDS) - .batchingMaxMessages(5); - } - Producer producer = producerBuilder.create(); - - PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topic).get(); - ManagedLedgerImpl ledger = (ManagedLedgerImpl) topicRef.getManagedLedger(); - - // reflection to set/get cache-backlog fields value: - final long maxMessageCacheRetentionTimeMillis = 100; - Field backlogThresholdField = ManagedLedgerImpl.class.getDeclaredField("maxActiveCursorBacklogEntries"); - backlogThresholdField.setAccessible(true); - Field field = ManagedLedgerImpl.class.getDeclaredField("maxMessageCacheRetentionTimeMillis"); - field.setAccessible(true); - Field modifiersField = Field.class.getDeclaredField("modifiers"); - modifiersField.setAccessible(true); - modifiersField.setInt(field, field.getModifiers() & ~Modifier.FINAL); - field.set(ledger, maxMessageCacheRetentionTimeMillis); - final long maxActiveCursorBacklogEntries = (long) backlogThresholdField.get(ledger); - - Message msg = null; - final int totalMsgs = (int) maxActiveCursorBacklogEntries + receiverSize + 1; - // 2. Produce messages - for (int i = 0; i < totalMsgs; i++) { - String message = "my-message-" + i; - producer.send(message.getBytes()); - } - // 3. Consume messages: at Faster subscriber - for (int i = 0; i < totalMsgs; i++) { - msg = subscriber1.receive(100, TimeUnit.MILLISECONDS); - subscriber1.acknowledge(msg); - } - - // wait : so message can be eligible to to be evict from cache - Thread.sleep(maxMessageCacheRetentionTimeMillis); - - // 4. deactivate subscriber which has built the backlog - ledger.checkBackloggedCursors(); - Thread.sleep(100); - - // 5. verify: active subscribers - Set activeSubscriber = Sets.newHashSet(); - ledger.getActiveCursors().forEach(c -> activeSubscriber.add(c.getName())); - assertTrue(activeSubscriber.contains(sub1)); - assertFalse(activeSubscriber.contains(sub2)); - - // 6. consume messages : at slower subscriber - for (int i = 0; i < totalMsgs; i++) { - msg = subscriber2.receive(100, TimeUnit.MILLISECONDS); - subscriber2.acknowledge(msg); - } - - ledger.checkBackloggedCursors(); - - activeSubscriber.clear(); - ledger.getActiveCursors().forEach(c -> activeSubscriber.add(c.getName())); - - assertTrue(activeSubscriber.contains(sub1)); - assertTrue(activeSubscriber.contains(sub2)); - } @Test(timeOut = 2000) public void testAsyncProducerAndConsumer() throws Exception { @@ -2911,16 +2829,16 @@ public void received(Consumer consumer, Message message) /** * This test verifies that broker activates fail-over consumer by considering priority-level as well. - * + * *
      * 1. Start two failover consumer with same priority level, broker selects consumer based on name-sorting (consumer1).
      * 2. Switch non-active consumer to active (consumer2): by giving it higher priority
      * Partitioned-topic with 9 partitions:
      * 1. C1 (priority=1)
      * 2. C2,C3,C4 (priority=0)
-     * So, broker should evenly distribute C2,C3,C4 active consumers among 9 partitions. 
+     * So, broker should evenly distribute C2,C3,C4 active consumers among 9 partitions.
      * 
- * + * * @throws Exception */ @Test diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java index 161c3742a8a06..b93d65c0fb4fc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java @@ -732,93 +732,6 @@ public void testActiveAndInActiveConsumerEntryCacheBehavior() throws Exception { log.info("-- Exiting {} test --", methodName); } - @Test - public void testDeactivatingBacklogConsumer() throws Exception { - log.info("-- Starting {} test --", methodName); - - final long batchMessageDelayMs = 100; - final int receiverSize = 10; - final String topicName = "cache-topic"; - final String topic = "persistent://my-property/use/my-ns/" + topicName; - final String sub1 = "faster-sub1"; - final String sub2 = "slower-sub2"; - - // 1. Subscriber Faster subscriber: let it consume all messages immediately - Consumer subscriber1 = pulsarClient.newConsumer() - .topic("persistent://my-property/use/my-ns/" + topicName) - .subscriptionName(sub1) - .subscriptionType(SubscriptionType.Shared) - .receiverQueueSize(receiverSize) - .subscribe(); - // 1.b. Subscriber Slow subscriber: - Consumer subscriber2 = pulsarClient.newConsumer() - .topic("persistent://my-property/use/my-ns/" + topicName) - .subscriptionName(sub2) - .subscriptionType(SubscriptionType.Shared) - .receiverQueueSize(receiverSize) - .subscribe(); - Producer producer = pulsarClient.newProducer() - .topic(topic) - .enableBatching(batchMessageDelayMs != 0) - .batchingMaxPublishDelay(batchMessageDelayMs, TimeUnit.MILLISECONDS) - .batchingMaxMessages(5) - .create(); - - PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topic).get(); - ManagedLedgerImpl ledger = (ManagedLedgerImpl) topicRef.getManagedLedger(); - - // reflection to set/get cache-backlog fields value: - final long maxMessageCacheRetentionTimeMillis = 100; - Field backlogThresholdField = ManagedLedgerImpl.class.getDeclaredField("maxActiveCursorBacklogEntries"); - backlogThresholdField.setAccessible(true); - Field field = ManagedLedgerImpl.class.getDeclaredField("maxMessageCacheRetentionTimeMillis"); - field.setAccessible(true); - Field modifiersField = Field.class.getDeclaredField("modifiers"); - modifiersField.setAccessible(true); - modifiersField.setInt(field, field.getModifiers() & ~Modifier.FINAL); - field.set(ledger, maxMessageCacheRetentionTimeMillis); - final long maxActiveCursorBacklogEntries = (long) backlogThresholdField.get(ledger); - - Messagemsg = null; - final int totalMsgs = (int) maxActiveCursorBacklogEntries + receiverSize + 1; - // 2. Produce messages - for (int i = 0; i < totalMsgs; i++) { - String message = "my-message-" + i; - producer.send(message.getBytes()); - } - // 3. Consume messages: at Faster subscriber - for (int i = 0; i < totalMsgs; i++) { - msg = subscriber1.receive(100, TimeUnit.MILLISECONDS); - subscriber1.acknowledge(msg); - } - - // wait : so message can be eligible to to be evict from cache - Thread.sleep(maxMessageCacheRetentionTimeMillis); - - // 4. deactivate subscriber which has built the backlog - ledger.checkBackloggedCursors(); - Thread.sleep(100); - - // 5. verify: active subscribers - Set activeSubscriber = Sets.newHashSet(); - ledger.getActiveCursors().forEach(c -> activeSubscriber.add(c.getName())); - assertTrue(activeSubscriber.contains(sub1)); - assertFalse(activeSubscriber.contains(sub2)); - - // 6. consume messages : at slower subscriber - for (int i = 0; i < totalMsgs; i++) { - msg = subscriber2.receive(100, TimeUnit.MILLISECONDS); - subscriber2.acknowledge(msg); - } - - ledger.checkBackloggedCursors(); - - activeSubscriber.clear(); - ledger.getActiveCursors().forEach(c -> activeSubscriber.add(c.getName())); - - assertTrue(activeSubscriber.contains(sub1)); - assertTrue(activeSubscriber.contains(sub2)); - } @Test(timeOut = 2000) public void testAsyncProducerAndConsumer() throws Exception { From 60cda0506fa80ff24c28fa06d23c4bc0aa920bb4 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Tue, 23 Apr 2019 16:42:31 -0700 Subject: [PATCH 04/15] Address comments --- conf/broker.conf | 4 ++ .../bookkeeper/mledger/ManagedLedger.java | 6 ++ .../mledger/ManagedLedgerFactoryConfig.java | 7 +- .../bookkeeper/mledger/impl/EntryImpl.java | 1 + .../impl/ManagedLedgerFactoryImpl.java | 10 +-- .../mledger/impl/ManagedLedgerImpl.java | 15 +++++ .../mledger/impl/ManagedLedgerTest.java | 66 +++++++++++++++++++ .../pulsar/broker/ServiceConfiguration.java | 4 ++ .../broker/ManagedLedgerClientFactory.java | 1 + 9 files changed, 104 insertions(+), 10 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 67c8cbcc1ce79..1a3e8753ec8ee 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -445,6 +445,10 @@ managedLedgerCacheEvictionWatermark=0.9 # Configure the cache eviction frequency for the managed ledger cache (evictions/sec) managedLedgerCacheEvictionFrequency=10.0 +# Configure the threshold (in number of entries) from where a cursor should be considered 'backlogged' +# and thus should be set as inactive. +managedLedgerCursorBackloggedThreshold=1000 + # Rate limit the amount of writes per second generated by consumer acking the messages managedLedgerDefaultMarkDeleteRateLimit=1.0 diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java index 36e542903df25..d51b0d803a7d0 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java @@ -322,6 +322,12 @@ public interface ManagedLedger { */ long getEstimatedBacklogSize(); + /** + * Activate cursors those caught up backlog-threshold entries and deactivate slow cursors which are creating + * backlog. + */ + void checkBackloggedCursors(); + void asyncTerminate(TerminateCallback callback, Object ctx); /** diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java index d7894792e7de0..4237127bb10b4 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java @@ -42,7 +42,8 @@ public class ManagedLedgerFactoryConfig { */ private double cacheEvictionFrequency = 10; - public long getMaxCacheSize() { - return maxCacheSize; - } + /** + * Threshould to consider a cursor as "backlogged" + */ + private long thresholdBackloggedCursor = 1000; } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java index 2a04f0e538569..58ea09923e64e 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java @@ -162,6 +162,7 @@ protected void deallocate() { // This method is called whenever the ref-count of the EntryImpl reaches 0, so that now we can recycle it data.release(); data = null; + timestamp = -1; ledgerId = -1; entryId = -1; recyclerHandle.recycle(this); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index 73a3be7d12a31..d5130198f4618 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -142,7 +142,7 @@ private ManagedLedgerFactoryImpl(BookKeeper bookKeeper, boolean isBookkeeperMana this.entryCacheManager = new EntryCacheManager(this); this.statsTask = scheduledExecutor.scheduleAtFixedRate(() -> refreshStats(), 0, StatsPeriodSeconds, TimeUnit.SECONDS); this.cacheEvictionTask = scheduledExecutor.scheduleAtFixedRate(() -> cacheEviction(), 0, - (long) (1000 / config.getCacheEvictionFrequency()), TimeUnit.MILLISECONDS); + (long) (1000 / Math.min(config.getCacheEvictionFrequency(), 1000.0)), TimeUnit.MILLISECONDS); } private synchronized void refreshStats() { @@ -163,20 +163,16 @@ private synchronized void refreshStats() { } private synchronized void cacheEviction() { - long now = System.nanoTime(); - long period = now - lastCacheEvictionTimestamp; - long maxTimestamp = now - period; - ledgers.values().forEach(mlfuture -> { if (mlfuture.isDone() && !mlfuture.isCompletedExceptionally()) { ManagedLedgerImpl ml = mlfuture.getNow(null); if (ml != null) { - ml.entryCache.invalidateEntriesBeforeTimestamp(maxTimestamp); + ml.entryCache.invalidateEntriesBeforeTimestamp(lastCacheEvictionTimestamp); } } }); - lastCacheEvictionTimestamp = now; + lastCacheEvictionTimestamp = System.nanoTime(); } /** 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 d883c740000d9..d584af51ae871 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 @@ -227,6 +227,8 @@ enum PositionBound { .newUpdater(ManagedLedgerImpl.class, "readOpCount"); private volatile long readOpCount = 0; + private final long backloggedCursorThresholdEntries; + /** * Queue of pending entries to be added to the managed ledger. Typically entries are queued when a new ledger is * created asynchronously and hence there is no ready ledger to write into. @@ -257,6 +259,7 @@ public ManagedLedgerImpl(ManagedLedgerFactoryImpl factory, BookKeeper bookKeeper this.waitingCursors = Queues.newConcurrentLinkedQueue(); this.uninitializedCursors = Maps.newHashMap(); this.clock = config.getClock(); + this.backloggedCursorThresholdEntries = factory.getConfig().getThresholdBackloggedCursor(); // Get the next rollover time. Add a random value upto 5% to avoid rollover multiple ledgers at the same time this.maximumRolloverTimeMs = (long) (config.getMaximumRolloverTimeMs() * (1 + random.nextDouble() * 5 / 100.0)); @@ -904,6 +907,18 @@ public long getTotalSize() { return TOTAL_SIZE_UPDATER.get(this); } + @Override + public void checkBackloggedCursors() { + // activate caught up cursors + cursors.forEach(cursor -> { + if (cursor.getNumberOfEntries() < backloggedCursorThresholdEntries) { + cursor.setActive(); + } else { + cursor.setInactive(); + } + }); + } + @Override public long getEstimatedBacklogSize() { diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index eefb0afb6639b..bd8b33ca093a9 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -39,6 +39,7 @@ import java.io.IOException; import java.lang.reflect.Field; +import java.lang.reflect.Modifier; import java.nio.charset.Charset; import java.security.GeneralSecurityException; import java.util.ArrayList; @@ -2041,6 +2042,71 @@ public void testCursorRecoveryForEmptyLedgers() throws Exception { assertEquals(c1.getMarkDeletedPosition(), ledger.lastConfirmedEntry); } + @Test + public void testBacklogCursor() throws Exception { + int backloggedThreshold = 10; + ManagedLedgerFactoryConfig factoryConf = new ManagedLedgerFactoryConfig(); + factoryConf.setThresholdBackloggedCursor(backloggedThreshold); + ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(bkc, zkc, factoryConf); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("cache_backlog_ledger"); + + // Open Cursor also adds cursor into activeCursor-container + ManagedCursor cursor1 = ledger.openCursor("c1"); + ManagedCursor cursor2 = ledger.openCursor("c2"); + + CountDownLatch latch = new CountDownLatch(backloggedThreshold); + for (int i = 0; i < backloggedThreshold + 1; i++) { + String content = "entry"; // 5 bytes + ByteBuf entry = getMessageWithMetadata(content.getBytes()); + ledger.asyncAddEntry(entry, new AddEntryCallback() { + @Override + public void addComplete(Position position, Object ctx) { + latch.countDown(); + entry.release(); + } + + @Override + public void addFailed(ManagedLedgerException exception, Object ctx) { + latch.countDown(); + entry.release(); + } + + }, null); + } + latch.await(); + + // Verify: cursors are active as :haven't started deactivateBacklogCursor scan + assertTrue(cursor1.isActive()); + assertTrue(cursor2.isActive()); + + // deactivate backlog cursors + ledger.checkBackloggedCursors(); + + // both cursors have to be inactive + assertFalse(cursor1.isActive()); + assertFalse(cursor2.isActive()); + + // read entries so, cursor1 reaches maxBacklog threshold again to be active again + List entries1 = cursor1.readEntries(50); + for (Entry entry : entries1) { + log.info("Read entry. Position={} Content='{}'", entry.getPosition(), new String(entry.getData())); + entry.release(); + } + + // activate cursors which caught up maxbacklog threshold + ledger.checkBackloggedCursors(); + + // verify: cursor1 has consumed messages so, under maxBacklog threshold => active + assertTrue(cursor1.isActive()); + + // verify: cursor2 has not consumed messages so, above maxBacklog threshold => inactive + assertFalse(cursor2.isActive()); + + ledger.close(); + + factory.shutdown(); + } + @Test public void testConcurrentOpenCursor() throws Exception { ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testConcurrentOpenCursor"); diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 5260f3b5cb335..41d0b30193d4f 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -774,6 +774,10 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext(category = CATEGORY_STORAGE_ML, doc = "Configure the cache eviction frequency for the managed ledger cache. Default is 100/s") private double managedLedgerCacheEvictionFrequency = 100.0; + @FieldContext(category = CATEGORY_STORAGE_ML, + doc = "Configure the threshold (in number of entries) from where a cursor should be considered 'backlogged'" + + " and thus should be set as inactive.") + private long managedLedgerCursorBackloggedThreshold = 1000; @FieldContext( category = CATEGORY_STORAGE_ML, doc = "Rate limit the amount of writes per second generated by consumer acking the messages" diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java index 2c1e78fc5b2be..e41015890e6ca 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java @@ -48,6 +48,7 @@ public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient, managedLedgerFactoryConfig.setNumManagedLedgerWorkerThreads(conf.getManagedLedgerNumWorkerThreads()); managedLedgerFactoryConfig.setNumManagedLedgerSchedulerThreads(conf.getManagedLedgerNumSchedulerThreads()); managedLedgerFactoryConfig.setCacheEvictionFrequency(conf.getManagedLedgerCacheEvictionFrequency()); + managedLedgerFactoryConfig.setThresholdBackloggedCursor(conf.getManagedLedgerCursorBackloggedThreshold()); this.managedLedgerFactory = new ManagedLedgerFactoryImpl(bkClient, zkClient, managedLedgerFactoryConfig); } From 6cca39b5ca7381835f52faccb71f8265630ed70b Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 24 Apr 2019 10:35:14 -0700 Subject: [PATCH 05/15] Apply eviction on slowest active reader by preference --- conf/broker.conf | 5 ++++- .../mledger/ManagedLedgerFactoryConfig.java | 9 +++++++-- .../mledger/impl/ManagedLedgerFactoryImpl.java | 12 ++++++++---- .../bookkeeper/mledger/impl/ManagedLedgerImpl.java | 11 +++++++++++ .../apache/pulsar/broker/ServiceConfiguration.java | 3 +++ 5 files changed, 33 insertions(+), 7 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 1a3e8753ec8ee..baa02b3ed74b0 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -443,7 +443,10 @@ managedLedgerCacheSizeMB= managedLedgerCacheEvictionWatermark=0.9 # Configure the cache eviction frequency for the managed ledger cache (evictions/sec) -managedLedgerCacheEvictionFrequency=10.0 +managedLedgerCacheEvictionFrequency=100.0 + +# All entries that have stayed in cache for more than the configured time, will be evicted +managedLedgerCacheEvictionTimeThresholdMillis=1000 # Configure the threshold (in number of entries) from where a cursor should be considered 'backlogged' # and thus should be set as inactive. diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java index 4237127bb10b4..40b815df76613 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java @@ -38,9 +38,14 @@ public class ManagedLedgerFactoryConfig { private int numManagedLedgerSchedulerThreads = Runtime.getRuntime().availableProcessors(); /** - * Frequency of cache eviction triggering. Default is 10 times per second. + * Frequency of cache eviction triggering. Default is 100 times per second. */ - private double cacheEvictionFrequency = 10; + private double cacheEvictionFrequency = 100; + + /** + * All entries that have stayed in cache for more than the configured time, will be evicted + */ + private long cacheEvictionTimeThresholdMillis = 1000; /** * Threshould to consider a cursor as "backlogged" diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index d5130198f4618..c228cc4626e00 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -88,9 +88,11 @@ public class ManagedLedgerFactoryImpl implements ManagedLedgerFactory { private final EntryCacheManager entryCacheManager; private long lastStatTimestamp = System.nanoTime(); - private long lastCacheEvictionTimestamp = System.nanoTime(); private final ScheduledFuture statsTask; private final ScheduledFuture cacheEvictionTask; + + private final long cacheEvictionTimeThresholdNanos; + private static final int StatsPeriodSeconds = 60; public ManagedLedgerFactoryImpl(ClientConfiguration bkClientConfiguration) throws Exception { @@ -143,6 +145,8 @@ private ManagedLedgerFactoryImpl(BookKeeper bookKeeper, boolean isBookkeeperMana this.statsTask = scheduledExecutor.scheduleAtFixedRate(() -> refreshStats(), 0, StatsPeriodSeconds, TimeUnit.SECONDS); this.cacheEvictionTask = scheduledExecutor.scheduleAtFixedRate(() -> cacheEviction(), 0, (long) (1000 / Math.min(config.getCacheEvictionFrequency(), 1000.0)), TimeUnit.MILLISECONDS); + this.cacheEvictionTimeThresholdNanos = TimeUnit.MILLISECONDS + .toNanos(config.getCacheEvictionTimeThresholdMillis()); } private synchronized void refreshStats() { @@ -163,16 +167,16 @@ private synchronized void refreshStats() { } private synchronized void cacheEviction() { + long maxTimestamp = System.nanoTime() - cacheEvictionTimeThresholdNanos; + ledgers.values().forEach(mlfuture -> { if (mlfuture.isDone() && !mlfuture.isCompletedExceptionally()) { ManagedLedgerImpl ml = mlfuture.getNow(null); if (ml != null) { - ml.entryCache.invalidateEntriesBeforeTimestamp(lastCacheEvictionTimestamp); + ml.doCacheEviction(maxTimestamp); } } }); - - lastCacheEvictionTimestamp = System.nanoTime(); } /** 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 d584af51ae871..e2784b407b956 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 @@ -1727,6 +1727,17 @@ void discardEntriesFromCache(ManagedCursorImpl cursor, PositionImpl newPosition) } } + void doCacheEviction(long maxTimestamp) { + // Always remove all entries already read by active cursors + PositionImpl slowestReaderPos = activeCursors.getSlowestReaderPosition(); + if (slowestReaderPos != null) { + entryCache.invalidateEntries(slowestReaderPos); + } + + // Remove entries older than the cutoff threshold + entryCache.invalidateEntriesBeforeTimestamp(maxTimestamp); + } + void updateCursor(ManagedCursorImpl cursor, PositionImpl newPosition) { Pair pair = cursors.cursorUpdated(cursor, newPosition); if (pair == null) { diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 41d0b30193d4f..ba790742f360e 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -774,6 +774,9 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext(category = CATEGORY_STORAGE_ML, doc = "Configure the cache eviction frequency for the managed ledger cache. Default is 100/s") private double managedLedgerCacheEvictionFrequency = 100.0; + @FieldContext(category = CATEGORY_STORAGE_ML, + doc = "All entries that have stayed in cache for more than the configured time, will be evicted") + private long managedLedgerCacheEvictionTimeThresholdMillis = 1000; @FieldContext(category = CATEGORY_STORAGE_ML, doc = "Configure the threshold (in number of entries) from where a cursor should be considered 'backlogged'" + " and thus should be set as inactive.") From 912799da4ed93c6730f2fc173c55335a5f4a1ffe Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 24 Apr 2019 12:09:43 -0700 Subject: [PATCH 06/15] Re-introduced backlogged subscriptions test --- .../api/SimpleProducerConsumerTest.java | 74 +++++++++++++++++++ 1 file changed, 74 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index dfd64a8c8cb71..dbe8f1963f372 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -838,6 +838,80 @@ public void testActiveAndInActiveConsumerEntryCacheBehavior() throws Exception { log.info("-- Exiting {} test --", methodName); } + @Test + public void testDeactivatingBacklogConsumer() throws Exception { + log.info("-- Starting {} test --", methodName); + + final long batchMessageDelayMs = 100; + final int receiverSize = 10; + final String topicName = "cache-topic"; + final String topic = "persistent://my-property/my-ns/" + topicName; + final String sub1 = "faster-sub1"; + final String sub2 = "slower-sub2"; + + // 1. Subscriber Faster subscriber: let it consume all messages immediately + Consumer subscriber1 = pulsarClient.newConsumer() + .topic("persistent://my-property/my-ns/" + topicName).subscriptionName(sub1) + .subscriptionType(SubscriptionType.Shared).receiverQueueSize(receiverSize).subscribe(); + // 1.b. Subscriber Slow subscriber: + Consumer subscriber2 = pulsarClient.newConsumer() + .topic("persistent://my-property/my-ns/" + topicName).subscriptionName(sub2) + .subscriptionType(SubscriptionType.Shared).receiverQueueSize(receiverSize).subscribe(); + + ProducerBuilder producerBuilder = pulsarClient.newProducer().topic(topic); + if (batchMessageDelayMs != 0) { + producerBuilder.enableBatching(true).batchingMaxPublishDelay(batchMessageDelayMs, TimeUnit.MILLISECONDS) + .batchingMaxMessages(5); + } + Producer producer = producerBuilder.create(); + + PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topic).get(); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) topicRef.getManagedLedger(); + + // reflection to set/get cache-backlog fields value: + final long maxMessageCacheRetentionTimeMillis = conf.getManagedLedgerCacheEvictionTimeThresholdMillis(); + final long maxActiveCursorBacklogEntries = conf.getManagedLedgerCursorBackloggedThreshold(); + + Message msg = null; + final int totalMsgs = (int) maxActiveCursorBacklogEntries + receiverSize + 1; + // 2. Produce messages + for (int i = 0; i < totalMsgs; i++) { + String message = "my-message-" + i; + producer.send(message.getBytes()); + } + // 3. Consume messages: at Faster subscriber + for (int i = 0; i < totalMsgs; i++) { + msg = subscriber1.receive(100, TimeUnit.MILLISECONDS); + subscriber1.acknowledge(msg); + } + + // wait : so message can be eligible to to be evict from cache + Thread.sleep(maxMessageCacheRetentionTimeMillis); + + // 4. deactivate subscriber which has built the backlog + ledger.checkBackloggedCursors(); + Thread.sleep(100); + + // 5. verify: active subscribers + Set activeSubscriber = Sets.newHashSet(); + ledger.getActiveCursors().forEach(c -> activeSubscriber.add(c.getName())); + assertTrue(activeSubscriber.contains(sub1)); + assertFalse(activeSubscriber.contains(sub2)); + + // 6. consume messages : at slower subscriber + for (int i = 0; i < totalMsgs; i++) { + msg = subscriber2.receive(100, TimeUnit.MILLISECONDS); + subscriber2.acknowledge(msg); + } + + ledger.checkBackloggedCursors(); + + activeSubscriber.clear(); + ledger.getActiveCursors().forEach(c -> activeSubscriber.add(c.getName())); + + assertTrue(activeSubscriber.contains(sub1)); + assertTrue(activeSubscriber.contains(sub2)); + } @Test(timeOut = 2000) public void testAsyncProducerAndConsumer() throws Exception { From f80ddf0cf4f12f6735c6cb64989e014178212d16 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 24 Apr 2019 15:54:04 -0700 Subject: [PATCH 07/15] Addressed comments --- .../org/apache/bookkeeper/mledger/impl/EntryCache.java | 2 +- .../apache/bookkeeper/mledger/impl/EntryCacheImpl.java | 6 +++--- .../mledger/impl/ManagedLedgerFactoryImpl.java | 4 +++- .../org/apache/bookkeeper/mledger/util/RangeCache.java | 10 ++++------ .../apache/bookkeeper/mledger/util/RangeCacheTest.java | 5 ++--- .../org/apache/pulsar/broker/service/PulsarStats.java | 3 +++ 6 files changed, 16 insertions(+), 14 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCache.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCache.java index 32ad3a0392e09..0e20a0facffb2 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCache.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCache.java @@ -50,7 +50,7 @@ public interface EntryCache extends Comparable { * Remove from cache all the entries related to a ledger up to lastPosition included. * * @param lastPosition - * the position of the last entry to be invalidated (inclusive) + * the position of the last entry to be invalidated (non-inclusive) */ void invalidateEntries(PositionImpl lastPosition); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java index 66c91c6d052b9..053e9002cc977 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java @@ -131,7 +131,7 @@ public boolean insert(EntryImpl entry) { public void invalidateEntries(final PositionImpl lastPosition) { final PositionImpl firstPosition = PositionImpl.get(-1, 0); - Pair removed = entries.removeRange(firstPosition, lastPosition, true); + Pair removed = entries.removeRange(firstPosition, lastPosition, false); int entriesRemoved = removed.getLeft(); long sizeRemoved = removed.getRight(); if (log.isDebugEnabled()) { @@ -342,8 +342,8 @@ public Pair evictEntries(long sizeToFree) { @Override public void invalidateEntriesBeforeTimestamp(long timestamp) { - Pair evicted = entries.evictLEntriesBeforeTimestamp(timestamp); - manager.entriesRemoved(evicted.getRight()); + long evictedSize = entries.evictLEntriesBeforeTimestamp(timestamp); + manager.entriesRemoved(evictedSize); } private static final Logger log = LoggerFactory.getLogger(EntryCacheImpl.class); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index c228cc4626e00..f9bf675ed1707 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -143,8 +143,10 @@ private ManagedLedgerFactoryImpl(BookKeeper bookKeeper, boolean isBookkeeperMana this.mbean = new ManagedLedgerFactoryMBeanImpl(this); this.entryCacheManager = new EntryCacheManager(this); this.statsTask = scheduledExecutor.scheduleAtFixedRate(() -> refreshStats(), 0, StatsPeriodSeconds, TimeUnit.SECONDS); + + double evictionFrequency = Math.max(Math.min(config.getCacheEvictionFrequency(), 1000.0), 0.001); this.cacheEvictionTask = scheduledExecutor.scheduleAtFixedRate(() -> cacheEviction(), 0, - (long) (1000 / Math.min(config.getCacheEvictionFrequency(), 1000.0)), TimeUnit.MILLISECONDS); + (long) (1000 / evictionFrequency), TimeUnit.MILLISECONDS); this.cacheEvictionTimeThresholdNanos = TimeUnit.MILLISECONDS .toNanos(config.getCacheEvictionTimeThresholdMillis()); } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java index dcf26b50048ed..e64a40de8aabb 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java @@ -179,12 +179,11 @@ public Pair evictLeastAccessedEntries(long minSize) { /** * - * @param minSize - * @return a pair containing the number of entries evicted and their total size + * @param maxTimestamp the max timestamp of the entries to be evicted + * @return the tota */ - public Pair evictLEntriesBeforeTimestamp(long maxTimestamp) { + public long evictLEntriesBeforeTimestamp(long maxTimestamp) { long removedSize = 0; - int removedEntries = 0; while (true) { Map.Entry entry = entries.firstEntry(); @@ -198,13 +197,12 @@ public Pair evictLEntriesBeforeTimestamp(long maxTimestamp) { } Value value = entry.getValue(); - ++removedEntries; removedSize += weighter.getSize(value); value.release(); } size.addAndGet(-removedSize); - return Pair.of(removedEntries, removedSize); + return removedSize; } /** diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java index 2cbeb6bced2fc..27e63815c0107 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java @@ -132,9 +132,8 @@ void customTimeExtraction() { assertEquals(cache.getSize(), 10); assertEquals(cache.getNumberOfEntries(), 4); - Pair p = cache.evictLEntriesBeforeTimestamp(3); - assertEquals(p.getLeft().intValue(), 3); - assertEquals(p.getRight().intValue(), 6); + long evictedSize = cache.evictLEntriesBeforeTimestamp(3); + assertEquals(evictedSize, 6); assertEquals(cache.getSize(), 4); assertEquals(cache.getNumberOfEntries(), 1); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarStats.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarStats.java index d99dcfcb8c5cb..678abf3bd5a67 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarStats.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarStats.java @@ -136,6 +136,9 @@ public synchronized void updateStats( } catch (Exception e) { log.error("Failed to generate topic stats for topic {}: {}", name, e.getMessage(), e); } + // this task: helps to activate inactive-backlog-cursors which have caught up and + // connected, also deactivate active-backlog-cursors which has backlog + ((PersistentTopic) topic).getManagedLedger().checkBackloggedCursors(); }else if (topic instanceof NonPersistentTopic) { tempNonPersistentTopics.add((NonPersistentTopic) topic); } else { From 7e635929c15269be2609a3d76e48c8658aeba0c4 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 24 Apr 2019 16:47:27 -0700 Subject: [PATCH 08/15] Use config option --- .../org/apache/pulsar/broker/ManagedLedgerClientFactory.java | 1 + 1 file changed, 1 insertion(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java index e41015890e6ca..e3a52d8d7d5fe 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java @@ -48,6 +48,7 @@ public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient, managedLedgerFactoryConfig.setNumManagedLedgerWorkerThreads(conf.getManagedLedgerNumWorkerThreads()); managedLedgerFactoryConfig.setNumManagedLedgerSchedulerThreads(conf.getManagedLedgerNumSchedulerThreads()); managedLedgerFactoryConfig.setCacheEvictionFrequency(conf.getManagedLedgerCacheEvictionFrequency()); + managedLedgerFactoryConfig.setCacheEvictionTimeThresholdMillis(conf.getManagedLedgerCacheEvictionTimeThresholdMillis()); managedLedgerFactoryConfig.setThresholdBackloggedCursor(conf.getManagedLedgerCursorBackloggedThreshold()); this.managedLedgerFactory = new ManagedLedgerFactoryImpl(bkClient, zkClient, managedLedgerFactoryConfig); From 332f2541871f61c4971b39cf798dfbd5e463f45d Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 24 Apr 2019 17:52:30 -0700 Subject: [PATCH 09/15] Fixed active/inactive logic and read position --- .../mledger/impl/ManagedLedgerImpl.java | 16 +++++++++++++++- .../service/persistent/PersistentTopic.java | 3 +++ 2 files changed, 18 insertions(+), 1 deletion(-) 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 54b37cf96d758..ca8610fdea85f 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 @@ -1717,7 +1717,7 @@ void discardEntriesFromCache(ManagedCursorImpl cursor, PositionImpl newPosition) void doCacheEviction(long maxTimestamp) { // Always remove all entries already read by active cursors - PositionImpl slowestReaderPos = activeCursors.getSlowestReaderPosition(); + PositionImpl slowestReaderPos = getEarlierReadPositionForActiveCursors(); if (slowestReaderPos != null) { entryCache.invalidateEntries(slowestReaderPos); } @@ -1726,6 +1726,20 @@ void doCacheEviction(long maxTimestamp) { entryCache.invalidateEntriesBeforeTimestamp(maxTimestamp); } + private PositionImpl getEarlierReadPositionForActiveCursors() { + PositionImpl smallest = null; + for (ManagedCursor cursor : activeCursors) { + PositionImpl p = (PositionImpl) cursor.getReadPosition(); + if (smallest == null) { + smallest = p; + } else if (p.compareTo(smallest) < 0) { + smallest = p; + } + } + + return smallest; + } + void updateCursor(ManagedCursorImpl cursor, PositionImpl newPosition) { Pair pair = cursors.cursorUpdated(cursor, newPosition); if (pair == null) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index e1a050979e282..31de1237fc83f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -556,9 +556,12 @@ public CompletableFuture subscribe(final ServerCnx cnx, String subscri subscriptionFuture.thenAccept(subscription -> { try { + ledger.checkBackloggedCursors(); + Consumer consumer = new Consumer(subscription, subType, topic, consumerId, priorityLevel, consumerName, maxUnackedMessages, cnx, cnx.getRole(), metadata, readCompacted, initialPosition); subscription.addConsumer(consumer); + if (!cnx.isActive()) { consumer.close(); if (log.isDebugEnabled()) { From 81b27dfe72ee0e10c28ff12576073aa78f272080 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Thu, 25 Apr 2019 12:14:12 -0700 Subject: [PATCH 10/15] Use dedicated thread for cache evictions --- .../impl/ManagedLedgerFactoryImpl.java | 38 ++++++++++++++++--- 1 file changed, 32 insertions(+), 6 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index f9bf675ed1707..7399233f1694f 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -24,6 +24,8 @@ import com.google.common.base.Predicates; import com.google.common.collect.Maps; +import io.netty.util.concurrent.DefaultThreadFactory; + import java.util.ArrayList; import java.util.List; import java.util.Map; @@ -32,6 +34,8 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentSkipListMap; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -82,6 +86,8 @@ public class ManagedLedgerFactoryImpl implements ManagedLedgerFactory { protected final OrderedScheduler scheduledExecutor; private final OrderedExecutor orderedExecutor; + private final ExecutorService cacheEvictionExecutor; + protected final ManagedLedgerFactoryMBeanImpl mbean; protected final ConcurrentHashMap> ledgers = new ConcurrentHashMap<>(); @@ -89,7 +95,6 @@ public class ManagedLedgerFactoryImpl implements ManagedLedgerFactory { private long lastStatTimestamp = System.nanoTime(); private final ScheduledFuture statsTask; - private final ScheduledFuture cacheEvictionTask; private final long cacheEvictionTimeThresholdNanos; @@ -134,6 +139,8 @@ private ManagedLedgerFactoryImpl(BookKeeper bookKeeper, boolean isBookkeeperMana .numThreads(config.getNumManagedLedgerWorkerThreads()) .name("bookkeeper-ml-workers") .build(); + cacheEvictionExecutor = Executors + .newSingleThreadExecutor(new DefaultThreadFactory("bookkeeper-ml-cache-eviction")); this.bookKeeper = bookKeeper; this.isBookkeeperManaged = isBookkeeperManaged; @@ -144,11 +151,12 @@ private ManagedLedgerFactoryImpl(BookKeeper bookKeeper, boolean isBookkeeperMana this.entryCacheManager = new EntryCacheManager(this); this.statsTask = scheduledExecutor.scheduleAtFixedRate(() -> refreshStats(), 0, StatsPeriodSeconds, TimeUnit.SECONDS); - double evictionFrequency = Math.max(Math.min(config.getCacheEvictionFrequency(), 1000.0), 0.001); - this.cacheEvictionTask = scheduledExecutor.scheduleAtFixedRate(() -> cacheEviction(), 0, - (long) (1000 / evictionFrequency), TimeUnit.MILLISECONDS); + this.cacheEvictionTimeThresholdNanos = TimeUnit.MILLISECONDS .toNanos(config.getCacheEvictionTimeThresholdMillis()); + + + cacheEvictionExecutor.execute(this::cacheEvictionTask); } private synchronized void refreshStats() { @@ -168,7 +176,25 @@ private synchronized void refreshStats() { lastStatTimestamp = now; } - private synchronized void cacheEviction() { + private void cacheEvictionTask() { + double evictionFrequency = Math.max(Math.min(config.getCacheEvictionFrequency(), 1000.0), 0.001); + long waitTimeMillis = (long) (1000 / evictionFrequency); + + while (true) { + try { + doCacheEviction(); + + Thread.sleep(waitTimeMillis); + } catch (InterruptedException e) { + // Factory is shutting down + return; + } catch (Throwable t) { + log.warn("Exception while performing cache eviction: {}", t.getMessage(), t); + } + } + } + + private synchronized void doCacheEviction() { long maxTimestamp = System.nanoTime() - cacheEvictionTimeThresholdNanos; ledgers.values().forEach(mlfuture -> { @@ -348,7 +374,6 @@ void close(ManagedLedger ledger) { @Override public void shutdown() throws InterruptedException, ManagedLedgerException { statsTask.cancel(true); - cacheEvictionTask.cancel(true); int numLedgers = ledgers.size(); final CountDownLatch latch = new CountDownLatch(numLedgers); @@ -392,6 +417,7 @@ public void closeFailed(ManagedLedgerException exception, Object ctx) { scheduledExecutor.shutdown(); orderedExecutor.shutdown(); + cacheEvictionExecutor.shutdownNow(); entryCacheManager.clear(); } From 038757a7834d4baf40b2afec4a8084b50ad0d94d Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Thu, 25 Apr 2019 12:18:45 -0700 Subject: [PATCH 11/15] Added config options in docs --- conf/standalone.conf | 10 ++++++++++ site2/docs/reference-configuration.md | 3 +++ 2 files changed, 13 insertions(+) diff --git a/conf/standalone.conf b/conf/standalone.conf index 4cbf931c75748..b8f13c51a0489 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -301,6 +301,16 @@ managedLedgerCacheSizeMB= # Threshold to which bring down the cache level when eviction is triggered managedLedgerCacheEvictionWatermark=0.9 +# Configure the cache eviction frequency for the managed ledger cache (evictions/sec) +managedLedgerCacheEvictionFrequency=100.0 + +# All entries that have stayed in cache for more than the configured time, will be evicted +managedLedgerCacheEvictionTimeThresholdMillis=1000 + +# Configure the threshold (in number of entries) from where a cursor should be considered 'backlogged' +# and thus should be set as inactive. +managedLedgerCursorBackloggedThreshold=1000 + # Rate limit the amount of writes generated by consumer acking the messages managedLedgerDefaultMarkDeleteRateLimit=0.1 diff --git a/site2/docs/reference-configuration.md b/site2/docs/reference-configuration.md index 959f546128c1b..1f828adc17a12 100644 --- a/site2/docs/reference-configuration.md +++ b/site2/docs/reference-configuration.md @@ -176,6 +176,9 @@ Pulsar brokers are responsible for handling incoming messages from producers, di |managedLedgerDefaultAckQuorum| Number of guaranteed copies (acks to wait before write is complete) |2| |managedLedgerCacheSizeMB| Amount of memory to use for caching data payload in managed ledger. This memory is allocated from JVM direct memory and it’s shared across all the topics running in the same broker. By default, uses 1/5th of available direct memory || |managedLedgerCacheEvictionWatermark| Threshold to which bring down the cache level when eviction is triggered |0.9| +|managedLedgerCacheEvictionFrequency| Configure the cache eviction frequency for the managed ledger cache (evictions/sec) | 100.0 | +|managedLedgerCacheEvictionTimeThresholdMillis| All entries that have stayed in cache for more than the configured time, will be evicted | 1000 | +|managedLedgerCursorBackloggedThreshold| Configure the threshold (in number of entries) from where a cursor should be considered 'backlogged' and thus should be set as inactive. | 1000| |managedLedgerDefaultMarkDeleteRateLimit| Rate limit the amount of writes per second generated by consumer acking the messages |1.0| |managedLedgerMaxEntriesPerLedger| Max number of entries to append to a ledger before triggering a rollover. A ledger rollover is triggered on these conditions:
  • Either the max rollover time has been reached
  • or max entries have been written to the ledged and at least min-time has passed
|50000| |managedLedgerMinLedgerRolloverTimeMinutes| Minimum time between ledger rollover for a topic |10| From 358d614b1e5d31d76b1ad9a91b6d2b5e3b7f54b4 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Thu, 25 Apr 2019 16:38:12 -0700 Subject: [PATCH 12/15] Fixed tests --- .../mledger/impl/EntryCacheManagerTest.java | 11 +++++------ .../bookkeeper/mledger/impl/ManagedLedgerTest.java | 12 ++++++------ 2 files changed, 11 insertions(+), 12 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java index d3f5733fff74c..bbd85f41f7b81 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java @@ -28,7 +28,6 @@ import org.apache.bookkeeper.common.util.OrderedScheduler; import org.apache.bookkeeper.mledger.Entry; -import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; import org.testng.annotations.BeforeClass; @@ -100,15 +99,15 @@ void simple() throws Exception { assertEquals(cacheManager.getSize(), 3); assertEquals(cache2.getSize(), 3); - // Should remove 2 entries + // Should remove 1 entry cache2.invalidateEntries(new PositionImpl(2, 1)); - assertEquals(cacheManager.getSize(), 1); - assertEquals(cache2.getSize(), 1); + assertEquals(cacheManager.getSize(), 2); + assertEquals(cache2.getSize(), 2); cacheManager.mlFactoryMBean.refreshStats(1, TimeUnit.SECONDS); assertEquals(cacheManager.mlFactoryMBean.getCacheMaxSize(), 10); - assertEquals(cacheManager.mlFactoryMBean.getCacheUsedSize(), 1); + assertEquals(cacheManager.mlFactoryMBean.getCacheUsedSize(), 2); assertEquals(cacheManager.mlFactoryMBean.getCacheHitsRate(), 0.0); assertEquals(cacheManager.mlFactoryMBean.getCacheMissesRate(), 0.0); assertEquals(cacheManager.mlFactoryMBean.getCacheHitsThroughput(), 0.0); @@ -265,7 +264,7 @@ void verifyHitsMisses() throws Exception { entries.forEach(e -> e.release()); cacheManager.mlFactoryMBean.refreshStats(1, TimeUnit.SECONDS); - assertEquals(cacheManager.mlFactoryMBean.getCacheUsedSize(), 0); + assertEquals(cacheManager.mlFactoryMBean.getCacheUsedSize(), 7); assertEquals(cacheManager.mlFactoryMBean.getCacheHitsRate(), 0.0); assertEquals(cacheManager.mlFactoryMBean.getCacheMissesRate(), 0.0); assertEquals(cacheManager.mlFactoryMBean.getCacheHitsThroughput(), 0.0); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index 0383f0cb4dc29..257d710b44b8f 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -1352,21 +1352,21 @@ public void invalidateConsumedEntriesFromCache() throws Exception { assertEquals(cacheManager.getSize(), entryCache.getSize()); c1.setReadPosition(p2); - ledger.discardEntriesFromCache(c1, p1); + ledger.discardEntriesFromCache(c1, p2); assertEquals(entryCache.getSize(), 7 * 3); assertEquals(cacheManager.getSize(), entryCache.getSize()); c1.setReadPosition(p3); - ledger.discardEntriesFromCache(c1, p2); - assertEquals(entryCache.getSize(), 7 * 2); + ledger.discardEntriesFromCache(c1, p3); + assertEquals(entryCache.getSize(), 7 * 3); assertEquals(cacheManager.getSize(), entryCache.getSize()); ledger.deactivateCursor(c1); - assertEquals(entryCache.getSize(), 7 * 2); // as c2.readPosition=p3 => Cache contains p3,p4 + assertEquals(entryCache.getSize(), 7 * 3); // as c2.readPosition=p3 => Cache contains p3,p4 assertEquals(cacheManager.getSize(), entryCache.getSize()); c2.setReadPosition(p4); - ledger.discardEntriesFromCache(c2, p3); + ledger.discardEntriesFromCache(c2, p4); assertEquals(entryCache.getSize(), 7); assertEquals(cacheManager.getSize(), entryCache.getSize()); @@ -2396,7 +2396,7 @@ private void setFieldValue(Class clazz, Object classObj, String fieldName, Objec field.setAccessible(true); field.set(classObj, fieldValue); } - + public static void retryStrategically(Predicate predicate, int retryCount, long intSleepTimeInMillis) throws Exception { for (int i = 0; i < retryCount; i++) { From 61f3f52c990a98498e63527b2e173abfd88eb4dc Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Tue, 30 Apr 2019 12:36:36 -0700 Subject: [PATCH 13/15] Added time triggered eviction test --- .../mledger/impl/EntryCacheManagerTest.java | 40 +++++++++++++++++++ 1 file changed, 40 insertions(+) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java index bbd85f41f7b81..69c23e323439e 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryCacheManagerTest.java @@ -28,6 +28,9 @@ import org.apache.bookkeeper.common.util.OrderedScheduler; import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.bookkeeper.mledger.ManagedLedger; +import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; import org.testng.annotations.BeforeClass; @@ -270,4 +273,41 @@ void verifyHitsMisses() throws Exception { assertEquals(cacheManager.mlFactoryMBean.getCacheHitsThroughput(), 0.0); assertEquals(cacheManager.mlFactoryMBean.getNumberOfCacheEvictions(), 0); } + + @Test + void verifyTimeBasedEviction() throws Exception { + ManagedLedgerFactoryConfig config = new ManagedLedgerFactoryConfig(); + config.setMaxCacheSize(1000); + config.setCacheEvictionFrequency(100); + config.setCacheEvictionTimeThresholdMillis(100); + + ManagedLedgerFactoryImpl factory = new ManagedLedgerFactoryImpl(bkc, bkc.getZkHandle(), config); + + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("test"); + ManagedCursor c1 = ledger.openCursor("c1"); + c1.setActive(); + ManagedCursor c2 = ledger.openCursor("c2"); + c2.setActive(); + + EntryCacheManager cacheManager = factory.getEntryCacheManager(); + assertEquals(cacheManager.getSize(), 0); + + EntryCache cache = cacheManager.getEntryCache(ledger); + assertEquals(cache.getSize(), 0); + + ledger.addEntry(new byte[4]); + ledger.addEntry(new byte[3]); + + // Cache eviction should happen every 10 millis and clean all the entries older that 100ms + Thread.sleep(1000); + + c1.close(); + c2.close(); + + assertEquals(cacheManager.getSize(), 0); + assertEquals(cache.getSize(), 0); + + factory.shutdown(); + } + } From ea2381fffad6b42b86f42ad9b2eb5206a0484580 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 1 May 2019 11:07:48 -0700 Subject: [PATCH 14/15] Fixed flaky test --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index 257d710b44b8f..35725e671bfcb 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -1960,6 +1960,7 @@ public void testActiveDeactiveCursorWithDiscardEntriesFromCache() throws Excepti entries1 = cursor1.readEntries(remainingEntries); // Acknowledge only on last entry cursor1.markDelete(entries1.get(entries1.size() - 1).getPosition()); + for (Entry entry : entries1) { log.info("Read entry. Position={} Content='{}'", entry.getPosition(), new String(entry.getData())); entry.release(); @@ -1969,10 +1970,11 @@ public void testActiveDeactiveCursorWithDiscardEntriesFromCache() throws Excepti // entries assertEquals((5 * totalInsertedEntries), entryCache.getSize()); + ledger.deactivateCursor(cursor1); ledger.deactivateCursor(cursor2); // (5) Validate: cursor2 is not active cursor now: cache should have removed all entries read by active cursor1 - assertEquals(0, entryCache.getSize()); + assertEquals(entryCache.getSize(), 0); log.info("Finished reading entries"); From a6eeefbc5b2d5c3b80794a933cb71325db6b3004 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 1 May 2019 15:26:07 -0700 Subject: [PATCH 15/15] Fixed tests --- .../org/apache/pulsar/client/api/ConsumerRedeliveryTest.java | 5 +++-- .../apache/pulsar/client/api/SimpleProducerConsumerTest.java | 3 --- .../apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java | 3 --- 3 files changed, 3 insertions(+), 8 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ConsumerRedeliveryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ConsumerRedeliveryTest.java index fc05d986a7c3b..63c974dec5ee6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ConsumerRedeliveryTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ConsumerRedeliveryTest.java @@ -36,6 +36,7 @@ public class ConsumerRedeliveryTest extends ProducerConsumerBase { @BeforeClass @Override protected void setup() throws Exception { + conf.setManagedLedgerCacheEvictionFrequency(0.1); super.internalSetup(); super.producerBaseSetup(); } @@ -62,7 +63,7 @@ public void testOrderedRedelivery() throws Exception { conf.setManagedLedgerMaxEntriesPerLedger(2); conf.setManagedLedgerMinLedgerRolloverTimeMinutes(0); - + ProducerBuilder producerBuilder = pulsarClient.newProducer().topic(topic) .producerName("my-producer-name"); Producer producer = producerBuilder.create(); @@ -99,7 +100,7 @@ public void testOrderedRedelivery() throws Exception { Message message = consumer1.receive(5, TimeUnit.SECONDS); MessageIdImpl msgId = (MessageIdImpl) message.getMessageId(); if (lastMsgId != null) { - assertTrue(lastMsgId.getLedgerId() <= msgId.getLedgerId()); + assertTrue(lastMsgId.getLedgerId() <= msgId.getLedgerId(), "lastMsgId: " + lastMsgId + " -- msgId: " + msgId); } lastMsgId = msgId; } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index 214ca4664eebb..7192e7fa6dd6a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -796,9 +796,6 @@ public void testActiveAndInActiveConsumerEntryCacheBehavior() throws Exception { producer.send("message".getBytes()); msg = subscriber1.receive(5, TimeUnit.SECONDS); - // Verify: cache has to be cleared as there is no message needs to be consumed by active subscriber - assertEquals(entryCache.getSize(), 0, 1); - /************ usecase-2: *************/ // 1.b Subscriber slower-subscriber Consumer subscriber2 = pulsarClient.newConsumer() diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java index b93d65c0fb4fc..fecef7946d5c1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java @@ -687,9 +687,6 @@ public void testActiveAndInActiveConsumerEntryCacheBehavior() throws Exception { producer.send("message".getBytes()); msg = subscriber1.receive(5, TimeUnit.SECONDS); - // Verify: cache has to be cleared as there is no message needs to be consumed by active subscriber - assertEquals(entryCache.getSize(), 0, 1); - /************ usecase-2: *************/ // 1.b Subscriber slower-subscriber Consumer subscriber2 = pulsarClient.newConsumer()