diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java index 3933db7194..c3d6bf9ede 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java @@ -22,7 +22,6 @@ import io.netty.handler.codec.LengthFieldPrepender; import io.netty.handler.ssl.SslHandler; import io.netty.handler.timeout.IdleStateHandler; -import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory; import io.streamnative.pulsar.handlers.kop.utils.ssl.SSLUtils; @@ -59,7 +58,7 @@ public class KafkaChannelInitializer extends ChannelInitializer { @Getter private final SslContextFactory.Server sslContextFactory; @Getter - private final StatsLogger statsLogger; + private final RequestStats requestStats; private final OrderedScheduler sendResponseScheduler; public KafkaChannelInitializer(PulsarService pulsarService, @@ -72,7 +71,7 @@ public KafkaChannelInitializer(PulsarService pulsarService, boolean enableTLS, EndPoint advertisedEndPoint, boolean skipMessagesWithoutIndex, - StatsLogger statsLogger, + RequestStats requestStats, OrderedScheduler sendResponseScheduler) { super(); this.pulsarService = pulsarService; @@ -85,7 +84,7 @@ public KafkaChannelInitializer(PulsarService pulsarService, this.enableTls = enableTLS; this.advertisedEndPoint = advertisedEndPoint; this.skipMessagesWithoutIndex = skipMessagesWithoutIndex; - this.statsLogger = statsLogger; + this.requestStats = requestStats; if (enableTls) { sslContextFactory = SSLUtils.createSslContextFactory(kafkaConfig); } else { @@ -116,15 +115,15 @@ public KafkaRequestHandler newCnx() throws Exception { return new KafkaRequestHandler(pulsarService, kafkaConfig, tenantContextManager, kopBrokerLookupManager, adminManager, producePurgatory, fetchPurgatory, - enableTls, advertisedEndPoint, skipMessagesWithoutIndex, statsLogger, sendResponseScheduler); + enableTls, advertisedEndPoint, skipMessagesWithoutIndex, requestStats, sendResponseScheduler); } @VisibleForTesting - public KafkaRequestHandler newCnx(final TenantContextManager tenantContextManager, - final StatsLogger statsLogger) throws Exception { + public KafkaRequestHandler newCnxWithoutStats(final TenantContextManager tenantContextManager) throws Exception { return new KafkaRequestHandler(pulsarService, kafkaConfig, tenantContextManager, kopBrokerLookupManager, adminManager, producePurgatory, fetchPurgatory, - enableTls, advertisedEndPoint, skipMessagesWithoutIndex, statsLogger, sendResponseScheduler); + enableTls, advertisedEndPoint, skipMessagesWithoutIndex, RequestStats.NULL_INSTANCE, + sendResponseScheduler); } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaCommandDecoder.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaCommandDecoder.java index 37b2912037..9cbb158a3a 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaCommandDecoder.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaCommandDecoder.java @@ -21,7 +21,6 @@ import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInboundHandlerAdapter; import io.netty.handler.timeout.IdleStateEvent; -import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; import java.io.Closeable; import java.net.SocketAddress; import java.nio.ByteBuffer; @@ -68,10 +67,10 @@ public abstract class KafkaCommandDecoder extends ChannelInboundHandlerAdapter { private final OrderedScheduler sendResponseScheduler; - public KafkaCommandDecoder(StatsLogger statsLogger, + public KafkaCommandDecoder(RequestStats requestStats, KafkaServiceConfiguration kafkaConfig, OrderedScheduler sendResponseScheduler) { - this.requestStats = new RequestStats(statsLogger); + this.requestStats = requestStats; this.kafkaConfig = kafkaConfig; this.requestQueue = new LinkedBlockingQueue<>(kafkaConfig.getMaxQueuedRequests()); this.sendResponseScheduler = sendResponseScheduler; @@ -193,12 +192,9 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception }; // Update handle request latency metrics - final BiConsumer registerRequestLatency = (apiName, startProcessTime) -> { - requestStats.getStatsLogger() - .scopeLabel(KopServerStats.REQUEST_SCOPE, apiName) - .getOpStatsLogger(KopServerStats.REQUEST_LATENCY) - .registerSuccessfulEvent(MathUtils.elapsedNanos(startProcessTime), - TimeUnit.NANOSECONDS); + final BiConsumer registerRequestLatency = (apiKey, startProcessTime) -> { + requestStats.getRequestStatsLogger(apiKey, KopServerStats.REQUEST_LATENCY) + .registerSuccessfulEvent(MathUtils.elapsedNanos(startProcessTime), TimeUnit.NANOSECONDS); }; // If kop is enabled for authentication and the client @@ -248,7 +244,7 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception return; } - registerRequestLatency.accept(kafkaHeaderAndRequest.getHeader().apiKey().name, + registerRequestLatency.accept(kafkaHeaderAndRequest.getHeader().apiKey(), startProcessRequestTimestamp); sendResponseScheduler.executeOrdered(channel.remoteAddress().hashCode(), () -> { @@ -377,6 +373,7 @@ protected void writeAndFlushResponseToClient(Channel channel) { } final CompletableFuture responseFuture = responseAndRequest.getResponseFuture(); + final ApiKeys apiKey = responseAndRequest.getApiKey(); final long nanoSecondsSinceCreated = responseAndRequest.nanoSecondsSinceCreated(); final boolean expired = (nanoSecondsSinceCreated > TimeUnit.MILLISECONDS.toNanos(kafkaConfig.getRequestTimeoutMs())); @@ -390,7 +387,7 @@ protected void writeAndFlushResponseToClient(Channel channel) { break; } else { if (requestQueue.remove(responseAndRequest)) { - responseAndRequest.updateStats(requestStats); + RequestStats.REQUEST_QUEUE_SIZE_INSTANCE.decrementAndGet(); } else { // it has been removed by another thread, skip this element continue; } @@ -409,12 +406,8 @@ protected void writeAndFlushResponseToClient(Channel channel) { log.error("[{}] request {} completed exceptionally", channel, request.getHeader(), e); channel.writeAndFlush(request.createErrorResponse(e)); - requestStats.getStatsLogger() - .scopeLabel(KopServerStats.REQUEST_SCOPE, - responseAndRequest.request.getHeader().apiKey().name) - .getOpStatsLogger(KopServerStats.REQUEST_QUEUED_LATENCY) - .registerFailedEvent(MathUtils.elapsedNanos(responseAndRequest.getCreatedTimestamp()), - TimeUnit.NANOSECONDS); + requestStats.getRequestStatsLogger(apiKey, KopServerStats.REQUEST_QUEUED_LATENCY) + .registerFailedEvent(nanoSecondsSinceCreated, TimeUnit.NANOSECONDS); return null; }); // send exception to client? continue; @@ -446,6 +439,8 @@ protected void writeAndFlushResponseToClient(Channel channel) { log.error("[{}] Failed to write {}", channel, request.getHeader(), future.cause()); } }); + requestStats.getRequestStatsLogger(apiKey, KopServerStats.REQUEST_QUEUED_LATENCY) + .registerSuccessfulEvent(nanoSecondsSinceCreated, TimeUnit.NANOSECONDS); }); continue; } @@ -457,12 +452,8 @@ protected void writeAndFlushResponseToClient(Channel channel) { responseFuture.cancel(true); channel.writeAndFlush( request.createErrorResponse(new ApiException("request is expired from server side"))); - - requestStats.getStatsLogger() - .scopeLabel(KopServerStats.REQUEST_SCOPE, responseAndRequest.request.getHeader().apiKey().name) - .getOpStatsLogger(KopServerStats.REQUEST_QUEUED_LATENCY) - .registerFailedEvent(MathUtils.elapsedNanos(responseAndRequest.getCreatedTimestamp()), - TimeUnit.NANOSECONDS); + requestStats.getRequestStatsLogger(apiKey, KopServerStats.REQUEST_QUEUED_LATENCY) + .registerFailedEvent(nanoSecondsSinceCreated, TimeUnit.NANOSECONDS); } } } @@ -472,7 +463,7 @@ protected void writeAndFlushResponseToClient(Channel channel) { protected abstract void channelPrepare(ChannelHandlerContext ctx, ByteBuf requestBuf, BiConsumer registerRequestParseLatency, - BiConsumer registerRequestLatency) + BiConsumer registerRequestLatency) throws AuthenticationException; protected abstract void maybeDelayCloseOnAuthenticationFailure(); @@ -702,16 +693,12 @@ public long nanoSecondsSinceCreated() { return MathUtils.elapsedNanos(createdTimestamp); } - public boolean expired(final int requestTimeoutMs) { - return MathUtils.elapsedNanos(createdTimestamp) > TimeUnit.MILLISECONDS.toNanos(requestTimeoutMs); + public ApiKeys getApiKey() { + return request.getHeader().apiKey(); } - public void updateStats(final RequestStats requestStats) { - RequestStats.REQUEST_QUEUE_SIZE_INSTANCE.decrementAndGet(); - requestStats.getStatsLogger() - .scopeLabel(KopServerStats.REQUEST_SCOPE, request.getHeader().apiKey().name) - .getOpStatsLogger(KopServerStats.REQUEST_QUEUED_LATENCY) - .registerSuccessfulEvent(MathUtils.elapsedNanos(createdTimestamp), TimeUnit.NANOSECONDS); + public boolean expired(final int requestTimeoutMs) { + return MathUtils.elapsedNanos(createdTimestamp) > TimeUnit.MILLISECONDS.toNanos(requestTimeoutMs); } ResponseAndRequest(CompletableFuture response, KafkaHeaderAndRequest request) { diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index 6e9d53e6bd..46e3770b45 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -74,8 +74,8 @@ public class KafkaProtocolHandler implements ProtocolHandler, TenantContextManag public static final String TLS_HANDLER = "tls"; private static final Map LOOKUP_CLIENT_MAP = new ConcurrentHashMap<>(); - private StatsLogger rootStatsLogger; - private StatsLogger scopeStatsLogger; + @Getter + private RequestStats requestStats; private PrometheusMetricsProvider statsProvider; @Getter private KopBrokerLookupManager kopBrokerLookupManager; @@ -435,8 +435,8 @@ public void initialize(ServiceConfiguration conf) throws Exception { } statsProvider = new PrometheusMetricsProvider(); - rootStatsLogger = statsProvider.getStatsLogger(""); - scopeStatsLogger = rootStatsLogger.scope(SERVER_SCOPE); + StatsLogger rootStatsLogger = statsProvider.getStatsLogger(""); + requestStats = new RequestStats(rootStatsLogger.scope(SERVER_SCOPE)); sendResponseScheduler = OrderedScheduler.newSchedulerBuilder() .name("send-response") .numThreads(kafkaConfig.getNumSendKafkaResponseThreads()) @@ -492,7 +492,7 @@ public void start(BrokerService service) { // init KopEventManager kopEventManager = new KopEventManager(adminManager, brokerService.getPulsar().getLocalMetadataStore(), - scopeStatsLogger, + requestStats.getStatsLogger(), kafkaConfig, groupCoordinatorsByTenant); kopEventManager.start(); @@ -583,7 +583,7 @@ private KafkaChannelInitializer newKafkaChannelInitializer(final EndPoint endPoi endPoint.isTlsEnabled(), endPoint, kafkaConfig.isSkipMessagesWithoutIndex(), - scopeStatsLogger, + requestStats, sendResponseScheduler); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index 774370b1c5..3a4ff4e61a 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -42,7 +42,6 @@ import io.streamnative.pulsar.handlers.kop.security.auth.Resource; import io.streamnative.pulsar.handlers.kop.security.auth.ResourceType; import io.streamnative.pulsar.handlers.kop.security.auth.SimpleAclAuthorizer; -import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; import io.streamnative.pulsar.handlers.kop.storage.AppendRecordsContext; import io.streamnative.pulsar.handlers.kop.storage.PartitionLog; import io.streamnative.pulsar.handlers.kop.storage.ReplicaManager; @@ -289,9 +288,9 @@ public KafkaRequestHandler(PulsarService pulsarService, Boolean tlsEnabled, EndPoint advertisedEndPoint, boolean skipMessagesWithoutIndex, - StatsLogger statsLogger, + RequestStats requestStats, OrderedScheduler sendResponseScheduler) throws Exception { - super(statsLogger, kafkaConfig, sendResponseScheduler); + super(requestStats, kafkaConfig, sendResponseScheduler); this.pulsarService = pulsarService; this.tenantContextManager = tenantContextManager; this.kopBrokerLookupManager = kopBrokerLookupManager; @@ -376,7 +375,7 @@ protected boolean hasAuthenticated() { protected void channelPrepare(ChannelHandlerContext ctx, ByteBuf requestBuf, BiConsumer registerRequestParseLatency, - BiConsumer registerRequestLatency) + BiConsumer registerRequestLatency) throws AuthenticationException { if (authenticator != null) { authenticator.authenticate(ctx, requestBuf, registerRequestParseLatency, registerRequestLatency, diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/RequestStats.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/RequestStats.java index 99e79faa48..cd57e2b04b 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/RequestStats.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/RequestStats.java @@ -30,7 +30,13 @@ import static io.streamnative.pulsar.handlers.kop.KopServerStats.SERVER_SCOPE; import static io.streamnative.pulsar.handlers.kop.KopServerStats.WAITING_FETCHES_TRIGGERED; +import com.google.common.annotations.VisibleForTesting; +import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; +import java.util.Map; +import java.util.Set; +import java.util.TreeSet; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; import lombok.Getter; import lombok.extern.slf4j.Slf4j; @@ -38,6 +44,7 @@ import org.apache.bookkeeper.stats.Gauge; import org.apache.bookkeeper.stats.OpStatsLogger; import org.apache.bookkeeper.stats.annotations.StatsDoc; +import org.apache.kafka.common.protocol.ApiKeys; /** * Kop request stats metric for prometheus metrics. @@ -57,6 +64,8 @@ public class RequestStats { public static final AtomicInteger ALIVE_CHANNEL_COUNT_INSTANCE = new AtomicInteger(0); public static final AtomicInteger ACTIVE_CHANNEL_COUNT_INSTANCE = new AtomicInteger(0); + public static final RequestStats NULL_INSTANCE = new RequestStats(NullStatsLogger.INSTANCE); + private final StatsLogger statsLogger; @StatsDoc( @@ -119,6 +128,8 @@ public class RequestStats { ) private final Counter waitingFetchesTriggered; + private final Map apiKeysToStatsLogger = new ConcurrentHashMap<>(); + public RequestStats(StatsLogger statsLogger) { this.statsLogger = statsLogger; @@ -184,4 +195,22 @@ public Number getSample() { } }); } + + /** + * Get the stats logger for Kafka requests. + * + * @param apiKey the {@link ApiKeys} object that represents the Kafka request's type + * @param statsName the stats name + * @return + */ + public OpStatsLogger getRequestStatsLogger(final ApiKeys apiKey, final String statsName) { + return apiKeysToStatsLogger.computeIfAbsent(apiKey, + __ -> statsLogger.scopeLabel(KopServerStats.REQUEST_SCOPE, apiKey.name) + ).getOpStatsLogger(statsName); + } + + @VisibleForTesting + public Set getApiKeysSet() { + return new TreeSet<>(apiKeysToStatsLogger.keySet()); + } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/SaslAuthenticator.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/SaslAuthenticator.java index 3d4e6cba6f..f74efd1e74 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/SaslAuthenticator.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/SaslAuthenticator.java @@ -177,7 +177,7 @@ public SaslAuthenticator(PulsarAdmin admin, public void authenticate(ChannelHandlerContext ctx, ByteBuf requestBuf, BiConsumer registerRequestParseLatency, - BiConsumer registerRequestLatency, + BiConsumer registerRequestLatency, Function tenantAccessValidationFunction) throws AuthenticationException { checkArgument(requestBuf.readableBytes() > 0); @@ -303,7 +303,7 @@ private AbstractRequest parseRequest(RequestHeader header, ByteBuffer nioBuffer) private void handleKafkaRequest(ChannelHandlerContext ctx, ByteBuf requestBuf, BiConsumer registerRequestParseLatency, - BiConsumer registerRequestLatency) + BiConsumer registerRequestLatency) throws AuthenticationException { final long beforeParseTime = MathUtils.nowInNano(); ByteBuffer nioBuffer = requestBuf.nioBuffer(); @@ -407,7 +407,7 @@ public static ByteBuf sizePrefixed(ByteBuffer buffer) { private void handleSaslToken(ChannelHandlerContext ctx, ByteBuf requestBuf, BiConsumer registerRequestParseLatency, - BiConsumer registerRequestLatency, + BiConsumer registerRequestLatency, Function tenantAccessValidationFunction) throws AuthenticationException { final long timeBeforeParse = MathUtils.nowInNano(); @@ -462,7 +462,7 @@ private void handleSaslToken(ChannelHandlerContext ctx, if (apiKey != ApiKeys.SASL_AUTHENTICATE) { AuthenticationException e = new AuthenticationException( "Unexpected Kafka request of type " + apiKey + " during SASL authentication"); - registerRequestLatency.accept(apiKey.name, startProcessTime); + registerRequestLatency.accept(apiKey, startProcessTime); buildResponseOnAuthenticateFailure(header, request, null, e); throw e; } @@ -482,11 +482,11 @@ private void handleSaslToken(ChannelHandlerContext ctx, (String) saslServer.getNegotiatedProperty(USER_NAME_PROP), (AuthenticationDataSource) saslServer.getNegotiatedProperty(AUTH_DATA_SOURCE_PROP)), header.clientId()); - registerRequestLatency.accept(apiKey.name, startProcessTime); + registerRequestLatency.accept(apiKey, startProcessTime); if (!tenantAccessValidationFunction.apply(session)) { AuthenticationException e = new AuthenticationException("User is not allowed to access this tenant"); - registerRequestLatency.accept(apiKey.name, startProcessTime); + registerRequestLatency.accept(apiKey, startProcessTime); buildResponseOnAuthenticateFailure(header, request, null, e); throw e; } @@ -503,7 +503,7 @@ private void handleSaslToken(ChannelHandlerContext ctx, saslServer.getNegotiatedProperty(AUTH_DATA_SOURCE_PROP)); } } catch (SaslException e) { - registerRequestLatency.accept(apiKey.name, startProcessTime); + registerRequestLatency.accept(apiKey, startProcessTime); buildResponseOnAuthenticateFailure(header, request, KafkaResponseUtils.newSaslAuthenticate(Errors.SASL_AUTHENTICATION_FAILED, e.getMessage()), e); sendAuthenticationFailureResponse(); @@ -519,20 +519,20 @@ private void handleApiVersionsRequest(ChannelHandlerContext ctx, RequestHeader header, ApiVersionsRequest request, Long startProcessTime, - BiConsumer registerRequestLatency) + BiConsumer registerRequestLatency) throws AuthenticationException { if (state != State.HANDSHAKE_OR_VERSIONS_REQUEST) { throw new IllegalStateException( "Receive ApiVersions request", state, State.HANDSHAKE_OR_VERSIONS_REQUEST); } if (request.hasUnsupportedRequestVersion()) { - registerRequestLatency.accept(header.apiKey().name, startProcessTime); + registerRequestLatency.accept(header.apiKey(), startProcessTime); sendKafkaResponse(ctx, header, request, request.getErrorResponse(0, Errors.UNSUPPORTED_VERSION.exception()), null); } else { ApiVersionsResponse versionsResponse = ApiVersionsResponse.defaultApiVersionsResponse(); - registerRequestLatency.accept(header.apiKey().name, startProcessTime); + registerRequestLatency.accept(header.apiKey(), startProcessTime); sendKafkaResponse(ctx, header, request, @@ -547,13 +547,13 @@ private void handleApiVersionsRequest(ChannelHandlerContext ctx, RequestHeader header, SaslHandshakeRequest request, Long startProcessTime, - BiConsumer registerRequestLatency) + BiConsumer registerRequestLatency) throws AuthenticationException { final String mechanism = request.mechanism(); if (mechanism == null) { AuthenticationException e = new AuthenticationException("client's mechanism is null"); - registerRequestLatency.accept(header.apiKey().name, startProcessTime); + registerRequestLatency.accept(header.apiKey(), startProcessTime); sendKafkaResponse(ctx, header, request, @@ -570,7 +570,7 @@ private void handleApiVersionsRequest(ChannelHandlerContext ctx, if (log.isDebugEnabled()) { log.debug("Using SASL mechanism '{}' provided by client", mechanism); } - registerRequestLatency.accept(header.apiKey().name, startProcessTime); + registerRequestLatency.accept(header.apiKey(), startProcessTime); sendKafkaResponse(ctx, header, request, @@ -581,7 +581,7 @@ private void handleApiVersionsRequest(ChannelHandlerContext ctx, if (log.isDebugEnabled()) { log.debug("SASL mechanism '{}' requested by client is not supported", mechanism); } - registerRequestLatency.accept(header.apiKey().name, startProcessTime); + registerRequestLatency.accept(header.apiKey(), startProcessTime); buildResponseOnAuthenticateFailure(header, request, KafkaResponseUtils.newSaslHandshake(Errors.UNSUPPORTED_SASL_MECHANISM, allowedMechanisms), null); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java index 9675029b6e..5615554fa3 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java @@ -25,7 +25,6 @@ import io.netty.channel.EventLoopGroup; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; -import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import io.streamnative.pulsar.handlers.kop.storage.ReplicaManager; import io.streamnative.pulsar.handlers.kop.utils.MetadataUtils; import java.io.Closeable; @@ -796,7 +795,7 @@ public KafkaRequestHandler newRequestHandler() throws Exception { handler.getReplicaManager(conf.getKafkaMetadataTenant()); return ((KafkaChannelInitializer) handler.getChannelInitializerMap().entrySet().iterator().next().getValue()) - .newCnx(new TenantContextManager() { + .newCnxWithoutStats(new TenantContextManager() { @Override public GroupCoordinator getGroupCoordinator(String tenant) { return groupCoordinator; @@ -811,6 +810,6 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) { public ReplicaManager getReplicaManager(String tenant) { return replicaManager; } - }, NullStatsLogger.INSTANCE); + }); } } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MetricsProviderTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MetricsProviderTest.java index 28221b2217..0bfeddfb7c 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MetricsProviderTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MetricsProviderTest.java @@ -18,9 +18,12 @@ import java.io.InputStreamReader; import java.nio.charset.StandardCharsets; import java.time.Duration; +import java.util.Arrays; import java.util.Collections; import java.util.List; import java.util.Properties; +import java.util.Set; +import java.util.TreeSet; import java.util.regex.Matcher; import java.util.regex.Pattern; import java.util.stream.Collectors; @@ -37,6 +40,7 @@ import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.protocol.ApiKeys; import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -61,6 +65,10 @@ protected void cleanup() throws Exception { super.internalCleanup(); } + private Set getApiKeysSet() { + return ((KafkaProtocolHandler) pulsar.getProtocolHandlers().protocol("kafka")) + .getRequestStats().getApiKeysSet(); + } @Test(timeOut = 30000) public void testMetricsProvider() throws Exception { @@ -92,6 +100,9 @@ public void testMetricsProvider() throws Exception { } } + Assert.assertEquals(getApiKeysSet(), new TreeSet<>( + Arrays.asList(ApiKeys.API_VERSIONS, ApiKeys.METADATA, ApiKeys.PRODUCE))); + // 2. consume messages with Kafka consumer @Cleanup KConsumer kConsumer = new KConsumer(kafkaTopicName, getKafkaBrokerPort()); @@ -115,8 +126,17 @@ public void testMetricsProvider() throws Exception { } Assert.assertEquals(msgs, totalMsgs); + Assert.assertEquals(getApiKeysSet(), new TreeSet<>(Arrays.asList( + ApiKeys.API_VERSIONS, ApiKeys.METADATA, ApiKeys.PRODUCE, ApiKeys.FIND_COORDINATOR, ApiKeys.LIST_OFFSETS, + ApiKeys.OFFSET_FETCH, ApiKeys.FETCH + ))); + // commit offsets kConsumer.getConsumer().commitSync(Duration.ofSeconds(5)); + Assert.assertEquals(getApiKeysSet(), new TreeSet<>(Arrays.asList( + ApiKeys.API_VERSIONS, ApiKeys.METADATA, ApiKeys.PRODUCE, ApiKeys.FIND_COORDINATOR, ApiKeys.LIST_OFFSETS, + ApiKeys.OFFSET_FETCH, ApiKeys.FETCH, ApiKeys.OFFSET_COMMIT + ))); try { Thread.sleep(1000);