From 9f00279662ac5898f194d7ddea2d425ebfbbf4ef Mon Sep 17 00:00:00 2001 From: Pradeep Nagaraju Date: Wed, 10 Nov 2021 00:25:49 -0800 Subject: [PATCH 1/8] Reorder the sequence of the bookkeeper server shutdown so that there are no read/ write ops while shutting down the bookie --- .../org/apache/bookkeeper/proto/BookieRequestProcessor.java | 2 +- .../main/java/org/apache/bookkeeper/proto/BookieServer.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) 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 67f83e9ce5f..789680f9afe 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 @@ -294,7 +294,7 @@ private OrderedExecutor createExecutor( private void shutdownExecutor(OrderedExecutor service) { if (null != service) { - service.shutdown(); + service.forceShutdown(1000, TimeUnit.SECONDS); } } 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 4e6a7821d44..50865f11eb0 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 @@ -216,12 +216,12 @@ public void resumeProcessing() { public synchronized void shutdown() { LOG.info("Shutting down BookieServer"); - this.nettyServer.shutdown(); if (!running) { return; } - exitCode = bookie.shutdown(); + this.nettyServer.shutdown(); this.requestProcessor.close(); + exitCode = bookie.shutdown(); running = false; } From 7fec12a1b0580d6094b481902be3d456af5eb00a Mon Sep 17 00:00:00 2001 From: Pradeep Nagaraju Date: Thu, 11 Nov 2021 23:38:44 -0800 Subject: [PATCH 2/8] Fix the force shutdown on the requestProcess; Check if the channel is active before sending response because it can be closed while responding; make bookie process as PID=1 in docker so that it can receive SIGINT --- .../proto/BookieRequestProcessor.java | 4 +- .../proto/PacketProcessorBaseV3.java | 37 ++++++++++--------- docker/scripts/entrypoint.sh | 2 +- 3 files changed, 23 insertions(+), 20 deletions(-) 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 789680f9afe..4af6eec8386 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 @@ -264,6 +264,7 @@ int maxReadsInProgressCount() { @Override public void close() { + LOG.info("Closing RequestProcessor"); shutdownExecutor(writeThreadPool); shutdownExecutor(readThreadPool); if (serverCfg.getNumLongPollWorkerThreads() > 0 || readThreadPool == null) { @@ -294,7 +295,8 @@ private OrderedExecutor createExecutor( private void shutdownExecutor(OrderedExecutor service) { if (null != service) { - service.forceShutdown(1000, TimeUnit.SECONDS); + service.shutdown(); + service.forceShutdown(10, TimeUnit.SECONDS); } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java index 15765a252b2..866f27a283c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java @@ -87,25 +87,26 @@ protected void sendResponse(StatusCode code, Object response, OpStatsLogger stat requestProcessor.invalidateBlacklist(channel); } } - - channel.writeAndFlush(response).addListener(new ChannelFutureListener() { - @Override - public void operationComplete(ChannelFuture future) throws Exception { - long writeElapsedNanos = MathUtils.elapsedNanos(writeNanos); - if (!future.isSuccess()) { - requestProcessor.getRequestStats().getChannelWriteStats() - .registerFailedEvent(writeElapsedNanos, TimeUnit.NANOSECONDS); - } else { - requestProcessor.getRequestStats().getChannelWriteStats() - .registerSuccessfulEvent(writeElapsedNanos, TimeUnit.NANOSECONDS); - } - if (StatusCode.EOK == code) { - statsLogger.registerSuccessfulEvent(MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS); - } else { - statsLogger.registerFailedEvent(MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS); + if (channel.isActive()) { + channel.writeAndFlush(response).addListener(new ChannelFutureListener() { + @Override + public void operationComplete(ChannelFuture future) throws Exception { + long writeElapsedNanos = MathUtils.elapsedNanos(writeNanos); + if (!future.isSuccess()) { + requestProcessor.getRequestStats().getChannelWriteStats() + .registerFailedEvent(writeElapsedNanos, TimeUnit.NANOSECONDS); + } else { + requestProcessor.getRequestStats().getChannelWriteStats() + .registerSuccessfulEvent(writeElapsedNanos, TimeUnit.NANOSECONDS); + } + if (StatusCode.EOK == code) { + statsLogger.registerSuccessfulEvent(MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS); + } else { + statsLogger.registerFailedEvent(MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS); + } } - } - }); + }); + } } protected boolean isVersionCompatible() { diff --git a/docker/scripts/entrypoint.sh b/docker/scripts/entrypoint.sh index 657eb6b918e..b71b326832f 100755 --- a/docker/scripts/entrypoint.sh +++ b/docker/scripts/entrypoint.sh @@ -41,7 +41,7 @@ function run_command() { chmod -R +x ${BINDIR} chmod -R +x ${SCRIPTS_DIR} echo "This is root, will use user $BK_USER to run command '$@'" - sudo -s -E -u "$BK_USER" /bin/bash "$@" + exec sudo -s -E -u "$BK_USER" /bin/bash "$@" exit else echo "Run command '$@'" From 66b2513031aea7347d0bb30392cc906b95d3ce2e Mon Sep 17 00:00:00 2001 From: Pradeep Nagaraju Date: Mon, 15 Nov 2021 17:28:56 -0800 Subject: [PATCH 3/8] Revert the Dockefile changes --- docker/scripts/entrypoint.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docker/scripts/entrypoint.sh b/docker/scripts/entrypoint.sh index b71b326832f..657eb6b918e 100755 --- a/docker/scripts/entrypoint.sh +++ b/docker/scripts/entrypoint.sh @@ -41,7 +41,7 @@ function run_command() { chmod -R +x ${BINDIR} chmod -R +x ${SCRIPTS_DIR} echo "This is root, will use user $BK_USER to run command '$@'" - exec sudo -s -E -u "$BK_USER" /bin/bash "$@" + sudo -s -E -u "$BK_USER" /bin/bash "$@" exit else echo "Run command '$@'" From 5a29761d8023658a46b3e68d17ecad65d27a949f Mon Sep 17 00:00:00 2001 From: Pradeep Nagaraju Date: Wed, 17 Nov 2021 18:00:11 -0800 Subject: [PATCH 4/8] Add channel.isActive() mock --- .../java/org/apache/bookkeeper/proto/PacketProcessorBase.java | 4 +++- .../apache/bookkeeper/proto/ForceLedgerProcessorV3Test.java | 1 + .../org/apache/bookkeeper/proto/WriteEntryProcessorTest.java | 1 + .../apache/bookkeeper/proto/WriteEntryProcessorV3Test.java | 1 + 4 files changed, 6 insertions(+), 1 deletion(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java index 4cc7176ede3..383716791b3 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java @@ -66,7 +66,9 @@ protected boolean isVersionCompatible() { } protected void sendResponse(int rc, Object response, OpStatsLogger statsLogger) { - channel.writeAndFlush(response, channel.voidPromise()); + if (channel.isActive()) { + channel.writeAndFlush(response, channel.voidPromise()); + } if (BookieProtocol.EOK == rc) { statsLogger.registerSuccessfulEvent(MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS); } else { diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3Test.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3Test.java index 3442adede08..6c6eea7adbf 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3Test.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ForceLedgerProcessorV3Test.java @@ -77,6 +77,7 @@ public void setup() { when(requestProcessor.getBookie()).thenReturn(bookie); when(requestProcessor.getWaitTimeoutOnBackpressureMillis()).thenReturn(-1L); when(requestProcessor.getRequestStats()).thenReturn(new RequestStats(NullStatsLogger.INSTANCE)); + when(channel.isActive()).thenReturn(true); processor = new ForceLedgerProcessorV3( request, channel, diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/WriteEntryProcessorTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/WriteEntryProcessorTest.java index 21bca29ff58..b249d4847c7 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/WriteEntryProcessorTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/WriteEntryProcessorTest.java @@ -71,6 +71,7 @@ public void setup() { requestProcessor = mock(BookieRequestProcessor.class); when(requestProcessor.getBookie()).thenReturn(bookie); when(requestProcessor.getRequestStats()).thenReturn(new RequestStats(NullStatsLogger.INSTANCE)); + when(channel.isActive()).thenReturn(true); processor = WriteEntryProcessor.create( request, channel, diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/WriteEntryProcessorV3Test.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/WriteEntryProcessorV3Test.java index 76ead525d01..3c81d73a5e7 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/WriteEntryProcessorV3Test.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/WriteEntryProcessorV3Test.java @@ -82,6 +82,7 @@ public void setup() { when(requestProcessor.getBookie()).thenReturn(bookie); when(requestProcessor.getWaitTimeoutOnBackpressureMillis()).thenReturn(-1L); when(requestProcessor.getRequestStats()).thenReturn(new RequestStats(NullStatsLogger.INSTANCE)); + when(channel.isActive()).thenReturn(true); processor = new WriteEntryProcessorV3( request, channel, From f8c2525e01611b008a2f6026d21dc1458ba0ab9d Mon Sep 17 00:00:00 2001 From: Pradeep Nagaraju Date: Wed, 17 Nov 2021 19:18:41 -0800 Subject: [PATCH 5/8] Revert nettyServer shutdown sequence --- .../src/main/java/org/apache/bookkeeper/proto/BookieServer.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 50865f11eb0..127a737efc3 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 @@ -216,10 +216,10 @@ public void resumeProcessing() { public synchronized void shutdown() { LOG.info("Shutting down BookieServer"); + this.nettyServer.shutdown(); if (!running) { return; } - this.nettyServer.shutdown(); this.requestProcessor.close(); exitCode = bookie.shutdown(); running = false; From b7908b9d8cfa757cac927303ea1cdb2dd7f93ff5 Mon Sep 17 00:00:00 2001 From: Pradeep Nagaraju Date: Thu, 18 Nov 2021 10:25:37 -0800 Subject: [PATCH 6/8] Added a log statement at the close of request processor for tracking time of the closure --- .../java/org/apache/bookkeeper/proto/BookieRequestProcessor.java | 1 + 1 file changed, 1 insertion(+) 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 4af6eec8386..902e2c1b419 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 @@ -272,6 +272,7 @@ public void close() { } shutdownExecutor(highPriorityThreadPool); requestTimer.stop(); + LOG.info("Closed RequestProcessor"); } private OrderedExecutor createExecutor( From e6b5501ecc3619ff127ccc970c6481ba3211e829 Mon Sep 17 00:00:00 2001 From: Pradeep Nagaraju Date: Tue, 23 Nov 2021 10:50:16 -0800 Subject: [PATCH 7/8] Adding logs when the netty channel is inactive to bypass the writeAndFlush --- .../java/org/apache/bookkeeper/proto/PacketProcessorBase.java | 2 ++ .../java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java | 2 ++ 2 files changed, 4 insertions(+) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java index 383716791b3..8ac71478b30 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java @@ -68,6 +68,8 @@ protected boolean isVersionCompatible() { protected void sendResponse(int rc, Object response, OpStatsLogger statsLogger) { if (channel.isActive()) { channel.writeAndFlush(response, channel.voidPromise()); + } else { + LOGGER.info("Netty channel is inactive, hence bypassing netty channel writeAndFlush during sendResponse"); } if (BookieProtocol.EOK == rc) { statsLogger.registerSuccessfulEvent(MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java index 866f27a283c..42e8e12b8b2 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java @@ -106,6 +106,8 @@ public void operationComplete(ChannelFuture future) throws Exception { } } }); + } else { + LOGGER.info("Netty channel is inactive, hence bypassing netty channel writeAndFlush during sendResponse"); } } From 15abe17c9c2e1761cc0ad8a225acf71a3c029649 Mon Sep 17 00:00:00 2001 From: Pradeep Nagaraju Date: Tue, 30 Nov 2021 12:29:36 -0800 Subject: [PATCH 8/8] Adding the channel information to the debug logs while the channel is inactive during sendResponse() --- .../java/org/apache/bookkeeper/proto/PacketProcessorBase.java | 3 ++- .../org/apache/bookkeeper/proto/PacketProcessorBaseV3.java | 3 ++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java index 8ac71478b30..d416b9f1417 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBase.java @@ -69,7 +69,8 @@ protected void sendResponse(int rc, Object response, OpStatsLogger statsLogger) if (channel.isActive()) { channel.writeAndFlush(response, channel.voidPromise()); } else { - LOGGER.info("Netty channel is inactive, hence bypassing netty channel writeAndFlush during sendResponse"); + LOGGER.debug("Netty channel {} is inactive, " + + "hence bypassing netty channel writeAndFlush during sendResponse", channel); } if (BookieProtocol.EOK == rc) { statsLogger.registerSuccessfulEvent(MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java index 42e8e12b8b2..d4ad65ba437 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PacketProcessorBaseV3.java @@ -107,7 +107,8 @@ public void operationComplete(ChannelFuture future) throws Exception { } }); } else { - LOGGER.info("Netty channel is inactive, hence bypassing netty channel writeAndFlush during sendResponse"); + LOGGER.debug("Netty channel {} is inactive, " + + "hence bypassing netty channel writeAndFlush during sendResponse", channel); } }