From 79105b6d01517d759085d614a69298a78feb841b Mon Sep 17 00:00:00 2001 From: horizonzy Date: Wed, 8 Mar 2023 16:21:16 +0800 Subject: [PATCH 01/12] Optimize group flush pending response. --- .../org/apache/bookkeeper/bookie/Journal.java | 61 ++++++++++++------- .../proto/BookieRequestHandler.java | 29 +++++++-- 2 files changed, 63 insertions(+), 27 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java index 7966f6d2abf..f034a2fdf29 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java @@ -23,6 +23,7 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Stopwatch; +import com.scurrilous.circe.Hash; import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBufAllocator; import io.netty.buffer.Unpooled; @@ -30,20 +31,6 @@ import io.netty.util.Recycler; import io.netty.util.Recycler.Handle; import io.netty.util.ReferenceCountUtil; -import java.io.File; -import java.io.FileInputStream; -import java.io.FileOutputStream; -import java.io.IOException; -import java.nio.ByteBuffer; -import java.nio.channels.FileChannel; -import java.util.ArrayDeque; -import java.util.ArrayList; -import java.util.Collections; -import java.util.List; -import java.util.concurrent.ArrayBlockingQueue; -import java.util.concurrent.BlockingQueue; -import java.util.concurrent.ThreadFactory; -import java.util.concurrent.TimeUnit; import org.apache.bookkeeper.bookie.LedgerDirsManager.NoWritableLedgerDirException; import org.apache.bookkeeper.bookie.stats.JournalStats; import org.apache.bookkeeper.common.collections.BlockingMpscQueue; @@ -52,6 +39,7 @@ import org.apache.bookkeeper.common.util.affinity.CpuAffinity; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.processor.RequestProcessor; +import org.apache.bookkeeper.proto.BookieRequestHandler; import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteCallback; import org.apache.bookkeeper.stats.Counter; import org.apache.bookkeeper.stats.NullStatsLogger; @@ -63,6 +51,23 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.io.File; +import java.io.FileInputStream; +import java.io.FileOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.channels.FileChannel; +import java.util.ArrayDeque; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.TimeUnit; + /** * Provide journal related management. */ @@ -332,7 +337,11 @@ public void run() { callbackTime.addLatency(MathUtils.elapsedNanos(startTime), TimeUnit.NANOSECONDS); recycle(); } - + + public Object getCtx() { + return ctx; + } + private final Handle recyclerHandle; private QueueEntry(Handle recyclerHandle) { @@ -388,7 +397,11 @@ private void flushFileToDisk() throws IOException { flushed = true; } } - + + public RecyclableArrayList getForceWriteWaiters() { + return forceWriteWaiters; + } + public void closeFileIfNecessary() { // Close if shouldClose is set if (shouldClose) { @@ -494,7 +507,8 @@ public void run() { } journalStats.getForceWriteQueueSize().addCount(-requestsCount); - + + Set writeHandlers = new HashSet<>(); // Sync and mark the journal up to the position of the last entry in the batch ForceWriteRequest lastRequest = localRequests.get(requestsCount - 1); syncJournal(lastRequest); @@ -503,17 +517,22 @@ public void run() { // responses for (int i = 0; i < requestsCount; i++) { ForceWriteRequest req = localRequests.get(i); + req.getForceWriteWaiters().forEach(ele -> { + Object ctx = ele.getCtx(); + if (ctx instanceof BookieRequestHandler) { + writeHandlers.add((BookieRequestHandler) ctx); + } + }); numReqInLastForceWrite += req.process(); req.recycle(); } journalStats.getForceWriteGroupingCountStats() .registerSuccessfulValue(numReqInLastForceWrite); - - if (requestProcessor != null) { - requestProcessor.flushPendingResponses(); + + for (BookieRequestHandler writeHandler : writeHandlers) { + writeHandler.flushPendingResponse(); } - } catch (IOException ioe) { LOG.error("I/O exception in ForceWrite thread", ioe); running = false; diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java index 50b7969023e..5b050c0c8f7 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java @@ -36,6 +36,8 @@ public class BookieRequestHandler extends ChannelInboundHandlerAdapter { static final Object EVENT_FLUSH_ALL_PENDING_RESPONSES = new Object(); + + private static final int DEFAULT_PENDING_RESPONSE_SIZE = 256; private final RequestProcessor requestProcessor; private final ChannelGroup allChannels; @@ -43,7 +45,8 @@ public class BookieRequestHandler extends ChannelInboundHandlerAdapter { private ChannelHandlerContext ctx; private ByteBuf pendingSendResponses = null; - private int maxPendingResponsesSize; + private int maxPendingResponsesSize = DEFAULT_PENDING_RESPONSE_SIZE; + BookieRequestHandler(ServerConfiguration conf, RequestProcessor processor, ChannelGroup allChannels) { this.requestProcessor = processor; @@ -92,20 +95,34 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception public synchronized void prepareSendResponseV2(int rc, BookieProtocol.ParsedAddRequest req) { if (pendingSendResponses == null) { - pendingSendResponses = ctx.alloc().directBuffer(maxPendingResponsesSize != 0 - ? maxPendingResponsesSize : 256); + pendingSendResponses = ctx.alloc().directBuffer(maxPendingResponsesSize); } BookieProtoEncoding.ResponseEnDeCoderPreV3.serializeAddResponseInto(rc, req, pendingSendResponses); } - + + public synchronized void flushPendingResponse() { + if (pendingSendResponses != null) { + maxPendingResponsesSize = (int) Math.max( + maxPendingResponsesSize * 0.9 + 0.1 * pendingSendResponses.readableBytes(), + DEFAULT_PENDING_RESPONSE_SIZE); + if (ctx.channel().isActive()) { + ctx.writeAndFlush(pendingSendResponses, ctx.voidPromise()); + } else { + pendingSendResponses.release(); + } + pendingSendResponses = null; + } + } + @Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt == EVENT_FLUSH_ALL_PENDING_RESPONSES) { synchronized (this) { if (pendingSendResponses != null) { - maxPendingResponsesSize = Math.max(maxPendingResponsesSize, - pendingSendResponses.readableBytes()); + maxPendingResponsesSize = (int) Math.max( + maxPendingResponsesSize * 0.9 + 0.1 * pendingSendResponses.readableBytes(), + DEFAULT_PENDING_RESPONSE_SIZE); if (ctx.channel().isActive()) { ctx.writeAndFlush(pendingSendResponses, ctx.voidPromise()); } else { From 048bceee3821691f42dbb67adf107b6ddcb69ec3 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Wed, 8 Mar 2023 17:11:30 +0800 Subject: [PATCH 02/12] 1.Optimize flushPendingResponses, only flush the channel which trigger prepareSendResponseV2. 2.Make maxPendingResponsesSize is more reasonable. --- .../org/apache/bookkeeper/bookie/Journal.java | 25 +++++++++++-------- .../processor/RequestProcessor.java | 5 ---- .../proto/BookieRequestHandler.java | 25 ------------------- .../proto/BookieRequestProcessor.java | 7 ------ 4 files changed, 15 insertions(+), 47 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java index f034a2fdf29..1139c753c62 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java @@ -338,7 +338,7 @@ public void run() { recycle(); } - public Object getCtx() { + private Object getCtx() { return ctx; } @@ -398,7 +398,7 @@ private void flushFileToDisk() throws IOException { } } - public RecyclableArrayList getForceWriteWaiters() { + private RecyclableArrayList getForceWriteWaiters() { return forceWriteWaiters; } @@ -519,7 +519,8 @@ public void run() { ForceWriteRequest req = localRequests.get(i); req.getForceWriteWaiters().forEach(ele -> { Object ctx = ele.getCtx(); - if (ctx instanceof BookieRequestHandler) { + if (ctx instanceof BookieRequestHandler + && ele.entryId != BookieImpl.METAENTRY_ID_FORCE_LEDGER) { writeHandlers.add((BookieRequestHandler) ctx); } }); @@ -924,7 +925,7 @@ public void logAddEntry(long ledgerId, long entryId, ByteBuf entry, memoryLimitController.reserveMemory(entry.readableBytes()); queue.put(QueueEntry.create( - entry, ackBeforeSync, ledgerId, entryId, cb, ctx, MathUtils.nowInNano(), + entry, ackBeforeSync, ledgerId, entryId, cb, ctx, MathUtils.nowInNano(), journalStats.getJournalAddEntryStats(), callbackTime)); } @@ -1119,20 +1120,24 @@ journalFormatVersionToWrite, getBufferedChannelBuilder(), } journalFlushWatcher.reset().start(); bc.flush(); - + + Set writeHandlers = new HashSet<>(); for (int i = 0; i < toFlush.size(); i++) { QueueEntry entry = toFlush.get(i); if (entry != null && (!syncData || entry.ackBeforeSync)) { toFlush.set(i, null); numEntriesToFlush--; + if (entry.getCtx() instanceof BookieRequestHandler + && entry.entryId != BookieImpl.METAENTRY_ID_FORCE_LEDGER) { + writeHandlers.add((BookieRequestHandler) entry.getCtx()); + } entry.run(); } - - if (forceWriteThread.requestProcessor != null) { - forceWriteThread.requestProcessor.flushPendingResponses(); - } } - + for (BookieRequestHandler writeHandler : writeHandlers) { + writeHandler.flushPendingResponse(); + } + lastFlushPosition = bc.position(); journalStats.getJournalFlushStats().registerSuccessfulEvent( journalFlushWatcher.stop().elapsed(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/processor/RequestProcessor.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/processor/RequestProcessor.java index 9f9a0daf682..5a4238e64d6 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/processor/RequestProcessor.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/processor/RequestProcessor.java @@ -42,9 +42,4 @@ public interface RequestProcessor extends AutoCloseable { * channel received the given request r */ void processRequest(Object r, BookieRequestHandler channel); - - /** - * Flush any pending response staged on all the client connections. - */ - void flushPendingResponses(); } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java index 5b050c0c8f7..d9e2796ce08 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java @@ -35,8 +35,6 @@ @Slf4j public class BookieRequestHandler extends ChannelInboundHandlerAdapter { - static final Object EVENT_FLUSH_ALL_PENDING_RESPONSES = new Object(); - private static final int DEFAULT_PENDING_RESPONSE_SIZE = 256; private final RequestProcessor requestProcessor; @@ -97,7 +95,6 @@ public synchronized void prepareSendResponseV2(int rc, BookieProtocol.ParsedAddR if (pendingSendResponses == null) { pendingSendResponses = ctx.alloc().directBuffer(maxPendingResponsesSize); } - BookieProtoEncoding.ResponseEnDeCoderPreV3.serializeAddResponseInto(rc, req, pendingSendResponses); } @@ -114,26 +111,4 @@ public synchronized void flushPendingResponse() { pendingSendResponses = null; } } - - @Override - public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { - if (evt == EVENT_FLUSH_ALL_PENDING_RESPONSES) { - synchronized (this) { - if (pendingSendResponses != null) { - maxPendingResponsesSize = (int) Math.max( - maxPendingResponsesSize * 0.9 + 0.1 * pendingSendResponses.readableBytes(), - DEFAULT_PENDING_RESPONSE_SIZE); - if (ctx.channel().isActive()) { - ctx.writeAndFlush(pendingSendResponses, ctx.voidPromise()); - } else { - pendingSendResponses.release(); - } - - pendingSendResponses = null; - } - } - } else { - super.userEventTriggered(ctx, evt); - } - } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java index d07aa9cffa0..6e7f5abcbb7 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java @@ -699,13 +699,6 @@ private void processReadRequest(final BookieProtocol.ReadRequest r, final Bookie } } - @Override - public void flushPendingResponses() { - for (Channel c : allChannels) { - c.pipeline().fireUserEventTriggered(BookieRequestHandler.EVENT_FLUSH_ALL_PENDING_RESPONSES); - } - } - public long getWaitTimeoutOnBackpressureMillis() { return waitTimeoutOnBackpressureMillis; } From 20e3db7cbf319454b0eb6715809830159b95862a Mon Sep 17 00:00:00 2001 From: horizonzy Date: Wed, 8 Mar 2023 18:04:07 +0800 Subject: [PATCH 03/12] Fix style. --- .../org/apache/bookkeeper/bookie/Bookie.java | 3 - .../apache/bookkeeper/bookie/BookieImpl.java | 8 --- .../org/apache/bookkeeper/bookie/Journal.java | 56 ++++++++----------- .../proto/BookieRequestHandler.java | 3 +- .../apache/bookkeeper/proto/BookieServer.java | 1 - 5 files changed, 25 insertions(+), 46 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Bookie.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Bookie.java index ac9df53cd22..90c8acf5af4 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Bookie.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Bookie.java @@ -23,7 +23,6 @@ import java.util.PrimitiveIterator; import java.util.concurrent.CompletableFuture; import org.apache.bookkeeper.common.util.Watcher; -import org.apache.bookkeeper.processor.RequestProcessor; import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteCallback; /** @@ -87,8 +86,6 @@ void cancelWaitForLastAddConfirmedUpdate(long ledgerId, // TODO: Should be constructed and passed in as a parameter LedgerStorage getLedgerStorage(); - void setRequestProcessor(RequestProcessor requestProcessor); - // TODO: Move this exceptions somewhere else /** * Exception is thrown when no such a ledger is found in this bookie. diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieImpl.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieImpl.java index 2b76488cbe9..0db230d9d3d 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieImpl.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieImpl.java @@ -69,7 +69,6 @@ import org.apache.bookkeeper.net.BookieId; import org.apache.bookkeeper.net.BookieSocketAddress; import org.apache.bookkeeper.net.DNS; -import org.apache.bookkeeper.processor.RequestProcessor; import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteCallback; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.bookkeeper.stats.StatsLogger; @@ -1282,11 +1281,4 @@ public OfLong getListOfEntriesOfLedger(long ledgerId) throws IOException, NoLedg } } } - - @Override - public void setRequestProcessor(RequestProcessor requestProcessor) { - for (Journal journal : journals) { - journal.setRequestProcessor(requestProcessor); - } - } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java index 1139c753c62..9154b2cc23d 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java @@ -23,7 +23,6 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Stopwatch; -import com.scurrilous.circe.Hash; import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBufAllocator; import io.netty.buffer.Unpooled; @@ -31,6 +30,22 @@ import io.netty.util.Recycler; import io.netty.util.Recycler.Handle; import io.netty.util.ReferenceCountUtil; +import java.io.File; +import java.io.FileInputStream; +import java.io.FileOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.channels.FileChannel; +import java.util.ArrayDeque; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.TimeUnit; import org.apache.bookkeeper.bookie.LedgerDirsManager.NoWritableLedgerDirException; import org.apache.bookkeeper.bookie.stats.JournalStats; import org.apache.bookkeeper.common.collections.BlockingMpscQueue; @@ -38,7 +53,6 @@ import org.apache.bookkeeper.common.util.MemoryLimitController; import org.apache.bookkeeper.common.util.affinity.CpuAffinity; import org.apache.bookkeeper.conf.ServerConfiguration; -import org.apache.bookkeeper.processor.RequestProcessor; import org.apache.bookkeeper.proto.BookieRequestHandler; import org.apache.bookkeeper.proto.BookkeeperInternalCallbacks.WriteCallback; import org.apache.bookkeeper.stats.Counter; @@ -51,23 +65,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.io.File; -import java.io.FileInputStream; -import java.io.FileOutputStream; -import java.io.IOException; -import java.nio.ByteBuffer; -import java.nio.channels.FileChannel; -import java.util.ArrayDeque; -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashSet; -import java.util.List; -import java.util.Set; -import java.util.concurrent.ArrayBlockingQueue; -import java.util.concurrent.BlockingQueue; -import java.util.concurrent.ThreadFactory; -import java.util.concurrent.TimeUnit; - /** * Provide journal related management. */ @@ -337,11 +334,11 @@ public void run() { callbackTime.addLatency(MathUtils.elapsedNanos(startTime), TimeUnit.NANOSECONDS); recycle(); } - + private Object getCtx() { return ctx; } - + private final Handle recyclerHandle; private QueueEntry(Handle recyclerHandle) { @@ -397,11 +394,11 @@ private void flushFileToDisk() throws IOException { flushed = true; } } - + private RecyclableArrayList getForceWriteWaiters() { return forceWriteWaiters; } - + public void closeFileIfNecessary() { // Close if shouldClose is set if (shouldClose) { @@ -467,7 +464,6 @@ private class ForceWriteThread extends BookieCriticalThread { // successful force write Thread threadToNotifyOnEx; - RequestProcessor requestProcessor; // should we group force writes private final boolean enableGroupForceWrites; private final Counter forceWriteThreadTime; @@ -507,7 +503,7 @@ public void run() { } journalStats.getForceWriteQueueSize().addCount(-requestsCount); - + Set writeHandlers = new HashSet<>(); // Sync and mark the journal up to the position of the last entry in the batch ForceWriteRequest lastRequest = localRequests.get(requestsCount - 1); @@ -530,7 +526,7 @@ public void run() { journalStats.getForceWriteGroupingCountStats() .registerSuccessfulValue(numReqInLastForceWrite); - + for (BookieRequestHandler writeHandler : writeHandlers) { writeHandler.flushPendingResponse(); } @@ -1120,7 +1116,7 @@ journalFormatVersionToWrite, getBufferedChannelBuilder(), } journalFlushWatcher.reset().start(); bc.flush(); - + Set writeHandlers = new HashSet<>(); for (int i = 0; i < toFlush.size(); i++) { QueueEntry entry = toFlush.get(i); @@ -1137,7 +1133,7 @@ journalFormatVersionToWrite, getBufferedChannelBuilder(), for (BookieRequestHandler writeHandler : writeHandlers) { writeHandler.flushPendingResponse(); } - + lastFlushPosition = bc.position(); journalStats.getJournalFlushStats().registerSuccessfulEvent( journalFlushWatcher.stop().elapsed(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS); @@ -1254,10 +1250,6 @@ public BufferedChannelBuilder getBufferedChannelBuilder() { return (FileChannel fc, int capacity) -> new BufferedChannel(allocator, fc, capacity); } - public void setRequestProcessor(RequestProcessor requestProcessor) { - forceWriteThread.requestProcessor = requestProcessor; - } - /** * Shuts down the journal. */ diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java index d9e2796ce08..a4ad265187d 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java @@ -44,7 +44,6 @@ public class BookieRequestHandler extends ChannelInboundHandlerAdapter { private ByteBuf pendingSendResponses = null; private int maxPendingResponsesSize = DEFAULT_PENDING_RESPONSE_SIZE; - BookieRequestHandler(ServerConfiguration conf, RequestProcessor processor, ChannelGroup allChannels) { this.requestProcessor = processor; @@ -97,7 +96,7 @@ public synchronized void prepareSendResponseV2(int rc, BookieProtocol.ParsedAddR } BookieProtoEncoding.ResponseEnDeCoderPreV3.serializeAddResponseInto(rc, req, pendingSendResponses); } - + public synchronized void flushPendingResponse() { if (pendingSendResponses != null) { maxPendingResponsesSize = (int) Math.max( diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieServer.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieServer.java index caff467db36..e50a09dba8c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieServer.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieServer.java @@ -106,7 +106,6 @@ public BookieServer(ServerConfiguration conf, this.requestProcessor = new BookieRequestProcessor(conf, bookie, statsLogger.scope(SERVER_SCOPE), shFactory, allocator, nettyServer.allChannels); this.nettyServer.setRequestProcessor(this.requestProcessor); - this.bookie.setRequestProcessor(this.requestProcessor); } /** From b46fe9c61c4268cefb2912bf5c63cbb70b63393b Mon Sep 17 00:00:00 2001 From: horizonzy Date: Wed, 8 Mar 2023 23:16:56 +0800 Subject: [PATCH 04/12] Fix findbugs problem. --- .../org/apache/bookkeeper/proto/BookieRequestHandler.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java index a4ad265187d..9336f24ab0a 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java @@ -56,7 +56,7 @@ public ChannelHandlerContext ctx() { @Override public void channelActive(ChannelHandlerContext ctx) throws Exception { - log.info("Channel connected {}", ctx.channel()); + log.info("Channel connected {}", ctx.channel()); this.ctx = ctx; super.channelActive(ctx); } @@ -92,7 +92,7 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception public synchronized void prepareSendResponseV2(int rc, BookieProtocol.ParsedAddRequest req) { if (pendingSendResponses == null) { - pendingSendResponses = ctx.alloc().directBuffer(maxPendingResponsesSize); + pendingSendResponses = ctx().alloc().directBuffer(maxPendingResponsesSize); } BookieProtoEncoding.ResponseEnDeCoderPreV3.serializeAddResponseInto(rc, req, pendingSendResponses); } @@ -102,8 +102,8 @@ public synchronized void flushPendingResponse() { maxPendingResponsesSize = (int) Math.max( maxPendingResponsesSize * 0.9 + 0.1 * pendingSendResponses.readableBytes(), DEFAULT_PENDING_RESPONSE_SIZE); - if (ctx.channel().isActive()) { - ctx.writeAndFlush(pendingSendResponses, ctx.voidPromise()); + if (ctx().channel().isActive()) { + ctx().writeAndFlush(pendingSendResponses, ctx.voidPromise()); } else { pendingSendResponses.release(); } From d4e59b416b9178f89380621be5799f9748f2ca52 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Thu, 9 Mar 2023 00:03:51 +0800 Subject: [PATCH 05/12] Fix ci failed. --- .../src/main/java/org/apache/bookkeeper/bookie/Journal.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java index 9154b2cc23d..43b68ec9266 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java @@ -514,6 +514,9 @@ public void run() { for (int i = 0; i < requestsCount; i++) { ForceWriteRequest req = localRequests.get(i); req.getForceWriteWaiters().forEach(ele -> { + if (ele == null) { + return; + } Object ctx = ele.getCtx(); if (ctx instanceof BookieRequestHandler && ele.entryId != BookieImpl.METAENTRY_ID_FORCE_LEDGER) { From bfd73e60756b9d644ed9e4230b10777bb405ded2 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Thu, 9 Mar 2023 19:46:11 +0800 Subject: [PATCH 06/12] Address the comments. --- bookkeeper-server/pom.xml | 4 ++ .../org/apache/bookkeeper/bookie/Journal.java | 50 +++++++++---------- .../proto/BookieRequestHandler.java | 2 +- pom.xml | 7 +++ 4 files changed, 35 insertions(+), 28 deletions(-) diff --git a/bookkeeper-server/pom.xml b/bookkeeper-server/pom.xml index 82dbadec473..a46b20c23d7 100644 --- a/bookkeeper-server/pom.xml +++ b/bookkeeper-server/pom.xml @@ -149,6 +149,10 @@ runtime true + + com.carrotsearch + hppc + org.apache.bookkeeper diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java index 43b68ec9266..176c7a8d47a 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java @@ -21,6 +21,8 @@ package org.apache.bookkeeper.bookie; +import com.carrotsearch.hppc.ObjectHashSet; +import com.carrotsearch.hppc.cursors.ObjectCursor; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Stopwatch; import io.netty.buffer.ByteBuf; @@ -374,13 +376,17 @@ public static class ForceWriteRequest { private long logId; private boolean flushed; - public int process() { + public int process(ObjectHashSet writeHandlers) { closeFileIfNecessary(); // Notify the waiters that the force write succeeded for (int i = 0; i < forceWriteWaiters.size(); i++) { QueueEntry qe = forceWriteWaiters.get(i); if (qe != null) { + if (qe.getCtx() instanceof BookieRequestHandler + && qe.entryId != BookieImpl.METAENTRY_ID_FORCE_LEDGER) { + writeHandlers.add((BookieRequestHandler) qe.getCtx()); + } qe.run(); } } @@ -395,10 +401,6 @@ private void flushFileToDisk() throws IOException { } } - private RecyclableArrayList getForceWriteWaiters() { - return forceWriteWaiters; - } - public void closeFileIfNecessary() { // Close if shouldClose is set if (shouldClose) { @@ -490,7 +492,8 @@ public void run() { } final List localRequests = new ArrayList<>(); - + final ObjectHashSet writeHandlers = new ObjectHashSet<>(); + while (running) { try { int numReqInLastForceWrite = 0; @@ -503,8 +506,7 @@ public void run() { } journalStats.getForceWriteQueueSize().addCount(-requestsCount); - - Set writeHandlers = new HashSet<>(); + // Sync and mark the journal up to the position of the last entry in the batch ForceWriteRequest lastRequest = localRequests.get(requestsCount - 1); syncJournal(lastRequest); @@ -513,25 +515,17 @@ public void run() { // responses for (int i = 0; i < requestsCount; i++) { ForceWriteRequest req = localRequests.get(i); - req.getForceWriteWaiters().forEach(ele -> { - if (ele == null) { - return; - } - Object ctx = ele.getCtx(); - if (ctx instanceof BookieRequestHandler - && ele.entryId != BookieImpl.METAENTRY_ID_FORCE_LEDGER) { - writeHandlers.add((BookieRequestHandler) ctx); - } - }); - numReqInLastForceWrite += req.process(); + numReqInLastForceWrite += req.process(writeHandlers); req.recycle(); } journalStats.getForceWriteGroupingCountStats() .registerSuccessfulValue(numReqInLastForceWrite); - - for (BookieRequestHandler writeHandler : writeHandlers) { - writeHandler.flushPendingResponse(); + if (writeHandlers.size() > 0) { + for (ObjectCursor writeHandler : writeHandlers) { + writeHandler.value.flushPendingResponse(); + } + writeHandlers.clear(); } } catch (IOException ioe) { LOG.error("I/O exception in ForceWrite thread", ioe); @@ -1002,7 +996,8 @@ public void run() { long busyStartTime = System.nanoTime(); ArrayDeque localQueueEntries = new ArrayDeque<>(); - + final ObjectHashSet writeHandlers = new ObjectHashSet<>(); + QueueEntry qe = null; while (true) { // new journal file to write @@ -1120,7 +1115,6 @@ journalFormatVersionToWrite, getBufferedChannelBuilder(), journalFlushWatcher.reset().start(); bc.flush(); - Set writeHandlers = new HashSet<>(); for (int i = 0; i < toFlush.size(); i++) { QueueEntry entry = toFlush.get(i); if (entry != null && (!syncData || entry.ackBeforeSync)) { @@ -1133,10 +1127,12 @@ journalFormatVersionToWrite, getBufferedChannelBuilder(), entry.run(); } } - for (BookieRequestHandler writeHandler : writeHandlers) { - writeHandler.flushPendingResponse(); + if (writeHandlers.size() > 0) { + for (ObjectCursor writeHandler : writeHandlers) { + writeHandler.value.flushPendingResponse(); + } + writeHandlers.clear(); } - lastFlushPosition = bc.position(); journalStats.getJournalFlushStats().registerSuccessfulEvent( journalFlushWatcher.stop().elapsed(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java index 9336f24ab0a..3d906dba449 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestHandler.java @@ -100,7 +100,7 @@ public synchronized void prepareSendResponseV2(int rc, BookieProtocol.ParsedAddR public synchronized void flushPendingResponse() { if (pendingSendResponses != null) { maxPendingResponsesSize = (int) Math.max( - maxPendingResponsesSize * 0.9 + 0.1 * pendingSendResponses.readableBytes(), + maxPendingResponsesSize * 0.5 + 0.5 * pendingSendResponses.readableBytes(), DEFAULT_PENDING_RESPONSE_SIZE); if (ctx().channel().isActive()) { ctx().writeAndFlush(pendingSendResponses, ctx.voidPromise()); diff --git a/pom.xml b/pom.xml index e8227092bb3..b249a6218ca 100644 --- a/pom.xml +++ b/pom.xml @@ -176,6 +176,7 @@ 3.8.1 1.1.7.7 2.1.2 + 0.9.1 0.12 2.7 @@ -801,6 +802,12 @@ rxjava ${rxjava.version} + + + com.carrotsearch + hppc + ${hppc.version} + From 4085c6d22390f5cef569a8f921d93f4b8ae2ceb5 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Thu, 9 Mar 2023 19:58:41 +0800 Subject: [PATCH 07/12] Fix style. --- .../main/java/org/apache/bookkeeper/bookie/Journal.java | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java index 176c7a8d47a..95bb7c7df9c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java @@ -41,9 +41,7 @@ import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Collections; -import java.util.HashSet; import java.util.List; -import java.util.Set; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; import java.util.concurrent.ThreadFactory; @@ -493,7 +491,7 @@ public void run() { final List localRequests = new ArrayList<>(); final ObjectHashSet writeHandlers = new ObjectHashSet<>(); - + while (running) { try { int numReqInLastForceWrite = 0; @@ -506,7 +504,7 @@ public void run() { } journalStats.getForceWriteQueueSize().addCount(-requestsCount); - + // Sync and mark the journal up to the position of the last entry in the batch ForceWriteRequest lastRequest = localRequests.get(requestsCount - 1); syncJournal(lastRequest); @@ -997,7 +995,7 @@ public void run() { long busyStartTime = System.nanoTime(); ArrayDeque localQueueEntries = new ArrayDeque<>(); final ObjectHashSet writeHandlers = new ObjectHashSet<>(); - + QueueEntry qe = null; while (true) { // new journal file to write From 82a849476d7c269b431657b15b83996d0acf2756 Mon Sep 17 00:00:00 2001 From: Yan Zhao Date: Fri, 10 Mar 2023 00:01:22 +0800 Subject: [PATCH 08/12] Update bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java address the comments. Co-authored-by: Matteo Merli --- .../main/java/org/apache/bookkeeper/bookie/Journal.java | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java index 95bb7c7df9c..9388a3c1c5d 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java @@ -519,12 +519,8 @@ public void run() { journalStats.getForceWriteGroupingCountStats() .registerSuccessfulValue(numReqInLastForceWrite); - if (writeHandlers.size() > 0) { - for (ObjectCursor writeHandler : writeHandlers) { - writeHandler.value.flushPendingResponse(); - } - writeHandlers.clear(); - } + writeHandlers.forEach(wh -> wh.flushPendingResponse()); + writeHandlers.clear(); } catch (IOException ioe) { LOG.error("I/O exception in ForceWrite thread", ioe); running = false; From 3d19988d04ba9a56916ede9b6f0da43c54a9aece Mon Sep 17 00:00:00 2001 From: horizonzy Date: Fri, 10 Mar 2023 00:01:57 +0800 Subject: [PATCH 09/12] add license. --- bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt | 1 + 1 file changed, 1 insertion(+) diff --git a/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt b/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt index b66bc8b749f..631c43159de 100644 --- a/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt +++ b/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt @@ -320,6 +320,7 @@ Apache Software License, Version 2. - lib/org.xerial.snappy-snappy-java-1.1.7.7.jar [50] - lib/io.reactivex.rxjava3-rxjava-3.0.1.jar [51] - lib/org.hdrhistogram-HdrHistogram-2.1.10.jar [52] +- lib/com.carrotsearch-hppc-0.9.1.jar [53] [1] Source available at https://github.com/FasterXML/jackson-annotations/tree/jackson-annotations-2.13.4 [2] Source available at https://github.com/FasterXML/jackson-core/tree/jackson-core-2.13.4 From f2604fc0272e5c13cbdb95815c53b48a4a0e1cf7 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Fri, 10 Mar 2023 00:09:08 +0800 Subject: [PATCH 10/12] address the comments. --- .../org/apache/bookkeeper/bookie/Journal.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java index 9388a3c1c5d..50eb6054146 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/Journal.java @@ -22,7 +22,7 @@ package org.apache.bookkeeper.bookie; import com.carrotsearch.hppc.ObjectHashSet; -import com.carrotsearch.hppc.cursors.ObjectCursor; +import com.carrotsearch.hppc.procedures.ObjectProcedure; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Stopwatch; import io.netty.buffer.ByteBuf; @@ -519,7 +519,9 @@ public void run() { journalStats.getForceWriteGroupingCountStats() .registerSuccessfulValue(numReqInLastForceWrite); - writeHandlers.forEach(wh -> wh.flushPendingResponse()); + writeHandlers.forEach( + (ObjectProcedure) + BookieRequestHandler::flushPendingResponse); writeHandlers.clear(); } catch (IOException ioe) { LOG.error("I/O exception in ForceWrite thread", ioe); @@ -1121,12 +1123,10 @@ journalFormatVersionToWrite, getBufferedChannelBuilder(), entry.run(); } } - if (writeHandlers.size() > 0) { - for (ObjectCursor writeHandler : writeHandlers) { - writeHandler.value.flushPendingResponse(); - } - writeHandlers.clear(); - } + writeHandlers.forEach( + (ObjectProcedure) + BookieRequestHandler::flushPendingResponse); + writeHandlers.clear(); lastFlushPosition = bc.position(); journalStats.getJournalFlushStats().registerSuccessfulEvent( journalFlushWatcher.stop().elapsed(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS); From 176aa25adc550e39d0522c3f4e2ad039a549e6ea Mon Sep 17 00:00:00 2001 From: horizonzy Date: Fri, 10 Mar 2023 01:06:26 +0800 Subject: [PATCH 11/12] add lisence. --- bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt | 1 + 1 file changed, 1 insertion(+) diff --git a/bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt b/bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt index cde305a40bf..670e81d7a99 100644 --- a/bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt +++ b/bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt @@ -316,6 +316,7 @@ Apache Software License, Version 2. - lib/org.conscrypt-conscrypt-openjdk-uber-2.5.1.jar [49] - lib/org.xerial.snappy-snappy-java-1.1.7.7.jar [50] - lib/io.reactivex.rxjava3-rxjava-3.0.1.jar [51] +- lib/com.carrotsearch-hppc-0.9.1.jar [52] [1] Source available at https://github.com/FasterXML/jackson-annotations/tree/jackson-annotations-2.13.4 [2] Source available at https://github.com/FasterXML/jackson-core/tree/jackson-core-2.13.4 From 29f1cf79390327446425ce970f5cb52b51f032e2 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Fri, 10 Mar 2023 10:24:31 +0800 Subject: [PATCH 12/12] Fix license. --- bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt | 1 + bookkeeper-dist/src/main/resources/LICENSE-bkctl.bin.txt | 2 ++ bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt | 1 + 3 files changed, 4 insertions(+) diff --git a/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt b/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt index 631c43159de..0fce3d3d325 100644 --- a/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt +++ b/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt @@ -370,6 +370,7 @@ Apache Software License, Version 2. [50] Source available at https://github.com/google/snappy/releases/tag/1.1.7.7 [51] Source available at https://github.com/ReactiveX/RxJava/tree/v3.0.1 [52] Source available at https://github.com/HdrHistogram/HdrHistogram/tree/HdrHistogram-2.1.10 +[53] Source available at https://github.com/carrotsearch/hppc/tree/0.9.1 ------------------------------------------------------------------------------------ lib/io.netty-netty-codec-4.1.89.Final.jar bundles some 3rd party dependencies diff --git a/bookkeeper-dist/src/main/resources/LICENSE-bkctl.bin.txt b/bookkeeper-dist/src/main/resources/LICENSE-bkctl.bin.txt index 013d46207a5..e8ac8cc2c4f 100644 --- a/bookkeeper-dist/src/main/resources/LICENSE-bkctl.bin.txt +++ b/bookkeeper-dist/src/main/resources/LICENSE-bkctl.bin.txt @@ -291,6 +291,7 @@ Apache Software License, Version 2. - lib/org.conscrypt-conscrypt-openjdk-uber-2.5.1.jar [49] - lib/org.xerial.snappy-snappy-java-1.1.7.7.jar [50] - lib/io.reactivex.rxjava3-rxjava-3.0.1.jar [51] +- lib/com.carrotsearch-hppc-0.9.1.jar [52] [1] Source available at https://github.com/FasterXML/jackson-annotations/tree/jackson-annotations-2.13.4 [2] Source available at https://github.com/FasterXML/jackson-core/tree/jackson-core-2.13.4 @@ -331,6 +332,7 @@ Apache Software License, Version 2. [49] Source available at https://github.com/google/conscrypt/releases/tag/2.5.1 [50] Source available at https://github.com/google/snappy/releases/tag/1.1.7.7 [51] Source available at https://github.com/ReactiveX/RxJava/tree/v3.0.1 +[52] Source available at https://github.com/carrotsearch/hppc/tree/0.9.1 ------------------------------------------------------------------------------------ lib/io.netty-netty-codec-4.1.89.Final.jar bundles some 3rd party dependencies diff --git a/bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt b/bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt index 670e81d7a99..51e2d3d5c81 100644 --- a/bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt +++ b/bookkeeper-dist/src/main/resources/LICENSE-server.bin.txt @@ -365,6 +365,7 @@ Apache Software License, Version 2. [49] Source available at https://github.com/google/conscrypt/releases/tag/2.5.1 [50] Source available at https://github.com/google/snappy/releases/tag/1.1.7.7 [51] Source available at https://github.com/ReactiveX/RxJava/tree/v3.0.1 +[52] Source available at https://github.com/carrotsearch/hppc/tree/0.9.1 ------------------------------------------------------------------------------------ lib/io.netty-netty-codec-4.1.89.Final.jar bundles some 3rd party dependencies