From 1743ff4879406ee7b9136de9aa04a95d2ccdef36 Mon Sep 17 00:00:00 2001 From: Mattison Date: Tue, 28 May 2024 17:20:33 +0800 Subject: [PATCH 1/3] [improve][broker] avoid creating a new object when intercept --- .../BrokerInterceptorWithClassLoader.java | 122 +++++++++++++++--- 1 file changed, 101 insertions(+), 21 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java index faee5799289d0..dfc06613de2db 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java @@ -58,9 +58,13 @@ public void beforeSendMessage(Subscription subscription, Entry entry, long[] ackSet, MessageMetadata msgMetadata) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.beforeSendMessage( subscription, entry, ackSet, msgMetadata); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @@ -70,25 +74,37 @@ public void beforeSendMessage(Subscription subscription, long[] ackSet, MessageMetadata msgMetadata, Consumer consumer) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.beforeSendMessage( subscription, entry, ackSet, msgMetadata, consumer); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void onMessagePublish(Producer producer, ByteBuf headersAndPayload, Topic.PublishContext publishContext) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.onMessagePublish(producer, headersAndPayload, publishContext); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void producerCreated(ServerCnx cnx, Producer producer, Map metadata){ - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.producerCreated(cnx, producer, metadata); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @@ -96,8 +112,12 @@ public void producerCreated(ServerCnx cnx, Producer producer, public void producerClosed(ServerCnx cnx, Producer producer, Map metadata) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.producerClosed(cnx, producer, metadata); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @@ -105,9 +125,12 @@ public void producerClosed(ServerCnx cnx, public void consumerCreated(ServerCnx cnx, Consumer consumer, Map metadata) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { - this.interceptor.consumerCreated( - cnx, consumer, metadata); + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); + this.interceptor.consumerCreated( cnx, consumer, metadata); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @@ -115,8 +138,12 @@ public void consumerCreated(ServerCnx cnx, public void consumerClosed(ServerCnx cnx, Consumer consumer, Map metadata) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.consumerClosed(cnx, consumer, metadata); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @@ -124,85 +151,138 @@ public void consumerClosed(ServerCnx cnx, @Override public void messageProduced(ServerCnx cnx, Producer producer, long startTimeNs, long ledgerId, long entryId, Topic.PublishContext publishContext) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.messageProduced(cnx, producer, startTimeNs, ledgerId, entryId, publishContext); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void messageDispatched(ServerCnx cnx, Consumer consumer, long ledgerId, long entryId, ByteBuf headersAndPayload) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.messageDispatched(cnx, consumer, ledgerId, entryId, headersAndPayload); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void messageAcked(ServerCnx cnx, Consumer consumer, CommandAck ackCmd) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.messageAcked(cnx, consumer, ackCmd); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void txnOpened(long tcId, String txnID) { - this.interceptor.txnOpened(tcId, txnID); + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); + this.interceptor.txnOpened(tcId, txnID); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); + } } @Override public void txnEnded(String txnID, long txnAction) { - this.interceptor.txnEnded(txnID, txnAction); + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); + this.interceptor.txnEnded(txnID, txnAction); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); + } } @Override public void onConnectionCreated(ServerCnx cnx) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.onConnectionCreated(cnx); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void onPulsarCommand(BaseCommand command, ServerCnx cnx) throws InterceptException { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.onPulsarCommand(command, cnx); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void onConnectionClosed(ServerCnx cnx) { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.onConnectionClosed(cnx); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void onWebserviceRequest(ServletRequest request) throws IOException, ServletException, InterceptException { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.onWebserviceRequest(request); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void onWebserviceResponse(ServletRequest request, ServletResponse response) throws IOException, ServletException { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.onWebserviceResponse(request, response); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void initialize(PulsarService pulsarService) throws Exception { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); this.interceptor.initialize(pulsarService); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } } @Override public void close() { - try (ClassLoaderSwitcher ignored = new ClassLoaderSwitcher(classLoader)) { + final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); + try { + Thread.currentThread().setContextClassLoader(classLoader); interceptor.close(); + } finally { + Thread.currentThread().setContextClassLoader(previousContext); } + try { classLoader.close(); } catch (IOException e) { From 9b3750be7704a5dfbbed475b018b802af529ba42 Mon Sep 17 00:00:00 2001 From: Mattison Date: Tue, 28 May 2024 17:26:44 +0800 Subject: [PATCH 2/3] rename class loader to nar class loader --- .../BrokerInterceptorWithClassLoader.java | 43 +++++++++---------- .../intercept/BrokerInterceptorUtilsTest.java | 2 +- .../BrokerInterceptorWithClassLoaderTest.java | 2 +- 3 files changed, 23 insertions(+), 24 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java index dfc06613de2db..68553e7430a3d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java @@ -29,7 +29,6 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.Entry; -import org.apache.pulsar.broker.ClassLoaderSwitcher; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.service.Consumer; import org.apache.pulsar.broker.service.Producer; @@ -51,7 +50,7 @@ public class BrokerInterceptorWithClassLoader implements BrokerInterceptor { private final BrokerInterceptor interceptor; - private final NarClassLoader classLoader; + private final NarClassLoader narClassLoader; @Override public void beforeSendMessage(Subscription subscription, @@ -60,7 +59,7 @@ public void beforeSendMessage(Subscription subscription, MessageMetadata msgMetadata) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.beforeSendMessage( subscription, entry, ackSet, msgMetadata); } finally { @@ -76,7 +75,7 @@ public void beforeSendMessage(Subscription subscription, Consumer consumer) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.beforeSendMessage( subscription, entry, ackSet, msgMetadata, consumer); } finally { @@ -89,7 +88,7 @@ public void onMessagePublish(Producer producer, ByteBuf headersAndPayload, Topic.PublishContext publishContext) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.onMessagePublish(producer, headersAndPayload, publishContext); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -101,7 +100,7 @@ public void producerCreated(ServerCnx cnx, Producer producer, Map metadata){ final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.producerCreated(cnx, producer, metadata); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -114,7 +113,7 @@ public void producerClosed(ServerCnx cnx, Map metadata) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.producerClosed(cnx, producer, metadata); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -127,7 +126,7 @@ public void consumerCreated(ServerCnx cnx, Map metadata) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.consumerCreated( cnx, consumer, metadata); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -140,7 +139,7 @@ public void consumerClosed(ServerCnx cnx, Map metadata) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.consumerClosed(cnx, consumer, metadata); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -153,7 +152,7 @@ public void messageProduced(ServerCnx cnx, Producer producer, long startTimeNs, long entryId, Topic.PublishContext publishContext) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.messageProduced(cnx, producer, startTimeNs, ledgerId, entryId, publishContext); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -165,7 +164,7 @@ public void messageDispatched(ServerCnx cnx, Consumer consumer, long ledgerId, long entryId, ByteBuf headersAndPayload) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.messageDispatched(cnx, consumer, ledgerId, entryId, headersAndPayload); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -177,7 +176,7 @@ public void messageAcked(ServerCnx cnx, Consumer consumer, CommandAck ackCmd) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.messageAcked(cnx, consumer, ackCmd); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -188,7 +187,7 @@ public void messageAcked(ServerCnx cnx, Consumer consumer, public void txnOpened(long tcId, String txnID) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.txnOpened(tcId, txnID); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -199,7 +198,7 @@ public void txnOpened(long tcId, String txnID) { public void txnEnded(String txnID, long txnAction) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.txnEnded(txnID, txnAction); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -210,7 +209,7 @@ public void txnEnded(String txnID, long txnAction) { public void onConnectionCreated(ServerCnx cnx) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.onConnectionCreated(cnx); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -221,7 +220,7 @@ public void onConnectionCreated(ServerCnx cnx) { public void onPulsarCommand(BaseCommand command, ServerCnx cnx) throws InterceptException { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.onPulsarCommand(command, cnx); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -232,7 +231,7 @@ public void onPulsarCommand(BaseCommand command, ServerCnx cnx) throws Intercept public void onConnectionClosed(ServerCnx cnx) { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.onConnectionClosed(cnx); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -243,7 +242,7 @@ public void onConnectionClosed(ServerCnx cnx) { public void onWebserviceRequest(ServletRequest request) throws IOException, ServletException, InterceptException { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.onWebserviceRequest(request); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -255,7 +254,7 @@ public void onWebserviceResponse(ServletRequest request, ServletResponse respons throws IOException, ServletException { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.onWebserviceResponse(request, response); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -266,7 +265,7 @@ public void onWebserviceResponse(ServletRequest request, ServletResponse respons public void initialize(PulsarService pulsarService) throws Exception { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); this.interceptor.initialize(pulsarService); } finally { Thread.currentThread().setContextClassLoader(previousContext); @@ -277,14 +276,14 @@ public void initialize(PulsarService pulsarService) throws Exception { public void close() { final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { - Thread.currentThread().setContextClassLoader(classLoader); + Thread.currentThread().setContextClassLoader(narClassLoader); interceptor.close(); } finally { Thread.currentThread().setContextClassLoader(previousContext); } try { - classLoader.close(); + narClassLoader.close(); } catch (IOException e) { log.warn("Failed to close the broker interceptor class loader", e); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorUtilsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorUtilsTest.java index 5abe8a69ee499..979bf6cd0d5db 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorUtilsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorUtilsTest.java @@ -65,7 +65,7 @@ public void testLoadBrokerEventListener() throws Exception { BrokerInterceptorWithClassLoader returnedPhWithCL = BrokerInterceptorUtils.load(metadata, ""); BrokerInterceptor returnedPh = returnedPhWithCL.getInterceptor(); - assertSame(mockLoader, returnedPhWithCL.getClassLoader()); + assertSame(mockLoader, returnedPhWithCL.getNarClassLoader()); assertTrue(returnedPh instanceof MockBrokerInterceptor); } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoaderTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoaderTest.java index a2f97e16a76ae..64d4b5ee6cca5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoaderTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoaderTest.java @@ -135,7 +135,7 @@ public void close() { new BrokerInterceptorWithClassLoader(interceptor, narLoader); ClassLoader curClassLoader = Thread.currentThread().getContextClassLoader(); // test class loader - assertEquals(brokerInterceptorWithClassLoader.getClassLoader(), narLoader); + assertEquals(brokerInterceptorWithClassLoader.getNarClassLoader(), narLoader); // test initialize brokerInterceptorWithClassLoader.initialize(mock(PulsarService.class)); assertEquals(Thread.currentThread().getContextClassLoader(), curClassLoader); From 22b65d5e61912ce044c86ce5c462b710ff58759c Mon Sep 17 00:00:00 2001 From: Mattison Date: Tue, 28 May 2024 18:04:50 +0800 Subject: [PATCH 3/3] fix checkstyle --- .../broker/intercept/BrokerInterceptorWithClassLoader.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java index 68553e7430a3d..3997e214f4316 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/intercept/BrokerInterceptorWithClassLoader.java @@ -127,7 +127,7 @@ public void consumerCreated(ServerCnx cnx, final ClassLoader previousContext = Thread.currentThread().getContextClassLoader(); try { Thread.currentThread().setContextClassLoader(narClassLoader); - this.interceptor.consumerCreated( cnx, consumer, metadata); + this.interceptor.consumerCreated(cnx, consumer, metadata); } finally { Thread.currentThread().setContextClassLoader(previousContext); }