From a04b741176759afe920f557fad0e00dbf3449cdc Mon Sep 17 00:00:00 2001 From: Yong Zhang Date: Fri, 25 Mar 2022 14:17:04 +0800 Subject: [PATCH 1/5] Add memory limiter for the add entry request to avoid OOM --- *Motivation* When there has a bookie is slowly to respond the add entry request, and the AQ is smaller than WQ, the client will hole the entry buffer and that will cause the memory won't be released. More context: https://github.com/apache/pulsar/issues/14861 *Modifications* Add the memory limit for the client to send add request. --- .../bookkeeper/conf/ClientConfiguration.java | 22 ++++++ .../bookkeeper/proto/BookieClientImpl.java | 75 +++++++++++++++++-- .../proto/BookkeeperInternalCallbacks.java | 4 + .../proto/PerChannelBookieClient.java | 24 ++++-- 4 files changed, 114 insertions(+), 11 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java index fb9b3d76d77..c3c8c170899 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java @@ -199,6 +199,10 @@ public class ClientConfiguration extends AbstractConfiguration memoryLimitController; public BookieClientImpl(ClientConfiguration conf, EventLoopGroup eventLoopGroup, ByteBufAllocator allocator, @@ -134,6 +138,16 @@ public BookieClientImpl(ClientConfiguration conf, EventLoopGroup eventLoopGroup, } else { this.timeoutFuture = null; } + + if (conf.getClientMemoryLimitEnabled()) { + memoryLimitController = Optional.of(new MemoryLimitController(conf.getClientMemoryLimitByBytes())); + } else { + memoryLimitController = Optional.empty(); + } + } + + public Optional getMemoryLimitController() { + return memoryLimitController; } private int getRc(int rc) { @@ -322,11 +336,33 @@ public void addEntry(final BookieId addr, // Retain the buffer, since the connection could be obtained after // the PendingApp might have already failed toSend.retain(); - + Optional callback = Optional.empty(); + try { + callback = setMemoryLimit(entryId, toSend.readableBytes()); + } catch (InterruptedException e) { + completeAdd(getRc(BKException.Code.IllegalOpException), ledgerId, entryId, addr, cb, ctx); + LOG.error("Failed to set memory limit when adding entry {}:{}", ledgerId, entryId, e); + return; + } client.obtain(ChannelReadyForAddEntryCallback.create( - this, toSend, ledgerId, entryId, addr, - ctx, cb, options, masterKey, allowFastFail, writeFlags), - ledgerId); + this, toSend, ledgerId, entryId, addr, + ctx, cb, options, masterKey, allowFastFail, writeFlags, callback), + ledgerId); + } + + private Optional setMemoryLimit(final long entryId, final long entrySize) throws InterruptedException { + if (getMemoryLimitController().isPresent()) { + MemoryLimitController mlc = getMemoryLimitController().get(); + mlc.reserveMemory(entrySize); + LOG.debug("Acquire memory size {} for entry {}, current usage {} ", entrySize, entryId, + mlc.currentUsage()); + WriteAndFlushCallbackImpl callback = new WriteAndFlushCallbackImpl(); + callback.setBookieClient(this); + callback.setSize(entrySize); + callback.setEntryId(entryId); + return Optional.of(callback); + } + return Optional.empty(); } @Override @@ -375,6 +411,31 @@ public void safeRun() { } } + private static class WriteAndFlushCallbackImpl implements WriteAndFlushCallback { + + private BookieClientImpl bookieClient; + private long size; + private long entryId; + + public void setBookieClient(BookieClientImpl bookieClient) { + this.bookieClient = bookieClient; + } + + public void setSize(long size) { + this.size = size; + } + + public void setEntryId(long entryId) { + this.entryId = entryId; + } + + @Override + public void complete() { + bookieClient.getMemoryLimitController().get().releaseMemory(size); + LOG.debug("Release memory size {} for entry {}", size, entryId); + } + } + private static class ChannelReadyForAddEntryCallback implements GenericCallback { private final Handle recyclerHandle; @@ -390,12 +451,13 @@ private static class ChannelReadyForAddEntryCallback private byte[] masterKey; private boolean allowFastFail; private EnumSet writeFlags; + private Optional writeAndFlushCallback; static ChannelReadyForAddEntryCallback create( BookieClientImpl bookieClient, ByteBufList toSend, long ledgerId, long entryId, BookieId addr, Object ctx, WriteCallback cb, int options, byte[] masterKey, boolean allowFastFail, - EnumSet writeFlags) { + EnumSet writeFlags, Optional writeAndFlushCallback) { ChannelReadyForAddEntryCallback callback = RECYCLER.get(); callback.bookieClient = bookieClient; callback.toSend = toSend; @@ -408,6 +470,7 @@ static ChannelReadyForAddEntryCallback create( callback.masterKey = masterKey; callback.allowFastFail = allowFastFail; callback.writeFlags = writeFlags; + callback.writeAndFlushCallback = writeAndFlushCallback; return callback; } @@ -418,7 +481,7 @@ public void operationComplete(final int rc, bookieClient.completeAdd(rc, ledgerId, entryId, addr, cb, ctx); } else { pcbc.addEntry(ledgerId, masterKey, entryId, - toSend, cb, ctx, options, allowFastFail, writeFlags); + toSend, cb, ctx, options, allowFastFail, writeFlags, writeAndFlushCallback); } toSend.release(); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookkeeperInternalCallbacks.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookkeeperInternalCallbacks.java index f42f7ff13a5..4a68225db63 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookkeeperInternalCallbacks.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookkeeperInternalCallbacks.java @@ -78,6 +78,10 @@ public interface WriteCallback { void writeComplete(int rc, long ledgerId, long entryId, BookieId addr, Object ctx); } + public interface WriteAndFlushCallback { + void complete(); + } + /** * A last-add-confirmed (LAC) reader callback interface. */ diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PerChannelBookieClient.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PerChannelBookieClient.java index ad2777ef4df..039076a317f 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PerChannelBookieClient.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PerChannelBookieClient.java @@ -108,6 +108,7 @@ import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.StartTLSCallback; import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteCallback; import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteLacCallback; +import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteAndFlushCallback; import org.apache.bookkeeper.proto.BookkeeperProtocol.AddRequest; import org.apache.bookkeeper.proto.BookkeeperProtocol.AddResponse; import org.apache.bookkeeper.proto.BookkeeperProtocol.BKPacketHeader; @@ -764,7 +765,8 @@ void forceLedger(final long ledgerId, ForceLedgerCallback cb, Object ctx) { * WriteFlags */ void addEntry(final long ledgerId, byte[] masterKey, final long entryId, ByteBufList toSend, WriteCallback cb, - Object ctx, final int options, boolean allowFastFail, final EnumSet writeFlags) { + Object ctx, final int options, boolean allowFastFail, final EnumSet writeFlags, + Optional wfc) { Object request = null; CompletionKey completionKey = null; if (useV2WireProtocol) { @@ -773,6 +775,7 @@ void addEntry(final long ledgerId, byte[] masterKey, final long entryId, ByteBuf executor.executeOrdered(ledgerId, () -> { cb.writeComplete(BKException.Code.IllegalOpException, ledgerId, entryId, bookieId, ctx); }); + wfc.ifPresent(WriteAndFlushCallback::complete); return; } completionKey = acquireV2Key(ledgerId, entryId, OperationType.ADD_ENTRY); @@ -832,10 +835,11 @@ void addEntry(final long ledgerId, byte[] masterKey, final long entryId, ByteBuf // because we need to release toSend. errorOut(completionKey); toSend.release(); + wfc.ifPresent(WriteAndFlushCallback::complete); return; } else { // addEntry times out on backpressure - writeAndFlush(c, completionKey, request, allowFastFail); + writeAndFlush(c, completionKey, request, allowFastFail, wfc); } } @@ -1117,9 +1121,18 @@ private void writeAndFlush(final Channel channel, } private void writeAndFlush(final Channel channel, - final CompletionKey key, - final Object request, - final boolean allowFastFail) { + final CompletionKey key, + final Object request, + final boolean allowFastFail) { + writeAndFlush(channel, key, request, allowFastFail, Optional.empty()); + + } + + private void writeAndFlush(final Channel channel, + final CompletionKey key, + final Object request, + final boolean allowFastFail, + final Optional wfc) { if (channel == null) { LOG.warn("Operation {} failed: channel == null", StringUtils.requestToString(request)); errorOut(key); @@ -1161,6 +1174,7 @@ private void writeAndFlush(final Channel channel, } else { nettyOpLogger.registerFailedEvent(MathUtils.elapsedNanos(startTime), TimeUnit.NANOSECONDS); } + wfc.ifPresent(WriteAndFlushCallback::complete); }); channel.writeAndFlush(request, promise); From a55b91b7bb5cb6a7213935dec3e43157b8e66456 Mon Sep 17 00:00:00 2001 From: zymap Date: Tue, 11 Oct 2022 16:42:22 +0800 Subject: [PATCH 2/5] Using water mark to control the add behaviors --- .../common/util/WritableListener.java | 25 ++++ .../common/util/WriteMemoryCounter.java | 67 ++++++++++ .../common/util/WriteWaterMark.java | 45 +++++++ .../apache/bookkeeper/client/BookKeeper.java | 18 +++ .../bookkeeper/client/ClientContext.java | 3 + .../bookkeeper/client/PendingAddOp.java | 3 +- .../bookkeeper/client/api/WriteHandle.java | 9 ++ .../bookkeeper/conf/ClientConfiguration.java | 20 +++ .../client/BookieClientMemoryCounterTest.java | 125 ++++++++++++++++++ .../client/MockBookKeeperTestCase.java | 9 +- .../bookkeeper/client/MockClientContext.java | 20 ++- 11 files changed, 339 insertions(+), 5 deletions(-) create mode 100644 bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WritableListener.java create mode 100644 bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java create mode 100644 bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteWaterMark.java create mode 100644 bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookieClientMemoryCounterTest.java diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WritableListener.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WritableListener.java new file mode 100644 index 00000000000..08e02a30deb --- /dev/null +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WritableListener.java @@ -0,0 +1,25 @@ +/** + * 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.common.util; + +public interface WritableListener { + + void onWriteStateChanged(boolean writable); + +} diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java new file mode 100644 index 00000000000..22b39e93800 --- /dev/null +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java @@ -0,0 +1,67 @@ +/** + * 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.common.util; + +import lombok.extern.slf4j.Slf4j; + +import java.util.LinkedList; +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; + +@Slf4j +public class WriteMemoryCounter { + private final WriteWaterMark writeWaterMark; + private AtomicLong sizeCounter = new AtomicLong(0); + private AtomicBoolean writeState = new AtomicBoolean(true); + private final List listeners = new LinkedList<>(); + + public WriteMemoryCounter(WriteWaterMark writeWaterMark) { + this.writeWaterMark = writeWaterMark; + } + + public WriteMemoryCounter() { + this.writeWaterMark = new WriteWaterMark(); + } + + public void register(WritableListener listener) { + listeners.add(listener); + } + + public void incrementPendingWriteBytes(long size) { + long usage = sizeCounter.addAndGet(size); + log.info("increment the size to {}", usage); + if (usage > writeWaterMark.high() && writeState.get()) { + setWritable(false); + } + } + + public void decrementPendingWriteBytes(long size) { + long usage = sizeCounter.addAndGet(-size); + log.info("decrement the size to {}", usage); + if (usage < writeWaterMark.low() && !writeState.get()) { + setWritable(true); + } + } + + public void setWritable(boolean state) { + writeState.set(state); + listeners.forEach(l -> l.onWriteStateChanged(state)); + } +} diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteWaterMark.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteWaterMark.java new file mode 100644 index 00000000000..4005d79ae44 --- /dev/null +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteWaterMark.java @@ -0,0 +1,45 @@ +/** + * 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.common.util; + +public class WriteWaterMark { + private static final int DEFAULT_LOW_WATER_MARK = 1; + private static final int DEFAULT_HIGH_WATER_MARK = 1; + + private final int low; + private final int high; + + public WriteWaterMark(int low, int high) { + this.low = low; + this.high = high; + } + + public WriteWaterMark() { + this.low = DEFAULT_LOW_WATER_MARK; + this.high = DEFAULT_HIGH_WATER_MARK; + } + + public int low() { + return low; + } + + public int high() { + return high; + } +} diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeper.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeper.java index a48e7d62a3d..ef1f243a1f7 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeper.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BookKeeper.java @@ -67,6 +67,8 @@ import org.apache.bookkeeper.common.util.OrderedExecutor; import org.apache.bookkeeper.common.util.OrderedScheduler; import org.apache.bookkeeper.common.util.ReflectionUtils; +import org.apache.bookkeeper.common.util.WriteMemoryCounter; +import org.apache.bookkeeper.common.util.WriteWaterMark; import org.apache.bookkeeper.conf.AbstractConfiguration; import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.feature.FeatureProvider; @@ -151,6 +153,8 @@ public class BookKeeper implements org.apache.bookkeeper.client.api.BookKeeper { boolean closed = false; final ReentrantReadWriteLock closeLock = new ReentrantReadWriteLock(); + final WriteMemoryCounter writeMemoryCounter; + /** * BookKeeper Client Builder to build client instances. * @@ -538,6 +542,10 @@ public BookKeeper(ClientConfiguration conf, ZooKeeper zk, EventLoopGroup eventLo this.ledgerIdGenerator = ledgerManagerFactory.newLedgerIdGenerator(); this.bookieQuarantineRatio = conf.getBookieQuarantineRatio(); + + this.writeMemoryCounter = new WriteMemoryCounter( + new WriteWaterMark(conf.getWriteMemoryLowWaterMark(), conf.getWriteMemoryHighWaterMark())); + scheduleBookieHealthCheckIfEnabled(conf); } @@ -566,6 +574,7 @@ public BookKeeper(ClientConfiguration conf, ZooKeeper zk, EventLoopGroup eventLo bookieClient = null; allocator = UnpooledByteBufAllocator.DEFAULT; bookieQuarantineRatio = 1.0; + writeMemoryCounter = null; } protected EnsemblePlacementPolicy initializeEnsemblePlacementPolicy(ClientConfiguration conf, @@ -705,6 +714,10 @@ public MetadataClientDriver getMetadataClientDriver() { return metadataDriver; } + WriteMemoryCounter getWriteMemoryCounter() { + return writeMemoryCounter; + } + /** * There are 3 digest types that can be used for verification. The CRC32 is * cheap to compute but does not protect against byzantine bookies (i.e., a @@ -1662,6 +1675,11 @@ public boolean isClientClosed() { public ByteBufAllocator getByteBufAllocator() { return allocator; } + + @Override + public WriteMemoryCounter getWriteMemoryCounter() { + return writeMemoryCounter; + } }; public ClientContext getClientCtx() { diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/ClientContext.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/ClientContext.java index 3b43502d96b..a76f3a2bf57 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/ClientContext.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/ClientContext.java @@ -23,6 +23,7 @@ import io.netty.buffer.ByteBufAllocator; import org.apache.bookkeeper.common.util.OrderedExecutor; import org.apache.bookkeeper.common.util.OrderedScheduler; +import org.apache.bookkeeper.common.util.WriteMemoryCounter; import org.apache.bookkeeper.meta.LedgerManager; import org.apache.bookkeeper.proto.BookieClient; @@ -43,4 +44,6 @@ public interface ClientContext { OrderedScheduler getScheduler(); BookKeeperClientStats getClientStats(); boolean isClientClosed(); + + WriteMemoryCounter getWriteMemoryCounter(); } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/PendingAddOp.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/PendingAddOp.java index d04f0d146c4..6a20b58f49e 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/PendingAddOp.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/PendingAddOp.java @@ -96,7 +96,7 @@ static PendingAddOp create(LedgerHandle lh, ClientContext clientCtx, op.currentLedgerLength = -1; op.payload = payload; op.entryLength = payload.readableBytes(); - + op.clientCtx.getWriteMemoryCounter().incrementPendingWriteBytes(op.entryLength); op.completed = false; op.ensemble = ensemble; op.ackSet = lh.getDistributionSchedule().getAckSet(); @@ -493,6 +493,7 @@ private void maybeRecycle() { } // only recycle a pending add op after it has been run. if (hasRun && toSend == null && pendingWriteRequests == 0) { + clientCtx.getWriteMemoryCounter().decrementPendingWriteBytes(entryLength); recyclePendAddOpObject(); } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/api/WriteHandle.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/api/WriteHandle.java index 28aded9cc68..c3bcf6c1841 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/api/WriteHandle.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/api/WriteHandle.java @@ -141,6 +141,15 @@ default long append(byte[] data, int offset, int length) throws BKException, Int */ long getLastAddPushed(); + /** + * + * + * @return + */ + default boolean isWritable() { + throw new UnsupportedOperationException("This operation is not supported for the current handler"); + } + /** * Asynchronous close the write handle, any adds in flight will return errors. * diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java index c3c8c170899..8d2670af7b5 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java @@ -202,6 +202,8 @@ public class ClientConfiguration extends AbstractConfiguration + * 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.client; + +import io.netty.buffer.UnpooledByteBufAllocator; +import lombok.extern.slf4j.Slf4j; +import org.apache.bookkeeper.bookie.Bookie; +import org.apache.bookkeeper.bookie.BookieImpl; +import org.apache.bookkeeper.bookie.Journal; +import org.apache.bookkeeper.bookie.SlowBufferedChannel; +import org.apache.bookkeeper.bookie.TestBookieImpl; +import org.apache.bookkeeper.common.util.WritableListener; +import org.apache.bookkeeper.conf.ServerConfiguration; +import org.apache.bookkeeper.proto.BookieServer; +import org.apache.bookkeeper.test.BookKeeperClusterTestCase; +import org.junit.Test; + +import java.lang.reflect.Field; +import java.nio.channels.FileChannel; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; + +@Slf4j +public class BookieClientMemoryCounterTest extends BookKeeperClusterTestCase { + + public BookieClientMemoryCounterTest() { + super(1); + baseClientConf.setAddEntryTimeout(10000); + baseClientConf.setAddEntryQuorumTimeout(10000); + baseClientConf.setWriteMemoryHighWaterMark(8 * 1024); + baseClientConf.setWriteMemoryLowWaterMark(2 * 1024); + } + + @Test + public void testPendingAddEntryMemory() throws Exception { + confByIndex(0).setMaxAddsInProgressLimit(30); + ServerConfiguration conf = killBookie(0); + BookieServer bks = startAndAddBookie(conf, + bookieWithMockedJournal(conf, 0, 1, 0)) + .getServer(); + + + AtomicBoolean writeState = new AtomicBoolean(true); + bkc.getWriteMemoryCounter().register(new WritableListener() { + @Override + public void onWriteStateChanged(boolean writable) { + log.info("Write state changed to {}", writeState); + writeState.set(writable); + } + }); + LedgerHandle lh = bkc.createLedger(1,1, BookKeeper.DigestType.CRC32, "".getBytes()); + CountDownLatch complete = new CountDownLatch(1); + byte[] msg = new byte[1024]; + + CountDownLatch latch = new CountDownLatch(100); + for (int i = 0; i < 100; i++) { + while (!writeState.get()) { + log.info("wait for the memory released"); + TimeUnit.SECONDS.sleep(1); + } + lh.asyncAddEntry(msg, new AsyncCallback.AddCallback() { + @Override + public void addComplete(int rc, LedgerHandle lh, long entryId, Object ctx) { + log.info("Add complete with rc {}", rc); + latch.countDown(); + } + }, null); + } + latch.await(); + lh.close(); + } + + private Bookie bookieWithMockedJournal(ServerConfiguration conf, + long getDelay, long addDelay, long flushDelay) throws Exception { + Bookie bookie = new TestBookieImpl(conf); + if (getDelay <= 0 && addDelay <= 0 && flushDelay <= 0) { + return bookie; + } + + List journals = getJournals(bookie); + for (int i = 0; i < journals.size(); i++) { + Journal mock = spy(journals.get(i)); + when(mock.getBufferedChannelBuilder()).thenReturn((FileChannel fc, int capacity) -> { + SlowBufferedChannel sbc = new SlowBufferedChannel(UnpooledByteBufAllocator.DEFAULT, fc, capacity); + sbc.setAddDelay(addDelay); + sbc.setGetDelay(getDelay); + sbc.setFlushDelay(flushDelay); + return sbc; + }); + + journals.set(i, mock); + } + return bookie; + } + + @SuppressWarnings("unchecked") + private List getJournals(Bookie bookie) throws NoSuchFieldException, IllegalAccessException { + Field f = BookieImpl.class.getDeclaredField("journals"); + f.setAccessible(true); + + return (List) f.get(bookie); + } + +} diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java index 2ee6feb3c8e..28a32e125e1 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java @@ -58,6 +58,7 @@ import org.apache.bookkeeper.client.api.OpenBuilder; import org.apache.bookkeeper.common.util.OrderedExecutor; import org.apache.bookkeeper.common.util.OrderedScheduler; +import org.apache.bookkeeper.common.util.WriteMemoryCounter; import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.meta.LedgerIdGenerator; import org.apache.bookkeeper.meta.LedgerManager; @@ -170,6 +171,7 @@ public void setup() throws Exception { when(bk.getMainWorkerPool()).thenReturn(executor); when(bk.getBookieClient()).thenReturn(bookieClient); when(bk.getScheduler()).thenReturn(scheduler); + when(bk.getWriteMemoryCounter()).thenCallRealMethod(); setBookKeeperConfig(new ClientConfiguration()); when(bk.getStatsLogger()).thenReturn(NullStatsLogger.INSTANCE); @@ -224,7 +226,12 @@ public boolean isClientClosed() { public ByteBufAllocator getByteBufAllocator() { return UnpooledByteBufAllocator.DEFAULT; } - }; + + @Override + public WriteMemoryCounter getWriteMemoryCounter() { + return bk.getWriteMemoryCounter(); + } + }; when(bk.getClientCtx()).thenReturn(clientCtx); when(bk.getLedgerManager()).thenReturn(ledgerManager); when(bk.getLedgerIdGenerator()).thenReturn(ledgerIdGenerator); diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java index 93078a05129..5078eb79094 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java @@ -27,6 +27,7 @@ import java.util.function.BooleanSupplier; import org.apache.bookkeeper.common.util.OrderedExecutor; import org.apache.bookkeeper.common.util.OrderedScheduler; +import org.apache.bookkeeper.common.util.WriteMemoryCounter; import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.discover.MockRegistrationClient; import org.apache.bookkeeper.meta.LedgerManager; @@ -53,6 +54,7 @@ public class MockClientContext implements ClientContext { private BooleanSupplier isClientClosed; private MockRegistrationClient regClient; private ByteBufAllocator allocator; + private WriteMemoryCounter writeMemoryCounter; static MockClientContext create(MockBookies mockBookies) throws Exception { ClientConfiguration conf = new ClientConfiguration(); @@ -64,7 +66,7 @@ static MockClientContext create(MockBookies mockBookies) throws Exception { new DefaultBookieAddressResolver(regClient), NullStatsLogger.INSTANCE); bookieWatcherImpl.initialBlockingBookieRead(); - + WriteMemoryCounter memoryCounter = new WriteMemoryCounter(); return new MockClientContext() .setConf(ClientInternalConf.fromConfig(conf)) .setLedgerManager(new MockLedgerManager()) @@ -76,7 +78,8 @@ static MockClientContext create(MockBookies mockBookies) throws Exception { .setMainWorkerPool(scheduler) .setScheduler(scheduler) .setClientStats(BookKeeperClientStats.newInstance(NullStatsLogger.INSTANCE)) - .setIsClientClosed(() -> false); + .setIsClientClosed(() -> false) + .setWriteMemoryCounter(memoryCounter); } static MockClientContext create() throws Exception { @@ -95,7 +98,8 @@ static MockClientContext copyOf(ClientContext other) { .setScheduler(other.getScheduler()) .setClientStats(other.getClientStats()) .setByteBufAllocator(other.getByteBufAllocator()) - .setIsClientClosed(other::isClientClosed); + .setIsClientClosed(other::isClientClosed) + .setWriteMemoryCounter(other.getWriteMemoryCounter()); } public MockRegistrationClient getMockRegistrationClient() { @@ -168,6 +172,11 @@ public MockClientContext setByteBufAllocator(ByteBufAllocator allocator) { return this; } + public MockClientContext setWriteMemoryCounter(WriteMemoryCounter writeMemoryCounter) { + this.writeMemoryCounter = writeMemoryCounter; + return this; + } + private static T maybeSpy(T orig) { if (Mockito.mockingDetails(orig).isSpy()) { return orig; @@ -225,4 +234,9 @@ public boolean isClientClosed() { public ByteBufAllocator getByteBufAllocator() { return allocator; } + + @Override + public WriteMemoryCounter getWriteMemoryCounter() { + return writeMemoryCounter; + } } From 8d747ee80d348d4770a0895bcdcce4bdb1ba1ed8 Mon Sep 17 00:00:00 2001 From: zymap Date: Tue, 11 Oct 2022 17:15:32 +0800 Subject: [PATCH 3/5] Revert "Add memory limiter for the add entry request to avoid OOM" This reverts commit a04b741176759afe920f557fad0e00dbf3449cdc. --- .../bookkeeper/conf/ClientConfiguration.java | 21 ------ .../bookkeeper/proto/BookieClientImpl.java | 75 ++----------------- .../proto/BookkeeperInternalCallbacks.java | 4 - .../proto/PerChannelBookieClient.java | 24 ++---- .../client/BookieClientMemoryCounterTest.java | 1 - 5 files changed, 11 insertions(+), 114 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java index 8d2670af7b5..eb276c10d02 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java @@ -199,9 +199,6 @@ public class ClientConfiguration extends AbstractConfiguration memoryLimitController; public BookieClientImpl(ClientConfiguration conf, EventLoopGroup eventLoopGroup, ByteBufAllocator allocator, @@ -138,16 +134,6 @@ public BookieClientImpl(ClientConfiguration conf, EventLoopGroup eventLoopGroup, } else { this.timeoutFuture = null; } - - if (conf.getClientMemoryLimitEnabled()) { - memoryLimitController = Optional.of(new MemoryLimitController(conf.getClientMemoryLimitByBytes())); - } else { - memoryLimitController = Optional.empty(); - } - } - - public Optional getMemoryLimitController() { - return memoryLimitController; } private int getRc(int rc) { @@ -336,33 +322,11 @@ public void addEntry(final BookieId addr, // Retain the buffer, since the connection could be obtained after // the PendingApp might have already failed toSend.retain(); - Optional callback = Optional.empty(); - try { - callback = setMemoryLimit(entryId, toSend.readableBytes()); - } catch (InterruptedException e) { - completeAdd(getRc(BKException.Code.IllegalOpException), ledgerId, entryId, addr, cb, ctx); - LOG.error("Failed to set memory limit when adding entry {}:{}", ledgerId, entryId, e); - return; - } - client.obtain(ChannelReadyForAddEntryCallback.create( - this, toSend, ledgerId, entryId, addr, - ctx, cb, options, masterKey, allowFastFail, writeFlags, callback), - ledgerId); - } - private Optional setMemoryLimit(final long entryId, final long entrySize) throws InterruptedException { - if (getMemoryLimitController().isPresent()) { - MemoryLimitController mlc = getMemoryLimitController().get(); - mlc.reserveMemory(entrySize); - LOG.debug("Acquire memory size {} for entry {}, current usage {} ", entrySize, entryId, - mlc.currentUsage()); - WriteAndFlushCallbackImpl callback = new WriteAndFlushCallbackImpl(); - callback.setBookieClient(this); - callback.setSize(entrySize); - callback.setEntryId(entryId); - return Optional.of(callback); - } - return Optional.empty(); + client.obtain(ChannelReadyForAddEntryCallback.create( + this, toSend, ledgerId, entryId, addr, + ctx, cb, options, masterKey, allowFastFail, writeFlags), + ledgerId); } @Override @@ -411,31 +375,6 @@ public void safeRun() { } } - private static class WriteAndFlushCallbackImpl implements WriteAndFlushCallback { - - private BookieClientImpl bookieClient; - private long size; - private long entryId; - - public void setBookieClient(BookieClientImpl bookieClient) { - this.bookieClient = bookieClient; - } - - public void setSize(long size) { - this.size = size; - } - - public void setEntryId(long entryId) { - this.entryId = entryId; - } - - @Override - public void complete() { - bookieClient.getMemoryLimitController().get().releaseMemory(size); - LOG.debug("Release memory size {} for entry {}", size, entryId); - } - } - private static class ChannelReadyForAddEntryCallback implements GenericCallback { private final Handle recyclerHandle; @@ -451,13 +390,12 @@ private static class ChannelReadyForAddEntryCallback private byte[] masterKey; private boolean allowFastFail; private EnumSet writeFlags; - private Optional writeAndFlushCallback; static ChannelReadyForAddEntryCallback create( BookieClientImpl bookieClient, ByteBufList toSend, long ledgerId, long entryId, BookieId addr, Object ctx, WriteCallback cb, int options, byte[] masterKey, boolean allowFastFail, - EnumSet writeFlags, Optional writeAndFlushCallback) { + EnumSet writeFlags) { ChannelReadyForAddEntryCallback callback = RECYCLER.get(); callback.bookieClient = bookieClient; callback.toSend = toSend; @@ -470,7 +408,6 @@ static ChannelReadyForAddEntryCallback create( callback.masterKey = masterKey; callback.allowFastFail = allowFastFail; callback.writeFlags = writeFlags; - callback.writeAndFlushCallback = writeAndFlushCallback; return callback; } @@ -481,7 +418,7 @@ public void operationComplete(final int rc, bookieClient.completeAdd(rc, ledgerId, entryId, addr, cb, ctx); } else { pcbc.addEntry(ledgerId, masterKey, entryId, - toSend, cb, ctx, options, allowFastFail, writeFlags, writeAndFlushCallback); + toSend, cb, ctx, options, allowFastFail, writeFlags); } toSend.release(); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookkeeperInternalCallbacks.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookkeeperInternalCallbacks.java index 4a68225db63..f42f7ff13a5 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookkeeperInternalCallbacks.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookkeeperInternalCallbacks.java @@ -78,10 +78,6 @@ public interface WriteCallback { void writeComplete(int rc, long ledgerId, long entryId, BookieId addr, Object ctx); } - public interface WriteAndFlushCallback { - void complete(); - } - /** * A last-add-confirmed (LAC) reader callback interface. */ diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PerChannelBookieClient.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PerChannelBookieClient.java index 039076a317f..ad2777ef4df 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PerChannelBookieClient.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PerChannelBookieClient.java @@ -108,7 +108,6 @@ import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.StartTLSCallback; import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteCallback; import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteLacCallback; -import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteAndFlushCallback; import org.apache.bookkeeper.proto.BookkeeperProtocol.AddRequest; import org.apache.bookkeeper.proto.BookkeeperProtocol.AddResponse; import org.apache.bookkeeper.proto.BookkeeperProtocol.BKPacketHeader; @@ -765,8 +764,7 @@ void forceLedger(final long ledgerId, ForceLedgerCallback cb, Object ctx) { * WriteFlags */ void addEntry(final long ledgerId, byte[] masterKey, final long entryId, ByteBufList toSend, WriteCallback cb, - Object ctx, final int options, boolean allowFastFail, final EnumSet writeFlags, - Optional wfc) { + Object ctx, final int options, boolean allowFastFail, final EnumSet writeFlags) { Object request = null; CompletionKey completionKey = null; if (useV2WireProtocol) { @@ -775,7 +773,6 @@ void addEntry(final long ledgerId, byte[] masterKey, final long entryId, ByteBuf executor.executeOrdered(ledgerId, () -> { cb.writeComplete(BKException.Code.IllegalOpException, ledgerId, entryId, bookieId, ctx); }); - wfc.ifPresent(WriteAndFlushCallback::complete); return; } completionKey = acquireV2Key(ledgerId, entryId, OperationType.ADD_ENTRY); @@ -835,11 +832,10 @@ void addEntry(final long ledgerId, byte[] masterKey, final long entryId, ByteBuf // because we need to release toSend. errorOut(completionKey); toSend.release(); - wfc.ifPresent(WriteAndFlushCallback::complete); return; } else { // addEntry times out on backpressure - writeAndFlush(c, completionKey, request, allowFastFail, wfc); + writeAndFlush(c, completionKey, request, allowFastFail); } } @@ -1121,18 +1117,9 @@ private void writeAndFlush(final Channel channel, } private void writeAndFlush(final Channel channel, - final CompletionKey key, - final Object request, - final boolean allowFastFail) { - writeAndFlush(channel, key, request, allowFastFail, Optional.empty()); - - } - - private void writeAndFlush(final Channel channel, - final CompletionKey key, - final Object request, - final boolean allowFastFail, - final Optional wfc) { + final CompletionKey key, + final Object request, + final boolean allowFastFail) { if (channel == null) { LOG.warn("Operation {} failed: channel == null", StringUtils.requestToString(request)); errorOut(key); @@ -1174,7 +1161,6 @@ private void writeAndFlush(final Channel channel, } else { nettyOpLogger.registerFailedEvent(MathUtils.elapsedNanos(startTime), TimeUnit.NANOSECONDS); } - wfc.ifPresent(WriteAndFlushCallback::complete); }); channel.writeAndFlush(request, promise); diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookieClientMemoryCounterTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookieClientMemoryCounterTest.java index 02566650d18..8d829c99e71 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookieClientMemoryCounterTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookieClientMemoryCounterTest.java @@ -70,7 +70,6 @@ public void onWriteStateChanged(boolean writable) { } }); LedgerHandle lh = bkc.createLedger(1,1, BookKeeper.DigestType.CRC32, "".getBytes()); - CountDownLatch complete = new CountDownLatch(1); byte[] msg = new byte[1024]; CountDownLatch latch = new CountDownLatch(100); From 19e3d5c36b8e8fdcb2df7edbd5ca5401176f5c2d Mon Sep 17 00:00:00 2001 From: zymap Date: Mon, 17 Oct 2022 15:58:42 +0800 Subject: [PATCH 4/5] Add java docs for new classes --- .../common/util/WritableListener.java | 6 ++++++ .../common/util/WriteMemoryCounter.java | 10 ++++++++++ .../bookkeeper/common/util/WriteWaterMark.java | 17 ++++++++++------- 3 files changed, 26 insertions(+), 7 deletions(-) diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WritableListener.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WritableListener.java index 08e02a30deb..325f8434330 100644 --- a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WritableListener.java +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WritableListener.java @@ -18,6 +18,12 @@ */ package org.apache.bookkeeper.common.util; +/** + * WritableListener used to listen the writable status changes. + * + * We use {@link WriteMemoryCounter} to listen on the memory usage when the client adds entries. The listener can + * take actions if they have been notified. + */ public interface WritableListener { void onWriteStateChanged(boolean writable); diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java index 22b39e93800..0b7b8bb2884 100644 --- a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java @@ -25,6 +25,16 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; +/** + * {@link WriteMemoryCounter} counts the memory usage on Adds request. + * When there has an Add request created, the {@link WriteMemoryCounter} will record the request content size. + * When the request is finished, the {@link WriteMemoryCounter} will decrease the record count. + * The range of the counter should in the range of {@link WriteWaterMark}'s high watermark and low watermark. + * + * If the record size is over to the high watermark, the registered listeners will receive writable state change + * to false and take actions. The listeners will receive writable state change to true until the record size is + * down to the low watermark. + */ @Slf4j public class WriteMemoryCounter { private final WriteWaterMark writeWaterMark; diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteWaterMark.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteWaterMark.java index 4005d79ae44..2199dec4699 100644 --- a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteWaterMark.java +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteWaterMark.java @@ -18,14 +18,17 @@ */ package org.apache.bookkeeper.common.util; +/** + * {@link WriteWaterMark} is used to configure the max value and the min value of the memory usage. + */ public class WriteWaterMark { - private static final int DEFAULT_LOW_WATER_MARK = 1; - private static final int DEFAULT_HIGH_WATER_MARK = 1; + private static final long DEFAULT_LOW_WATER_MARK = 1610612736L; // 1.5GB + private static final long DEFAULT_HIGH_WATER_MARK = 2147483648L; // 2GB - private final int low; - private final int high; + private final long low; + private final long high; - public WriteWaterMark(int low, int high) { + public WriteWaterMark(long low, long high) { this.low = low; this.high = high; } @@ -35,11 +38,11 @@ public WriteWaterMark() { this.high = DEFAULT_HIGH_WATER_MARK; } - public int low() { + public long low() { return low; } - public int high() { + public long high() { return high; } } From 33f8fbbcbeb2df517c790733de786d50ccb629c5 Mon Sep 17 00:00:00 2001 From: zymap Date: Tue, 18 Oct 2022 20:14:05 +0800 Subject: [PATCH 5/5] Update the tests and fix the style --- .../common/util/WriteMemoryCounter.java | 7 +- .../bookkeeper/conf/ClientConfiguration.java | 8 +- .../client/BookieClientMemoryCounterTest.java | 130 ++++++++---------- 3 files changed, 65 insertions(+), 80 deletions(-) diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java index 0b7b8bb2884..aa199ced51b 100644 --- a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/WriteMemoryCounter.java @@ -18,12 +18,11 @@ */ package org.apache.bookkeeper.common.util; -import lombok.extern.slf4j.Slf4j; - import java.util.LinkedList; import java.util.List; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; +import lombok.extern.slf4j.Slf4j; /** * {@link WriteMemoryCounter} counts the memory usage on Adds request. @@ -74,4 +73,8 @@ public void setWritable(boolean state) { writeState.set(state); listeners.forEach(l -> l.onWriteStateChanged(state)); } + + public long getSize() { + return sizeCounter.get(); + } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java index eb276c10d02..f8b47cbdbf7 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ClientConfiguration.java @@ -2061,21 +2061,21 @@ public long getClientConnectBookieUnavailableLogThrottlingMs() { return getLong(CLIENT_CONNECT_BOOKIE_UNAVAILABLE_LOG_THROTTLING, 5_000L); } - public ClientConfiguration setWriteMemoryLowWaterMark(int bytes) { + public ClientConfiguration setWriteMemoryLowWaterMark(long bytes) { setProperty(WRITE_MEMORY_LOW_WATER_MARK, bytes); return this; } - public int getWriteMemoryLowWaterMark() { + public long getWriteMemoryLowWaterMark() { return getInt(WRITE_MEMORY_LOW_WATER_MARK, 64 * 1024 * 1024); } - public ClientConfiguration setWriteMemoryHighWaterMark(int bytes) { + public ClientConfiguration setWriteMemoryHighWaterMark(long bytes) { setProperty(WRITE_MEMORY_HIGH_WATER_MARK, bytes); return this; } - public int getWriteMemoryHighWaterMark() { + public long getWriteMemoryHighWaterMark() { return getInt(WRITE_MEMORY_HIGH_WATER_MARK, 256 * 1024 * 1024); } diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookieClientMemoryCounterTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookieClientMemoryCounterTest.java index 8d829c99e71..ffc300a42ed 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookieClientMemoryCounterTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/BookieClientMemoryCounterTest.java @@ -18,107 +18,89 @@ */ package org.apache.bookkeeper.client; -import io.netty.buffer.UnpooledByteBufAllocator; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import lombok.extern.slf4j.Slf4j; -import org.apache.bookkeeper.bookie.Bookie; -import org.apache.bookkeeper.bookie.BookieImpl; -import org.apache.bookkeeper.bookie.Journal; -import org.apache.bookkeeper.bookie.SlowBufferedChannel; -import org.apache.bookkeeper.bookie.TestBookieImpl; +import org.apache.bookkeeper.client.api.BKException; import org.apache.bookkeeper.common.util.WritableListener; -import org.apache.bookkeeper.conf.ServerConfiguration; -import org.apache.bookkeeper.proto.BookieServer; import org.apache.bookkeeper.test.BookKeeperClusterTestCase; import org.junit.Test; -import java.lang.reflect.Field; -import java.nio.channels.FileChannel; -import java.util.List; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; - -import static org.mockito.Mockito.spy; -import static org.mockito.Mockito.when; @Slf4j public class BookieClientMemoryCounterTest extends BookKeeperClusterTestCase { + static final int MESSAGE_SIZE = 1024; + static final long LOW_WATER_MARK = 10 * 1024; + static final long HIGH_WATER_MARK = 20 * 1024; + public BookieClientMemoryCounterTest() { super(1); - baseClientConf.setAddEntryTimeout(10000); - baseClientConf.setAddEntryQuorumTimeout(10000); - baseClientConf.setWriteMemoryHighWaterMark(8 * 1024); - baseClientConf.setWriteMemoryLowWaterMark(2 * 1024); + baseClientConf.setWriteMemoryHighWaterMark(HIGH_WATER_MARK); + baseClientConf.setWriteMemoryLowWaterMark(LOW_WATER_MARK); } @Test public void testPendingAddEntryMemory() throws Exception { - confByIndex(0).setMaxAddsInProgressLimit(30); - ServerConfiguration conf = killBookie(0); - BookieServer bks = startAndAddBookie(conf, - bookieWithMockedJournal(conf, 0, 1, 0)) - .getServer(); - - + // listen to the write state change events AtomicBoolean writeState = new AtomicBoolean(true); bkc.getWriteMemoryCounter().register(new WritableListener() { @Override public void onWriteStateChanged(boolean writable) { - log.info("Write state changed to {}", writeState); + long usage = bkc.getWriteMemoryCounter().getSize(); + log.info("Write state changed to {}, current memory usage is {}", writeState, usage); + // when the writable change to ture, the usage should under the LowWaterMark. + // when the writable change to false, the usage should over than the HighWaterMark. + if (writable) { + assertEquals(LOW_WATER_MARK - MESSAGE_SIZE, usage); + } else { + assertEquals(HIGH_WATER_MARK + MESSAGE_SIZE, usage); + } writeState.set(writable); } }); - LedgerHandle lh = bkc.createLedger(1,1, BookKeeper.DigestType.CRC32, "".getBytes()); - byte[] msg = new byte[1024]; - CountDownLatch latch = new CountDownLatch(100); - for (int i = 0; i < 100; i++) { - while (!writeState.get()) { - log.info("wait for the memory released"); - TimeUnit.SECONDS.sleep(1); - } - lh.asyncAddEntry(msg, new AsyncCallback.AddCallback() { - @Override - public void addComplete(int rc, LedgerHandle lh, long entryId, Object ctx) { - log.info("Add complete with rc {}", rc); - latch.countDown(); - } - }, null); - } - latch.await(); - lh.close(); - } + LedgerHandle lh = bkc.createLedger(1, 1, BookKeeper.DigestType.CRC32, "".getBytes()); + byte[] msg = new byte[MESSAGE_SIZE]; - private Bookie bookieWithMockedJournal(ServerConfiguration conf, - long getDelay, long addDelay, long flushDelay) throws Exception { - Bookie bookie = new TestBookieImpl(conf); - if (getDelay <= 0 && addDelay <= 0 && flushDelay <= 0) { - return bookie; - } + int testMessagesNum = 1000; - List journals = getJournals(bookie); - for (int i = 0; i < journals.size(); i++) { - Journal mock = spy(journals.get(i)); - when(mock.getBufferedChannelBuilder()).thenReturn((FileChannel fc, int capacity) -> { - SlowBufferedChannel sbc = new SlowBufferedChannel(UnpooledByteBufAllocator.DEFAULT, fc, capacity); - sbc.setAddDelay(addDelay); - sbc.setGetDelay(getDelay); - sbc.setFlushDelay(flushDelay); - return sbc; - }); + // start a thread to send message + AtomicInteger addCount = new AtomicInteger(testMessagesNum); + new Thread(() -> { + for (int i = 0; i < testMessagesNum; i++) { + while (!writeState.get()) { + log.info("wait for the memory released"); + try { + TimeUnit.MILLISECONDS.sleep(10); + } catch (InterruptedException e) { + // ignore + } + } + lh.asyncAddEntry(msg, new AsyncCallback.AddCallback() { + @Override + public void addComplete(int rc, LedgerHandle lh, long entryId, Object ctx) { + if (rc == BKException.Code.OK) { + log.info("Add complete with rc {}", rc); + addCount.getAndDecrement(); + } + } + }, null); + } + }).start(); - journals.set(i, mock); + // while sending messages, we listen on the memory counter size. The size should never over than the + // (highWaterMark + 1 message) bytes. + while (addCount.get() != 0) { + long size = bkc.getWriteMemoryCounter().getSize(); + assertTrue(size >= 0 && size < baseClientConf.getWriteMemoryHighWaterMark() + msg.length + 1); + TimeUnit.MILLISECONDS.sleep(10); } - return bookie; - } - - @SuppressWarnings("unchecked") - private List getJournals(Bookie bookie) throws NoSuchFieldException, IllegalAccessException { - Field f = BookieImpl.class.getDeclaredField("journals"); - f.setAccessible(true); - return (List) f.get(bookie); + lh.close(); } - }