{"record":{"id":"1fa1b4af42bd1f9d","repo":"apache/pulsar","slug":"close-interrupted-1fa1b4","errorCode":null,"errorMessage":"Close interrupted","messagePattern":"Close interrupted","errorType":"exception","errorClass":"PulsarClientException","httpStatus":null,"severity":"warning","filePath":"pulsar-client-v5/src/main/java/org/apache/pulsar/client/impl/v5/ScalableStreamConsumer.java","lineNumber":386,"sourceCode":"        if (!drainingConsumers.containsKey(segmentId)) {\n            return;\n        }\n        pendingDrainAcks.computeIfAbsent(segmentId, __ -> new ConcurrentLinkedQueue<>())\n                .add(ackFuture.exceptionally(ex -> null));\n    }\n\n    @Override\n    public AsyncStreamConsumer<T> async() {\n        return asyncView;\n    }\n\n    @Override\n    public void close() throws PulsarClientException {\n        try {\n            closeAsync().get();\n        } catch (InterruptedException e) {\n            Thread.currentThread().interrupt();\n            throw new PulsarClientException(\"Close interrupted\", e);\n        } catch (ExecutionException e) {\n            throw new PulsarClientException(e.getCause());\n        }\n    }\n\n    // --- Async internals ---\n\n    CompletableFuture<Message<T>> receiveAsync() {\n        return receiveQueue.receiveAsync();\n    }\n\n    CompletableFuture<Message<T>> receiveAsync(Duration timeout) {\n        return receiveQueue.receiveAsync(timeout);\n    }\n\n    CompletableFuture<List<Message<T>>> receiveMultiAsync(int maxNumMessages, Duration timeout) {\n        return receiveQueue.receiveMultiAsync(maxNumMessages, timeout);\n    }","sourceCodeStart":368,"sourceCodeEnd":404,"githubUrl":"https://github.com/apache/pulsar/blob/820761864ed8e2a7d2e52dd9763ad2ae117c1395/pulsar-client-v5/src/main/java/org/apache/pulsar/client/impl/v5/ScalableStreamConsumer.java#L368-L404","documentation":"ScalableStreamConsumer.close() blocks on closeAsync().get(); if the thread is interrupted while waiting, the interrupt flag is restored and PulsarClientException(\"Close interrupted\") is thrown. The underlying close of all stream segments may still continue in the background.","triggerScenarios":"Interruption of the thread blocked in closeAsync().get(): ExecutorService.shutdownNow() during app shutdown, future.cancel(true), or timeout frameworks that interrupt threads.","commonSituations":"SIGTERM handlers interrupting shutdown threads before the multi-segment close finishes; shutdown timeouts that call shutdownNow(); interrupt-based cancellation in worker loops.","solutions":["Close consumers before shutting down executors / interrupting threads","Use closeAsync().get(timeout, TimeUnit) instead of interruption for deadlines","Catch the exception and verify Thread.interrupted() to confirm this scenario","If interrupted, call closeAsync() again from a non-interrupted thread to confirm completion"],"exampleFix":"// before\n// in shutdown hook, with pool already shutdownNow()\nstreamConsumer.close();\n// after\nstreamConsumer.closeAsync().get(30, TimeUnit.SECONDS);\nexecutor.shutdownNow(); // interrupt only after close finished","handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"try {\n    streamConsumer.close();\n} catch (PulsarClientException e) {\n    if (\"Close interrupted\".equals(e.getMessage())) {\n        log.warn(\"stream close interrupted; segments may still be closing\");\n        streamConsumer.closeAsync(); // finish in background\n    }\n}","preventionTips":["Complete consumer close before interrupting shutdown threads","Prefer closeAsync().get(timeout) over thread interruption","Avoid shutdownNow() while multi-segment close is in flight"],"tags":["pulsar","close","interrupt","shutdown"],"backgroundTag":"close-interrupted","analyzedSha":"820761864ed8e2a7d2e52dd9763ad2ae117c1395","analyzedAt":"2026-09-06T00:14:20.138Z","contentChangedAt":"2026-09-06T00:14:20.138Z","schemaVersion":2},"datasetVersion":"2026-09-14T00:17:10.932Z"}