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 381da6ab9eeae..3cb360efd1701 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 @@ -57,8 +57,10 @@ protected EntryImpl newObject(Handle handle) { ByteBuf data; private EntryReadCountHandler readCountHandler; private boolean decreaseReadCountOnRelease = true; + // Cache readers publish metadata lazily; entry copies must see a fully initialized instance. @Getter @Setter - private MessageMetadata messageMetadata; + private volatile MessageMetadata messageMetadata; + private boolean messageMetadataInitializationFailed; private Runnable onDeallocate; @@ -290,6 +292,7 @@ protected void deallocate() { readCountHandler = null; decreaseReadCountOnRelease = true; messageMetadata = null; + messageMetadataInitializationFailed = false; recyclerHandle.recycle(this); } @@ -308,12 +311,14 @@ public void setDecreaseReadCountOnRelease(boolean enabled) { } public synchronized void initializeMessageMetadataIfNeeded(String managedLedgerName) { - if (messageMetadata == null) { + if (messageMetadata == null && !messageMetadataInitializationFailed) { try { MessageMetadata msgMetadata = new MessageMetadata(); Commands.parseMessageMetadata(data.duplicate(), msgMetadata); this.messageMetadata = msgMetadata; } catch (Throwable t) { + // The entry bytes are immutable; another cache reader cannot make a failed parse succeed. + messageMetadataInitializationFailed = true; log.warn().attr("managedLedgerName", managedLedgerName) .attr("ledgerId", ledgerId) .attr("entryId", entryId) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java index 66b01c851e065..6375465a13eb0 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java @@ -169,10 +169,6 @@ public boolean insert(Entry entry, boolean copy) { EntryImpl cacheEntry = EntryImpl.createWithRetainedDuplicate(position, cachedData, entry.getReadCountHandler(), copy ? null : entry.getMessageMetadata()); - if (ml.getConfig().isPulsarMessageEntries()) { - // Parse the message metadata once at insert time so that cache reads don't have to do it lazily - cacheEntry.initializeMessageMetadataIfNeeded(ml.getName()); - } cachedData.release(); if (entries.put(position, cacheEntry, entryLength)) { totalAddedEntriesSize.add(entryLength); @@ -404,7 +400,8 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx2) { void doAsyncReadEntriesByPosition(ReadHandle lh, Position firstPosition, Position lastPosition, int numberOfEntries, long maxSizeBytes, IntSupplier expectedReadCount, final ReadEntriesCallback callback, Object ctx) { - CachedEntries cachedEntries = new CachedEntries(firstPosition.getEntryId(), numberOfEntries); + CachedEntries cachedEntries = new CachedEntries(firstPosition.getEntryId(), numberOfEntries, + ml.getConfig().isPulsarMessageEntries() ? ml.getName() : null); if (firstPosition.compareTo(lastPosition) == 0) { ReferenceCountedEntry cachedEntry = entries.get(firstPosition); if (cachedEntry != null) { @@ -508,13 +505,15 @@ void doAsyncReadEntriesByPosition(ReadHandle lh, Position firstPosition, Positio static final class CachedEntries implements Consumer { private final long firstEntryId; private final int numberOfEntries; + private final String managedLedgerName; List entries; private int count; private long totalSize; - CachedEntries(long firstEntryId, int numberOfEntries) { + CachedEntries(long firstEntryId, int numberOfEntries, String managedLedgerName) { this.firstEntryId = firstEntryId; this.numberOfEntries = numberOfEntries; + this.managedLedgerName = managedLedgerName; } @Override @@ -525,6 +524,11 @@ public void accept(ReferenceCountedEntry entry) { entries.add(null); } } + // The visitor retains the cached entry while parsing. Initialize on the shared cached entry + // before copying, so fanout readers reuse one instance backed by the cache-owned buffer. + if (managedLedgerName != null && entry.getMessageMetadata() == null) { + ((EntryImpl) entry).initializeMessageMetadataIfNeeded(managedLedgerName); + } int index = (int) (entry.getPosition().getEntryId() - firstEntryId); entries.set(index, EntryImpl.create(entry)); count++; diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryImplTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryImplTest.java index eac61f5d88ec3..e8742b92811ad 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryImplTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryImplTest.java @@ -18,7 +18,11 @@ */ package org.apache.bookkeeper.mledger.impl; +import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotNull; @@ -34,6 +38,24 @@ public class EntryImplTest { + @Test + public void testFailedMetadataInitializationIsNotRetried() { + ByteBuf bytes = Unpooled.buffer(4).writeInt(-1); + EntryImpl entry = EntryImpl.create(1, 0, bytes); + bytes.release(); + entry.data = spy(entry.data); + try { + entry.initializeMessageMetadataIfNeeded("ledger"); + entry.initializeMessageMetadataIfNeeded("ledger"); + assertThat(entry.getMessageMetadata()).isNull(); + assertThat(entry.getDataBuffer().readerIndex()).isZero(); + assertThat(entry.getDataBuffer().getInt(0)).isEqualTo(-1); + verify(entry.data, times(1)).duplicate(); + } finally { + entry.release(); + } + } + @Test public void testCreateWithLedgerIdEntryIdAndByteBuf() { // Given diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImplTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImplTest.java index d6712e613d553..b0f3ca2197909 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImplTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImplTest.java @@ -39,6 +39,10 @@ import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.IntSupplier; import org.apache.bookkeeper.client.LedgerHandle; @@ -173,7 +177,7 @@ private static ByteBuf serializeMessage(String producerName) { } @Test - public void testInsertParsesMessageMetadata() { + public void testInsertDefersMetadataUntilFirstReadAndSharesIt() { managedLedgerConfig.setPulsarMessageEntries(true); ByteBuf headersAndPayload = serializeMessage("producer"); EntryImpl entry = EntryImpl.create(1, 50, headersAndPayload); @@ -183,14 +187,78 @@ public void testInsertParsesMessageMetadata() { assertThat(rangeEntryCache.insert(entry)).isTrue(); entry.release(); - // the metadata is parsed once at insert time. Reading the entry back out of the cache doesn't parse - // anything any more, so this asserts what insert actually stored ReferenceCountedEntry cached = rangeEntryCache.getEntries().get(PositionFactory.create(1, 50)); assertThat(cached).isNotNull(); - assertThat(cached.getMessageMetadata()).isNotNull(); - assertThat(cached.getMessageMetadata().getProducerName()).isEqualTo("producer"); - assertThat(cached.getMessageMetadata().getSequenceId()).isEqualTo(7); - cached.release(); + assertThat(cached.getMessageMetadata()).isNull(); + Entry first = readSingleEntryFromCache(1, 50); + Entry second = readSingleEntryFromCache(1, 50); + try { + try { + assertThat(first.getMessageMetadata()).isNotNull().isSameAs(cached.getMessageMetadata()); + assertThat(second.getMessageMetadata()).isSameAs(first.getMessageMetadata()); + rangeEntryCache.clear(); + } finally { + cached.release(); + first.release(); + } + // The second read still retains the cache-owned buffer after eviction and the first read's release. + assertThat(second.getMessageMetadata().getProducerName()).isEqualTo("producer"); + assertThat(second.getMessageMetadata().getSequenceId()).isEqualTo(7); + } finally { + second.release(); + } + } + + @Test(timeOut = 30_000) + public void testConcurrentCacheReadsShareMetadata() throws Exception { + managedLedgerConfig.setPulsarMessageEntries(true); + ByteBuf bytes = serializeMessage("producer"); + EntryImpl source = EntryImpl.create(1, 50, bytes); + bytes.release(); + assertThat(rangeEntryCache.insert(source)).isTrue(); + source.release(); + int readers = 8; + CountDownLatch ready = new CountDownLatch(readers); + CountDownLatch start = new CountDownLatch(1); + List> reads = new ArrayList<>(); + ExecutorService executor = Executors.newFixedThreadPool(readers); + try { + try { + for (int i = 0; i < readers; i++) { + reads.add(CompletableFuture.supplyAsync(() -> { + ready.countDown(); + try { + start.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } + return readSingleEntryFromCache(1, 50); + }, executor)); + } + assertThat(ready.await(5, TimeUnit.SECONDS)).isTrue(); + } finally { + start.countDown(); + } + CompletableFuture.allOf(reads.toArray(CompletableFuture[]::new)).get(5, TimeUnit.SECONDS); + MessageMetadata metadata = reads.get(0).join().getMessageMetadata(); + assertThat(metadata).isNotNull(); + for (CompletableFuture read : reads) { + assertThat(read.join().getMessageMetadata()).isSameAs(metadata); + assertThat(read.join().getMessageMetadata().getSequenceId()).isEqualTo(7); + } + rangeEntryCache.clear(); + assertThat(metadata.getProducerName()).isEqualTo("producer"); + } finally { + executor.shutdownNow(); + executor.awaitTermination(5, TimeUnit.SECONDS); + for (CompletableFuture read : reads) { + if (read.isDone() && !read.isCompletedExceptionally()) { + read.join().release(); + } + } + rangeEntryCache.clear(); + } } @Test @@ -268,8 +336,10 @@ public void testCachedEntryMetadataStaysReadableWhenEntriesAreCopied() { ReferenceCountedEntry cached = copyingCache.getEntries().get(PositionFactory.create(1, 50)); assertThat(cached).isNotNull(); - // MessageMetadata decodes its string and bytes fields lazily from the buffer it was parsed from, so the - // cached entry must not share metadata that was parsed from the now released source buffer + assertThat(cached.getMessageMetadata()).isNull(); + Entry readBack = readSingleEntryFromCache(copyingCache, 1, 50); + readBack.release(); + // Metadata is initialized from the retained cache copy, after the source buffer has been released. assertThat(cached.getMessageMetadata()).isNotNull(); assertThat(cached.getMessageMetadata().getProducerName()).isEqualTo("producer"); assertThat(cached.getMessageMetadata().getSequenceId()).isEqualTo(7); @@ -310,7 +380,7 @@ public void testReadFromStorageDoesNotShareSourceMetadataWithTheCopiedCacheEntry assertThat(cached).isNotNull(); // the cached entry is backed by a copy of the payload, so it must not share the metadata that was parsed // from the source buffer - assertThat(cached.getMessageMetadata()).isNotNull().isNotSameAs(sourceMetadata); + assertThat(cached.getMessageMetadata()).isNull(); // overwrite the source payload while it is still referenced, the way the pooled buffer behind it gets // overwritten once it has been recycled. This turns a leftover dependency on the source buffer into a @@ -328,7 +398,10 @@ public void testReadFromStorageDoesNotShareSourceMetadataWithTheCopiedCacheEntry assertThat(headersAndPayload.refCnt()).isZero(); // MessageMetadata decodes its string and bytes fields lazily from the buffer it was parsed from, so the - // cached entry stays readable only because its metadata was parsed from the buffer the cache owns + // cached entry stays readable only because its metadata is parsed from the buffer the cache owns + Entry readBack = readSingleEntryFromCache(copyingCache, 1, 0); + assertThat(readBack.getMessageMetadata()).isNotNull().isNotSameAs(sourceMetadata); + readBack.release(); assertThat(cached.getMessageMetadata().getProducerName()).isEqualTo("producer"); assertThat(cached.getMessageMetadata().getSequenceId()).isEqualTo(7); cached.release(); @@ -352,9 +425,7 @@ public void testInsertDoesNotParseMessageMetadataWhenTheEntriesArentPulsarMessag assertThat(cached.getMessageMetadata()).isNull(); cached.release(); - // reading the entry back through the cache must not parse it either. This is what pins the removal of - // the lazy initialization that RangeCacheEntryWrapper used to do under its write lock, which would have - // defeated skipping the parse at insert time + // Reading back through the cache must also skip metadata initialization for raw ledger entries. Entry readBack = readSingleEntryFromCache(1, 50); assertThat(readBack.getMessageMetadata()).isNull(); readBack.release(); @@ -373,8 +444,11 @@ public void testInsertDoesNotParseMessageMetadataWhenTheEntriesArentPulsarMessag ReferenceCountedEntry cachedControl = rangeEntryCache.getEntries().get(PositionFactory.create(1, 51)); assertThat(cachedControl).isNotNull(); - assertThat(cachedControl.getMessageMetadata()).isNotNull(); - assertThat(cachedControl.getMessageMetadata().getProducerName()).isEqualTo("producer"); + assertThat(cachedControl.getMessageMetadata()).isNull(); + Entry controlRead = readSingleEntryFromCache(1, 51); + assertThat(controlRead.getMessageMetadata()).isNotNull().isSameAs(cachedControl.getMessageMetadata()); + assertThat(controlRead.getMessageMetadata().getProducerName()).isEqualTo("producer"); + controlRead.release(); cachedControl.release(); } @@ -404,8 +478,12 @@ public void testReadFromStorageDoesNotParseMessageMetadataWhenTheEntriesArentPul * @apiNote the returned entry must be released by the caller */ private Entry readSingleEntryFromCache(long ledgerId, long entryId) { + return readSingleEntryFromCache(rangeEntryCache, ledgerId, entryId); + } + + private Entry readSingleEntryFromCache(RangeEntryCacheImpl cache, long ledgerId, long entryId) { CompletableFuture future = new CompletableFuture<>(); - rangeEntryCache.asyncReadEntry(lh, PositionFactory.create(ledgerId, entryId), + cache.asyncReadEntry(lh, PositionFactory.create(ledgerId, entryId), new AsyncCallbacks.ReadEntryCallback() { @Override public void readEntryComplete(Entry entry, Object ctx) { diff --git a/microbench/src/main/java/org/apache/bookkeeper/mledger/impl/cache/EntryCacheMetadataBenchmark.java b/microbench/src/main/java/org/apache/bookkeeper/mledger/impl/cache/EntryCacheMetadataBenchmark.java new file mode 100644 index 0000000000000..fd69663886c13 --- /dev/null +++ b/microbench/src/main/java/org/apache/bookkeeper/mledger/impl/cache/EntryCacheMetadataBenchmark.java @@ -0,0 +1,119 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.impl.cache; + +import io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; +import java.util.concurrent.TimeUnit; +import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.Position; +import org.apache.bookkeeper.mledger.PositionFactory; +import org.apache.bookkeeper.mledger.impl.EntryImpl; +import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.apache.pulsar.common.protocol.Commands; +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Fork; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Measurement; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Param; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.TearDown; +import org.openjdk.jmh.annotations.Warmup; +import org.openjdk.jmh.infra.Blackhole; + +/** + * Measures real cache-entry copy, insertion, metadata initialization and fanout read/release costs. + * Eager initialization emulates parsing before insertion; deferred initialization uses CachedEntries. + * This serial lifecycle benchmark measures total work, not the benefit of moving work between threads. + */ +@State(Scope.Thread) +@BenchmarkMode(Mode.AverageTime) +@OutputTimeUnit(TimeUnit.NANOSECONDS) +@Warmup(iterations = 3, time = 1) +@Measurement(iterations = 5, time = 1) +@Fork(2) +public class EntryCacheMetadataBenchmark { + @Param({"true", "false"}) + public boolean eagerMetadata; + + @Param({"0", "1", "20"}) + public int readers; + + private ByteBuf serialized; + private RangeCache cache; + private RangeCacheRemovalQueue removalQueue; + private final Position position = PositionFactory.create(1, 0); + + @Setup + public void setup() { + MessageMetadata metadata = new MessageMetadata().setProducerName("producer-123") + .setSequenceId(7).setPublishTime(123456789L).setPartitionKey("random-key-123"); + ByteBuf payload = Unpooled.buffer(128).writeZero(128); + try { + serialized = Commands.serializeMetadataAndPayload(Commands.ChecksumType.Crc32c, metadata, payload); + } finally { + payload.release(); + } + removalQueue = new RangeCacheRemovalQueue(0, false); + cache = new RangeCache(removalQueue); + } + + @Benchmark + public void insertReadAndEvict(Blackhole blackhole) { + ByteBuf copied = serialized.copy(); + EntryImpl cached = EntryImpl.createWithRetainedDuplicate(position, copied, null, null); + copied.release(); + if (eagerMetadata) { + cached.initializeMessageMetadataIfNeeded("benchmark"); + } + if (!cache.put(position, cached)) { + cached.release(); + throw new IllegalStateException("Failed to insert entry"); + } + for (int i = 0; i < readers; i++) { + RangeEntryCacheImpl.CachedEntries result = new RangeEntryCacheImpl.CachedEntries(0, 1, "benchmark"); + cache.forEachInRange(position, position, result); + Entry entry = result.entries.get(0); + blackhole.consume(entry.getMessageMetadata().getSequenceId()); + blackhole.consume(entry.getMessageMetadata().getPartitionKey()); + entry.release(); + } + // Clearing only the map leaves wrappers queued for eviction and makes the fixture grow + // throughout the trial. Exercise the real removal queue so each operation is a full cycle. + removalQueue.evictLeastAccessedEntries(Long.MAX_VALUE); + } + + @TearDown(Level.Iteration) + public void verifyEmpty() { + if (cache.getSize() != 0 || !removalQueue.isEmpty()) { + throw new IllegalStateException("Cache lifecycle did not finish eviction"); + } + } + + @TearDown + public void tearDown() { + cache.clear(); + serialized.release(); + } +} diff --git a/microbench/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeCacheReadBenchmark.java b/microbench/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeCacheReadBenchmark.java index 61aaf74dce84e..cdc122d0ade14 100644 --- a/microbench/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeCacheReadBenchmark.java +++ b/microbench/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeCacheReadBenchmark.java @@ -109,7 +109,7 @@ public int readAndCopyRange() { if (range == RANGE_COUNT) { range = 0; } - RangeEntryCacheImpl.CachedEntries result = new RangeEntryCacheImpl.CachedEntries(0, batchSize); + RangeEntryCacheImpl.CachedEntries result = new RangeEntryCacheImpl.CachedEntries(0, batchSize, null); if (visit) { cache.forEachInRange(firstPositions[currentRange], lastPositions[currentRange], result); } else {