From 57fc85a08d1735c7dabbdf7228dee5e53253f70c Mon Sep 17 00:00:00 2001 From: Ali Ahmed Date: Fri, 31 Jul 2026 02:30:45 -0700 Subject: [PATCH] [cleanup][misc] Fix all javac warnings in main sources MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `compileJava` emitted 77 warnings across ten modules, mostly deprecations left by the Netty 4.2, Jetty 12.1, Guava 33 and commons-lang3 3.20 upgrades — including Jetty's `onWebSocketClose(int, String)`, which is marked for removal. Moved to the successor APIs: Guava's `Duration` cache overloads, `RECVBUF_ALLOCATOR`, `ThreadLocalRandom.current()`, `MutableObject.get()`, the 4-arg `PersistentSubscription` constructor and Jetty's 3-arg `onWebSocketClose`. Dropped the `EPOLL_MODE` call that Netty 4.2 documents as a no-op. Where no successor exists, suppressed narrowly at the declaration with a reason. Added `compileOnly(swagger-annotations)` to the three modules that read `pulsar-client`'s `@Schema` annotations without the types on their compile classpath. `compileJava` is now warning-free and `sanityCheck` passes. Test sources still emit warnings; those are left for a separate change. Assisted-by: Claude Code (Opus 5) Co-Authored-By: Claude Opus 5 --- .../broker/loadbalance/impl/SimpleLoadManagerImpl.java | 3 ++- .../org/apache/pulsar/broker/service/BrokerService.java | 4 ++-- .../broker/service/LegacyAwareTopicPoliciesService.java | 1 + .../broker/service/MetadataStoreTopicPoliciesService.java | 1 + .../service/SystemTopicBasedTopicPoliciesService.java | 3 ++- .../pulsar/broker/service/persistent/PersistentTopic.java | 5 +++-- .../broker/service/scalable/ScalableTopicController.java | 6 +++--- pulsar-client-admin/build.gradle.kts | 4 ++++ .../java/org/apache/pulsar/client/api/ConsumerBuilder.java | 2 ++ .../java/org/apache/pulsar/client/api/ProducerBuilder.java | 2 ++ .../java/org/apache/pulsar/client/api/ReaderBuilder.java | 2 ++ pulsar-client-v5/build.gradle.kts | 3 +++ .../pulsar/client/impl/v5/AuthenticationAdapter.java | 3 +++ .../org/apache/pulsar/client/impl/PulsarClientImpl.java | 2 +- .../pulsar/client/impl/PulsarServiceNameResolver.java | 3 ++- .../schema/generic/MultiVersionSchemaInfoProvider.java | 4 ++-- .../impl/schema/reader/AbstractMultiVersionReader.java | 4 ++-- .../org/apache/pulsar/common/naming/NamespaceName.java | 4 ++-- .../packages/management/core/common/PackageName.java | 4 ++-- .../org/apache/pulsar/proxy/server/DirectProxyHandler.java | 7 ++----- .../java/org/apache/pulsar/proxy/server/ProxyService.java | 4 ++-- pulsar-testclient/build.gradle.kts | 4 ++++ .../apache/pulsar/websocket/AbstractWebSocketHandler.java | 3 ++- .../java/org/apache/pulsar/websocket/ConsumerHandler.java | 3 ++- .../pulsar/websocket/AbstractWebSocketHandlerTest.java | 4 +++- 25 files changed, 56 insertions(+), 29 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/SimpleLoadManagerImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/SimpleLoadManagerImpl.java index b1f68e07f2fb2..f4314ba4f7a08 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/SimpleLoadManagerImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/SimpleLoadManagerImpl.java @@ -26,6 +26,7 @@ import com.google.common.collect.Multimap; import com.google.common.collect.Sets; import com.google.common.collect.TreeMultimap; +import java.time.Duration; import java.time.LocalDateTime; import java.util.ArrayList; import java.util.Collection; @@ -238,7 +239,7 @@ public void initialize(final PulsarService pulsar) { }); int entryExpiryTime = (int) pulsar.getConfiguration().getLoadBalancerSheddingGracePeriodMinutes(); - unloadedHotNamespaceCache = CacheBuilder.newBuilder().expireAfterWrite(entryExpiryTime, TimeUnit.MINUTES) + unloadedHotNamespaceCache = CacheBuilder.newBuilder().expireAfterWrite(Duration.ofMinutes(entryExpiryTime)) .build(new CacheLoader() { @Override public Long load(String key) throws Exception { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 496667ef0e3f9..78fa0bce8fd64 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -578,7 +578,7 @@ private void startProtocolHandler(String protocol, bootstrap.option(ChannelOption.SO_REUSEADDR, true); bootstrap.childOption(ChannelOption.ALLOCATOR, PulsarByteBufAllocator.DEFAULT); bootstrap.childOption(ChannelOption.TCP_NODELAY, true); - bootstrap.childOption(ChannelOption.RCVBUF_ALLOCATOR, + bootstrap.childOption(ChannelOption.RECVBUF_ALLOCATOR, new AdaptiveRecvByteBufAllocator(1024, 16 * 1024, 1 * 1024 * 1024)); EventLoopUtil.enableTriggeredMode(bootstrap); DefaultThreadFactory defaultThreadFactory = @@ -606,7 +606,7 @@ private ServerBootstrap defaultServerBootstrap() { bootstrap.childOption(ChannelOption.ALLOCATOR, PulsarByteBufAllocator.DEFAULT); bootstrap.group(acceptorGroup, workerGroup); bootstrap.childOption(ChannelOption.TCP_NODELAY, true); - bootstrap.childOption(ChannelOption.RCVBUF_ALLOCATOR, + bootstrap.childOption(ChannelOption.RECVBUF_ALLOCATOR, new AdaptiveRecvByteBufAllocator(1024, 16 * 1024, 1 * 1024 * 1024)); bootstrap.channel(EventLoopUtil.getServerSocketChannelClass(workerGroup)); EventLoopUtil.enableTriggeredMode(bootstrap); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/LegacyAwareTopicPoliciesService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/LegacyAwareTopicPoliciesService.java index 20f7b20799128..9e56ec9a4af87 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/LegacyAwareTopicPoliciesService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/LegacyAwareTopicPoliciesService.java @@ -125,6 +125,7 @@ public CompletableFuture registerListenerAsync(TopicName topicName, Top .thenCompose(service -> service.registerListenerAsync(topicName, listener)); } + @Deprecated @Override public boolean registerListener(TopicName topicName, TopicPolicyListener listener) { throw new RuntimeException("should not be called"); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/MetadataStoreTopicPoliciesService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/MetadataStoreTopicPoliciesService.java index 56319a44ac330..91cf100cb9939 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/MetadataStoreTopicPoliciesService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/MetadataStoreTopicPoliciesService.java @@ -135,6 +135,7 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam .thenApply(policies -> policies.map(policy -> cloneWithScope(policy, global))); } + @Deprecated @Override public boolean registerListener(TopicName topicName, TopicPolicyListener listener) { listeners.compute(normalizeTopicName(topicName), (__, topicListeners) -> { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java index 8b4cae18bc483..4f7b6e3b37a3f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java @@ -951,7 +951,7 @@ void cleanPoliciesCacheInitMap(@NonNull NamespaceName namespace) { // Complete the removed init future outside the compute() remapping function: completing it can run the // awaiting topic-load callbacks synchronously, and doing that while holding the ConcurrentHashMap bin lock // risks a recursive map update / deadlock (see #24977). - failPendingPolicyCacheInit(namespace, removedInitFuture.getValue()); + failPendingPolicyCacheInit(namespace, removedInitFuture.get()); if (readerFuture != null && !readerFuture.isCompletedExceptionally()) { readerFuture .thenCompose(SystemTopicClient.Reader::closeAsync) @@ -1152,6 +1152,7 @@ public CompletableFuture getPoliciesCacheInit(NamespaceName namespaceName) return policyCacheInitMap.get(namespaceName); } + @Deprecated @Override public boolean registerListener(TopicName topicName, TopicPolicyListener listener) { listeners.compute(topicName, (k, topicListeners) -> { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 3c4a404c32c12..bc8c58f2d1ea8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -625,7 +625,8 @@ public CompletableFuture unloadSubscription(@NonNull String subName) { * @param subscriptionName the name of the subscription * @param cursor the cursor to use for the subscription * @param replicated the subscription replication flag - * @param subscriptionProperties the subscription properties + * @param subscriptionProperties the subscription properties. No longer used to build the subscription — the cursor + * already carries them — but retained so overriding implementations keep working. * @return the subscription instance */ protected PersistentSubscription createPersistentSubscription(String subscriptionName, ManagedCursor cursor, @@ -636,7 +637,7 @@ protected PersistentSubscription createPersistentSubscription(String subscriptio CompactedTopicImpl compactedTopic = pulsarTopicCompactionService.getCompactedTopic(); return new PulsarCompactorSubscription(this, compactedTopic, subscriptionName, cursor); } else { - return new PersistentSubscription(this, subscriptionName, cursor, replicated, subscriptionProperties); + return new PersistentSubscription(this, subscriptionName, cursor, replicated); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/scalable/ScalableTopicController.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/scalable/ScalableTopicController.java index aebeca97c6ab8..a392ec6bea0d6 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/scalable/ScalableTopicController.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/scalable/ScalableTopicController.java @@ -1264,10 +1264,10 @@ private CompletableFuture prunable(SegmentInfo seg, List subs) // No subscribers ever attached / all unsubscribed → nothing left to drain. return CompletableFuture.completedFuture(true); } - CompletableFuture[] checks = subs.stream() + List> checks = subs.stream() .map(sub -> isSegmentDrained(seg, sub)) - .toArray(CompletableFuture[]::new); - return CompletableFuture.allOf(checks) + .toList(); + return CompletableFuture.allOf(checks.toArray(CompletableFuture[]::new)) .thenApply(__ -> { for (CompletableFuture c : checks) { if (!c.join()) { diff --git a/pulsar-client-admin/build.gradle.kts b/pulsar-client-admin/build.gradle.kts index b82bdf95708e9..efb621f0ec453 100644 --- a/pulsar-client-admin/build.gradle.kts +++ b/pulsar-client-admin/build.gradle.kts @@ -43,5 +43,9 @@ dependencies { implementation(libs.commons.lang3) implementation(libs.completable.futures) + // pulsar-client's configuration-data classes carry runtime-retained @Schema annotations, so javac needs + // the OpenAPI annotation types on the compile classpath to read their class files without warning. + compileOnly(libs.swagger.annotations) + testImplementation(libs.wiremock) } diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java index d4c0b0c939750..a6c2c88cfcaa3 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java @@ -637,6 +637,8 @@ public interface ConsumerBuilder extends Cloneable { * * @param interceptors the list of interceptors to intercept the consumer created by this builder. */ + // @SafeVarargs cannot be applied to an abstract method; implementations declare it on their final override. + @SuppressWarnings("unchecked") ConsumerBuilder intercept(ConsumerInterceptor ...interceptors); /** diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java index 7b35432da1d6a..d8240edae9b85 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java @@ -561,6 +561,8 @@ public interface ProducerBuilder extends Cloneable { * @return the producer builder instance */ @Deprecated + // @SafeVarargs cannot be applied to an abstract method; implementations declare it on their final override. + @SuppressWarnings("unchecked") ProducerBuilder intercept(ProducerInterceptor ... interceptors); /** diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ReaderBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ReaderBuilder.java index c87d91744e8ff..7203eab7617bd 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ReaderBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ReaderBuilder.java @@ -342,6 +342,8 @@ public interface ReaderBuilder extends Cloneable { * @param interceptors the list of interceptors to intercept the reader created by this builder. * @return the reader builder instance */ + // @SafeVarargs cannot be applied to an abstract method; implementations declare it on their final override. + @SuppressWarnings("unchecked") ReaderBuilder intercept(ReaderInterceptor... interceptors); /** diff --git a/pulsar-client-v5/build.gradle.kts b/pulsar-client-v5/build.gradle.kts index 8e564ca4593f1..08daebb51a8cc 100644 --- a/pulsar-client-v5/build.gradle.kts +++ b/pulsar-client-v5/build.gradle.kts @@ -35,6 +35,9 @@ dependencies { compileOnly(libs.protobuf.java) implementation(libs.netty.handler) implementation(libs.jackson.annotations) + // pulsar-client's configuration-data classes carry runtime-retained @Schema annotations, so javac needs + // the OpenAPI annotation types on the compile classpath to read their class files without warning. + compileOnly(libs.swagger.annotations) compileOnly(libs.lombok) annotationProcessor(libs.lombok) testImplementation(libs.testng) diff --git a/pulsar-client-v5/src/main/java/org/apache/pulsar/client/impl/v5/AuthenticationAdapter.java b/pulsar-client-v5/src/main/java/org/apache/pulsar/client/impl/v5/AuthenticationAdapter.java index bb51d2f9107e6..762b7175def74 100644 --- a/pulsar-client-v5/src/main/java/org/apache/pulsar/client/impl/v5/AuthenticationAdapter.java +++ b/pulsar-client-v5/src/main/java/org/apache/pulsar/client/impl/v5/AuthenticationAdapter.java @@ -110,6 +110,9 @@ public String authMethodName() { } @Override + // The v4 no-arg getAuthData() is deprecated in favour of the broker-host variant, but this bridge has to + // preserve the host-less semantics of the v5 no-arg authData() it implements. + @SuppressWarnings("deprecation") public AuthenticationData authData() throws PulsarClientException { try { AuthenticationDataProvider v4Data = v4Auth.getAuthData(); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java index 9d0a814c5eb1e..bc17f3228d9a5 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java @@ -154,7 +154,7 @@ public enum State { private final LoadingCache schemaProviderLoadingCache = CacheBuilder.newBuilder().maximumSize(100000) - .expireAfterAccess(30, TimeUnit.MINUTES) + .expireAfterAccess(Duration.ofMinutes(30)) .build(new CacheLoader() { @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarServiceNameResolver.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarServiceNameResolver.java index 0d5d977c35721..d056ce98b56de 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarServiceNameResolver.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarServiceNameResolver.java @@ -29,6 +29,7 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.stream.Collectors; @@ -155,7 +156,7 @@ public synchronized void updateServiceUrl(String serviceUrl) throws InvalidServi private static int randomIndex(int numAddresses) { return numAddresses == 1 ? - 0 : io.netty.util.internal.PlatformDependent.threadLocalRandom().nextInt(numAddresses); + 0 : ThreadLocalRandom.current().nextInt(numAddresses); } /** diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/MultiVersionSchemaInfoProvider.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/MultiVersionSchemaInfoProvider.java index 4c29c420b3ed9..36628722a955e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/MultiVersionSchemaInfoProvider.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/generic/MultiVersionSchemaInfoProvider.java @@ -22,9 +22,9 @@ import com.google.common.cache.CacheLoader; import com.google.common.cache.LoadingCache; import java.nio.charset.StandardCharsets; +import java.time.Duration; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; import lombok.CustomLog; import org.apache.pulsar.client.api.schema.SchemaInfoProvider; import org.apache.pulsar.client.impl.PulsarClientImpl; @@ -44,7 +44,7 @@ public class MultiVersionSchemaInfoProvider implements SchemaInfoProvider { private final LoadingCache> cache = CacheBuilder.newBuilder() .maximumSize(100000) - .expireAfterAccess(30, TimeUnit.MINUTES) + .expireAfterAccess(Duration.ofMinutes(30)) .build(new CacheLoader>() { @Override public CompletableFuture load(BytesSchemaVersion schemaVersion) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/reader/AbstractMultiVersionReader.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/reader/AbstractMultiVersionReader.java index de971b1343e43..fd48b0a8f0860 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/reader/AbstractMultiVersionReader.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/schema/reader/AbstractMultiVersionReader.java @@ -22,8 +22,8 @@ import com.google.common.cache.CacheLoader; import com.google.common.cache.LoadingCache; import java.io.InputStream; +import java.time.Duration; import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; import lombok.CustomLog; import org.apache.avro.AvroTypeException; import org.apache.commons.codec.binary.Hex; @@ -45,7 +45,7 @@ public abstract class AbstractMultiVersionReader implements SchemaReader { protected SchemaInfoProvider schemaInfoProvider; LoadingCache> readerCache = CacheBuilder.newBuilder().maximumSize(100000) - .expireAfterAccess(30, TimeUnit.MINUTES).build(new CacheLoader>() { + .expireAfterAccess(Duration.ofMinutes(30)).build(new CacheLoader>() { @Override public SchemaReader load(BytesSchemaVersion schemaVersion) { return loadReader(schemaVersion); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/naming/NamespaceName.java b/pulsar-common/src/main/java/org/apache/pulsar/common/naming/NamespaceName.java index bda21502dfbee..d3058e0b455ea 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/naming/NamespaceName.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/naming/NamespaceName.java @@ -23,10 +23,10 @@ import com.google.common.cache.LoadingCache; import com.google.common.util.concurrent.UncheckedExecutionException; import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; +import java.time.Duration; import java.util.Objects; import java.util.Optional; import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; /** * Parser of a value from the namespace field provided in configuration. @@ -39,7 +39,7 @@ public class NamespaceName implements ServiceUnitId { private final String localName; private static final LoadingCache cache = CacheBuilder.newBuilder().maximumSize(100000) - .expireAfterAccess(30, TimeUnit.MINUTES).build(new CacheLoader() { + .expireAfterAccess(Duration.ofMinutes(30)).build(new CacheLoader() { @Override public NamespaceName load(String name) throws Exception { return new NamespaceName(name); diff --git a/pulsar-package-management/core/src/main/java/org/apache/pulsar/packages/management/core/common/PackageName.java b/pulsar-package-management/core/src/main/java/org/apache/pulsar/packages/management/core/common/PackageName.java index 668bb71654884..7fa215d91d0a1 100644 --- a/pulsar-package-management/core/src/main/java/org/apache/pulsar/packages/management/core/common/PackageName.java +++ b/pulsar-package-management/core/src/main/java/org/apache/pulsar/packages/management/core/common/PackageName.java @@ -25,10 +25,10 @@ import com.google.common.cache.CacheLoader; import com.google.common.cache.LoadingCache; import com.google.common.net.UrlEscapers; +import java.time.Duration; import java.util.List; import java.util.Objects; import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; /** * A package name has five parts, type, tenant, namespace, package-name, and version. @@ -49,7 +49,7 @@ public class PackageName { private static final LoadingCache cache = CacheBuilder.newBuilder() .maximumSize(100000) - .expireAfterAccess(30, TimeUnit.MINUTES) + .expireAfterAccess(Duration.ofMinutes(30)) .build(new CacheLoader() { @Override public PackageName load(String name) throws Exception { diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index b748a7aa4f77f..3d6d428848d62 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -30,8 +30,6 @@ import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInitializer; import io.netty.channel.ChannelOption; -import io.netty.channel.epoll.EpollChannelOption; -import io.netty.channel.epoll.EpollMode; import io.netty.channel.epoll.EpollSocketChannel; import io.netty.channel.socket.SocketChannel; import io.netty.handler.codec.haproxy.HAProxyCommand; @@ -160,9 +158,8 @@ public void connect(String brokerHostAndPort, InetSocketAddress targetBrokerAddr b.group(inboundChannel.eventLoop()) .channel(inboundChannel.getClass()); - if (service.proxyZeroCopyModeEnabled && EpollSocketChannel.class.isAssignableFrom(inboundChannel.getClass())) { - b.option(EpollChannelOption.EPOLL_MODE, EpollMode.LEVEL_TRIGGERED); - } + // Zero-copy (splice) requires the epoll transport to be level-triggered. Netty always uses + // level-triggered mode since 4.2, so no channel option needs to be set for it anymore. b.handler(new ChannelInitializer() { @Override diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyService.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyService.java index c62246708ed43..a3721810e6edc 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyService.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyService.java @@ -255,7 +255,7 @@ public void start() throws Exception { bootstrap.childOption(ChannelOption.ALLOCATOR, PulsarByteBufAllocator.DEFAULT); bootstrap.group(acceptorGroup, workerGroup); bootstrap.childOption(ChannelOption.TCP_NODELAY, true); - bootstrap.childOption(ChannelOption.RCVBUF_ALLOCATOR, + bootstrap.childOption(ChannelOption.RECVBUF_ALLOCATOR, new AdaptiveRecvByteBufAllocator(1024, 16 * 1024, 1 * 1024 * 1024)); Class serverSocketChannelClass = @@ -361,7 +361,7 @@ private void startProxyExtension(String extensionName, bootstrap.option(ChannelOption.SO_REUSEADDR, true); bootstrap.childOption(ChannelOption.ALLOCATOR, PulsarByteBufAllocator.DEFAULT); bootstrap.childOption(ChannelOption.TCP_NODELAY, true); - bootstrap.childOption(ChannelOption.RCVBUF_ALLOCATOR, + bootstrap.childOption(ChannelOption.RECVBUF_ALLOCATOR, new AdaptiveRecvByteBufAllocator(1024, 16 * 1024, 1 * 1024 * 1024)); EventLoopUtil.enableTriggeredMode(bootstrap); diff --git a/pulsar-testclient/build.gradle.kts b/pulsar-testclient/build.gradle.kts index 9fd025af6f1df..c1b29ea98d29f 100644 --- a/pulsar-testclient/build.gradle.kts +++ b/pulsar-testclient/build.gradle.kts @@ -54,6 +54,10 @@ dependencies { implementation(libs.jetty.util) implementation(libs.opentelemetry.sdk.extension.autoconfigure) + // pulsar-client's configuration-data classes carry runtime-retained @Schema annotations, so javac needs + // the OpenAPI annotation types on the compile classpath to read their class files without warning. + compileOnly(libs.swagger.annotations) + testImplementation(project(":pulsar-broker")) testImplementation(project(path = ":pulsar-broker", configuration = "testJar")) } diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java index 930cb4f571a1d..e9df837b62ae7 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java @@ -262,7 +262,7 @@ public void onWebSocketError(Throwable cause) { } @Override - public void onWebSocketClose(int statusCode, String reason) { + public void onWebSocketClose(int statusCode, String reason, Callback callback) { log.info() .attr("remoteAddress", getSession().getRemoteSocketAddress()) .attr("topic", topic) @@ -279,6 +279,7 @@ public void onWebSocketClose(int statusCode, String reason) { .exception(e) .log("Failed to close handler for topic"); } + callback.succeed(); } public void close(WebSocketError error) { diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/ConsumerHandler.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/ConsumerHandler.java index b0ee29e48ec5e..580d7413d4beb 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/ConsumerHandler.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/ConsumerHandler.java @@ -27,6 +27,7 @@ import com.google.common.cache.CacheBuilder; import jakarta.servlet.http.HttpServletRequest; import java.io.IOException; +import java.time.Duration; import java.util.Base64; import java.util.List; import java.util.concurrent.TimeUnit; @@ -96,7 +97,7 @@ public class ConsumerHandler extends AbstractWebSocketHandler { // Make sure use the same BatchMessageIdImpl to acknowledge the batch message, otherwise the BatchMessageAcker // of the BatchMessageIdImpl will not complete. private Cache messageIdCache = CacheBuilder.newBuilder() - .expireAfterWrite(1, TimeUnit.HOURS) + .expireAfterWrite(Duration.ofHours(1)) .build(); public ConsumerHandler(WebSocketService service, HttpServletRequest request, JettyServerUpgradeResponse response) { diff --git a/pulsar-websocket/src/test/java/org/apache/pulsar/websocket/AbstractWebSocketHandlerTest.java b/pulsar-websocket/src/test/java/org/apache/pulsar/websocket/AbstractWebSocketHandlerTest.java index aa0ed89c72cb7..0ffa10fb1474d 100644 --- a/pulsar-websocket/src/test/java/org/apache/pulsar/websocket/AbstractWebSocketHandlerTest.java +++ b/pulsar-websocket/src/test/java/org/apache/pulsar/websocket/AbstractWebSocketHandlerTest.java @@ -57,6 +57,7 @@ import org.apache.pulsar.websocket.service.WebSocketProxyConfiguration; import org.eclipse.jetty.ee10.websocket.server.JettyServerUpgradeResponse; import org.eclipse.jetty.http.HttpStatus; +import org.eclipse.jetty.websocket.api.Callback; import org.eclipse.jetty.websocket.api.Session; import org.mockito.Answers; import org.mockito.Mock; @@ -395,7 +396,8 @@ public void testPingFuture() throws IOException { assertNotNull(pingFuture); assertFalse(pingFuture.isDone()); - webSocketHandler.onWebSocketClose(HttpStatus.INTERNAL_SERVER_ERROR_500, "INTERNAL_SERVER_ERROR_500"); + webSocketHandler.onWebSocketClose(HttpStatus.INTERNAL_SERVER_ERROR_500, "INTERNAL_SERVER_ERROR_500", + Callback.NOOP); assertTrue(pingFuture.isDone());