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..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 @@ -264,6 +264,7 @@ int maxReadsInProgressCount() { @Override public void close() { + LOG.info("Closing RequestProcessor"); shutdownExecutor(writeThreadPool); shutdownExecutor(readThreadPool); if (serverCfg.getNumLongPollWorkerThreads() > 0 || readThreadPool == null) { @@ -271,6 +272,7 @@ public void close() { } shutdownExecutor(highPriorityThreadPool); requestTimer.stop(); + LOG.info("Closed RequestProcessor"); } private OrderedExecutor createExecutor( @@ -295,6 +297,7 @@ private OrderedExecutor createExecutor( private void shutdownExecutor(OrderedExecutor service) { if (null != service) { service.shutdown(); + service.forceShutdown(10, 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..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 @@ -220,8 +220,8 @@ public synchronized void shutdown() { if (!running) { return; } - exitCode = bookie.shutdown(); this.requestProcessor.close(); + exitCode = bookie.shutdown(); running = false; } 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..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 @@ -66,7 +66,12 @@ 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()); + } else { + 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); } else { 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..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 @@ -87,25 +87,29 @@ 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); + } } - } - }); + }); + } else { + LOGGER.debug("Netty channel {} is inactive, " + + "hence bypassing netty channel writeAndFlush during sendResponse", channel); + } } protected boolean isVersionCompatible() { 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,