diff --git a/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt b/bookkeeper-dist/src/main/resources/LICENSE-all.bin.txt
index b66bc8b749f..0fce3d3d325 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
@@ -369,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 cde305a40bf..51e2d3d5c81 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
@@ -364,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
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/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 7966f6d2abf..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
@@ -21,6 +21,8 @@
package org.apache.bookkeeper.bookie;
+import com.carrotsearch.hppc.ObjectHashSet;
+import com.carrotsearch.hppc.procedures.ObjectProcedure;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Stopwatch;
import io.netty.buffer.ByteBuf;
@@ -51,7 +53,7 @@
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;
import org.apache.bookkeeper.stats.NullStatsLogger;
@@ -333,6 +335,10 @@ public void run() {
recycle();
}
+ private Object getCtx() {
+ return ctx;
+ }
+
private final Handle recyclerHandle;
private QueueEntry(Handle recyclerHandle) {
@@ -368,13 +374,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();
}
}
@@ -454,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;
@@ -481,6 +490,7 @@ public void run() {
}
final List localRequests = new ArrayList<>();
+ final ObjectHashSet writeHandlers = new ObjectHashSet<>();
while (running) {
try {
@@ -503,17 +513,16 @@ public void run() {
// responses
for (int i = 0; i < requestsCount; i++) {
ForceWriteRequest req = localRequests.get(i);
- numReqInLastForceWrite += req.process();
+ numReqInLastForceWrite += req.process(writeHandlers);
req.recycle();
}
journalStats.getForceWriteGroupingCountStats()
.registerSuccessfulValue(numReqInLastForceWrite);
-
- if (requestProcessor != null) {
- requestProcessor.flushPendingResponses();
- }
-
+ writeHandlers.forEach(
+ (ObjectProcedure super BookieRequestHandler>)
+ BookieRequestHandler::flushPendingResponse);
+ writeHandlers.clear();
} catch (IOException ioe) {
LOG.error("I/O exception in ForceWrite thread", ioe);
running = false;
@@ -905,7 +914,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));
}
@@ -983,6 +992,7 @@ public void run() {
long busyStartTime = System.nanoTime();
ArrayDeque localQueueEntries = new ArrayDeque<>();
+ final ObjectHashSet writeHandlers = new ObjectHashSet<>();
QueueEntry qe = null;
while (true) {
@@ -1106,14 +1116,17 @@ journalFormatVersionToWrite, getBufferedChannelBuilder(),
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();
- }
}
-
+ writeHandlers.forEach(
+ (ObjectProcedure super BookieRequestHandler>)
+ BookieRequestHandler::flushPendingResponse);
+ writeHandlers.clear();
lastFlushPosition = bc.position();
journalStats.getJournalFlushStats().registerSuccessfulEvent(
journalFlushWatcher.stop().elapsed(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS);
@@ -1230,10 +1243,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/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 50b7969023e..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
@@ -35,7 +35,7 @@
@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;
private final ChannelGroup allChannels;
@@ -43,7 +43,7 @@ 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;
@@ -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,31 +92,22 @@ 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);
}
- @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());
- if (ctx.channel().isActive()) {
- ctx.writeAndFlush(pendingSendResponses, ctx.voidPromise());
- } else {
- pendingSendResponses.release();
- }
-
- pendingSendResponses = null;
- }
+ public synchronized void flushPendingResponse() {
+ if (pendingSendResponses != null) {
+ maxPendingResponsesSize = (int) Math.max(
+ maxPendingResponsesSize * 0.5 + 0.5 * pendingSendResponses.readableBytes(),
+ DEFAULT_PENDING_RESPONSE_SIZE);
+ if (ctx().channel().isActive()) {
+ ctx().writeAndFlush(pendingSendResponses, ctx.voidPromise());
+ } else {
+ pendingSendResponses.release();
}
- } else {
- super.userEventTriggered(ctx, evt);
+ pendingSendResponses = null;
}
}
}
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;
}
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);
}
/**
diff --git a/pom.xml b/pom.xml
index b66c0a725e6..f1926121fa1 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}
+