From b6b70027209cee63346b3be94980add3cd21c201 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 4 Nov 2021 11:09:13 +0800 Subject: [PATCH 1/5] Maintain a StatsLogger map to avoid creating new object every time --- .../handlers/kop/KafkaCommandDecoder.java | 67 +++++++++++-------- .../handlers/kop/KafkaRequestHandler.java | 2 +- .../kop/security/SaslAuthenticator.java | 28 ++++---- 3 files changed, 54 insertions(+), 43 deletions(-) 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..aef0b67342 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 @@ -25,8 +25,10 @@ import java.io.Closeable; import java.net.SocketAddress; import java.nio.ByteBuffer; +import java.util.Map; import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -36,6 +38,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.util.MathUtils; import org.apache.bookkeeper.common.util.OrderedScheduler; +import org.apache.bookkeeper.stats.OpStatsLogger; import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.errors.AuthenticationException; import org.apache.kafka.common.protocol.ApiKeys; @@ -59,6 +62,7 @@ public abstract class KafkaCommandDecoder extends ChannelInboundHandlerAdapter { protected SocketAddress remoteAddress; @Getter protected AtomicBoolean isActive = new AtomicBoolean(false); + private final Map apiKeysToStatsLogger = new ConcurrentHashMap<>(); // Queue to make response get responseFuture in order and limit the max request size private final LinkedBlockingQueue requestQueue; @@ -77,6 +81,20 @@ public KafkaCommandDecoder(StatsLogger statsLogger, this.sendResponseScheduler = sendResponseScheduler; } + private void registerOperationLatency(final ApiKeys apiKey, + final long startProcessTimeNs, + final String operationName, + final boolean success) { + final OpStatsLogger opStatsLogger = apiKeysToStatsLogger.computeIfAbsent(apiKey, __ -> + requestStats.getStatsLogger().scopeLabel(KopServerStats.REQUEST_SCOPE, apiKey.name()) + ).getOpStatsLogger(operationName); + if (success) { + opStatsLogger.registerSuccessfulEvent(MathUtils.elapsedNanos(startProcessTimeNs), TimeUnit.NANOSECONDS); + } else { + opStatsLogger.registerFailedEvent(MathUtils.elapsedNanos(startProcessTimeNs), TimeUnit.NANOSECONDS); + } + } + @Override public void channelActive(ChannelHandlerContext ctx) throws Exception { super.channelActive(ctx); @@ -193,12 +211,8 @@ 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) -> { + registerOperationLatency(apiKey, startProcessTime, KopServerStats.REQUEST_LATENCY, true); }; // If kop is enabled for authentication and the client @@ -248,7 +262,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(), () -> { @@ -390,7 +404,11 @@ protected void writeAndFlushResponseToClient(Channel channel) { break; } else { if (requestQueue.remove(responseAndRequest)) { - responseAndRequest.updateStats(requestStats); + RequestStats.REQUEST_QUEUE_SIZE_INSTANCE.decrementAndGet(); + registerOperationLatency(responseAndRequest.getApiKey(), + responseAndRequest.getCreatedTimestamp(), + KopServerStats.REQUEST_QUEUED_LATENCY, + true); } else { // it has been removed by another thread, skip this element continue; } @@ -409,12 +427,10 @@ 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); + registerOperationLatency(responseAndRequest.getApiKey(), + responseAndRequest.getCreatedTimestamp(), + KopServerStats.REQUEST_QUEUED_LATENCY, + false); return null; }); // send exception to client? continue; @@ -458,11 +474,10 @@ protected void writeAndFlushResponseToClient(Channel channel) { 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); + registerOperationLatency(responseAndRequest.getApiKey(), + responseAndRequest.getCreatedTimestamp(), + KopServerStats.REQUEST_QUEUED_LATENCY, + false); } } } @@ -472,7 +487,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 +717,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/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index 774370b1c5..57d980436b 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 @@ -376,7 +376,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/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); From 639d88be6c7fe14ce381b003dd695b023e5dfb54 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Mon, 17 Jan 2022 16:26:39 +0800 Subject: [PATCH 2/5] Fix tests --- .../streamnative/pulsar/handlers/kop/KafkaCommandDecoder.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 aef0b67342..5b0a5db559 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 @@ -86,7 +86,7 @@ private void registerOperationLatency(final ApiKeys apiKey, final String operationName, final boolean success) { final OpStatsLogger opStatsLogger = apiKeysToStatsLogger.computeIfAbsent(apiKey, __ -> - requestStats.getStatsLogger().scopeLabel(KopServerStats.REQUEST_SCOPE, apiKey.name()) + requestStats.getStatsLogger().scopeLabel(KopServerStats.REQUEST_SCOPE, apiKey.name) ).getOpStatsLogger(operationName); if (success) { opStatsLogger.registerSuccessfulEvent(MathUtils.elapsedNanos(startProcessTimeNs), TimeUnit.NANOSECONDS); From 9f2a4c21a5d34a8b05d755dbd76e8667948a5b9f Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 18 Jan 2022 17:23:43 +0800 Subject: [PATCH 3/5] Make all connections share the same RequestStats instance --- .../handlers/kop/KafkaChannelInitializer.java | 15 +++++++-------- .../pulsar/handlers/kop/KafkaCommandDecoder.java | 4 ++-- .../pulsar/handlers/kop/KafkaProtocolHandler.java | 11 +++++------ .../pulsar/handlers/kop/KafkaRequestHandler.java | 5 ++--- .../pulsar/handlers/kop/RequestStats.java | 3 +++ .../handlers/kop/KopProtocolHandlerTestBase.java | 5 ++--- 6 files changed, 21 insertions(+), 22 deletions(-) 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 5b0a5db559..ab9c02696e 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 @@ -72,10 +72,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; 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..fd63c115ef 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,7 @@ 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; + private RequestStats requestStats; private PrometheusMetricsProvider statsProvider; @Getter private KopBrokerLookupManager kopBrokerLookupManager; @@ -435,8 +434,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 +491,7 @@ public void start(BrokerService service) { // init KopEventManager kopEventManager = new KopEventManager(adminManager, brokerService.getPulsar().getLocalMetadataStore(), - scopeStatsLogger, + requestStats.getStatsLogger(), kafkaConfig, groupCoordinatorsByTenant); kopEventManager.start(); @@ -583,7 +582,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 57d980436b..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; 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..ce42322809 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,6 +30,7 @@ import static io.streamnative.pulsar.handlers.kop.KopServerStats.SERVER_SCOPE; import static io.streamnative.pulsar.handlers.kop.KopServerStats.WAITING_FETCHES_TRIGGERED; +import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; import java.util.concurrent.atomic.AtomicInteger; import lombok.Getter; @@ -57,6 +58,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( 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); + }); } } From f97cc32ab89f738b74735949b5ff2a88ddac0d2c Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 18 Jan 2022 18:26:30 +0800 Subject: [PATCH 4/5] Remove the wrapper method --- .../handlers/kop/KafkaCommandDecoder.java | 42 ++++--------------- .../pulsar/handlers/kop/RequestStats.java | 18 ++++++++ 2 files changed, 27 insertions(+), 33 deletions(-) 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 ab9c02696e..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,14 +21,11 @@ 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; -import java.util.Map; import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -38,7 +35,6 @@ import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.util.MathUtils; import org.apache.bookkeeper.common.util.OrderedScheduler; -import org.apache.bookkeeper.stats.OpStatsLogger; import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.errors.AuthenticationException; import org.apache.kafka.common.protocol.ApiKeys; @@ -62,7 +58,6 @@ public abstract class KafkaCommandDecoder extends ChannelInboundHandlerAdapter { protected SocketAddress remoteAddress; @Getter protected AtomicBoolean isActive = new AtomicBoolean(false); - private final Map apiKeysToStatsLogger = new ConcurrentHashMap<>(); // Queue to make response get responseFuture in order and limit the max request size private final LinkedBlockingQueue requestQueue; @@ -81,20 +76,6 @@ public KafkaCommandDecoder(RequestStats requestStats, this.sendResponseScheduler = sendResponseScheduler; } - private void registerOperationLatency(final ApiKeys apiKey, - final long startProcessTimeNs, - final String operationName, - final boolean success) { - final OpStatsLogger opStatsLogger = apiKeysToStatsLogger.computeIfAbsent(apiKey, __ -> - requestStats.getStatsLogger().scopeLabel(KopServerStats.REQUEST_SCOPE, apiKey.name) - ).getOpStatsLogger(operationName); - if (success) { - opStatsLogger.registerSuccessfulEvent(MathUtils.elapsedNanos(startProcessTimeNs), TimeUnit.NANOSECONDS); - } else { - opStatsLogger.registerFailedEvent(MathUtils.elapsedNanos(startProcessTimeNs), TimeUnit.NANOSECONDS); - } - } - @Override public void channelActive(ChannelHandlerContext ctx) throws Exception { super.channelActive(ctx); @@ -212,7 +193,8 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception // Update handle request latency metrics final BiConsumer registerRequestLatency = (apiKey, startProcessTime) -> { - registerOperationLatency(apiKey, startProcessTime, KopServerStats.REQUEST_LATENCY, true); + requestStats.getRequestStatsLogger(apiKey, KopServerStats.REQUEST_LATENCY) + .registerSuccessfulEvent(MathUtils.elapsedNanos(startProcessTime), TimeUnit.NANOSECONDS); }; // If kop is enabled for authentication and the client @@ -391,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())); @@ -405,10 +388,6 @@ protected void writeAndFlushResponseToClient(Channel channel) { } else { if (requestQueue.remove(responseAndRequest)) { RequestStats.REQUEST_QUEUE_SIZE_INSTANCE.decrementAndGet(); - registerOperationLatency(responseAndRequest.getApiKey(), - responseAndRequest.getCreatedTimestamp(), - KopServerStats.REQUEST_QUEUED_LATENCY, - true); } else { // it has been removed by another thread, skip this element continue; } @@ -427,10 +406,8 @@ protected void writeAndFlushResponseToClient(Channel channel) { log.error("[{}] request {} completed exceptionally", channel, request.getHeader(), e); channel.writeAndFlush(request.createErrorResponse(e)); - registerOperationLatency(responseAndRequest.getApiKey(), - responseAndRequest.getCreatedTimestamp(), - KopServerStats.REQUEST_QUEUED_LATENCY, - false); + requestStats.getRequestStatsLogger(apiKey, KopServerStats.REQUEST_QUEUED_LATENCY) + .registerFailedEvent(nanoSecondsSinceCreated, TimeUnit.NANOSECONDS); return null; }); // send exception to client? continue; @@ -462,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; } @@ -473,11 +452,8 @@ protected void writeAndFlushResponseToClient(Channel channel) { responseFuture.cancel(true); channel.writeAndFlush( request.createErrorResponse(new ApiException("request is expired from server side"))); - - registerOperationLatency(responseAndRequest.getApiKey(), - responseAndRequest.getCreatedTimestamp(), - KopServerStats.REQUEST_QUEUED_LATENCY, - false); + requestStats.getRequestStatsLogger(apiKey, KopServerStats.REQUEST_QUEUED_LATENCY) + .registerFailedEvent(nanoSecondsSinceCreated, TimeUnit.NANOSECONDS); } } } 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 ce42322809..07e12ce5fe 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 @@ -32,6 +32,8 @@ import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; import lombok.Getter; import lombok.extern.slf4j.Slf4j; @@ -39,6 +41,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. @@ -122,6 +125,8 @@ public class RequestStats { ) private final Counter waitingFetchesTriggered; + private final Map apiKeysToStatsLogger = new ConcurrentHashMap<>(); + public RequestStats(StatsLogger statsLogger) { this.statsLogger = statsLogger; @@ -187,4 +192,17 @@ 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); + } } From 66fb7ac714085092e47979d9f0fc682710550904 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 18 Jan 2022 21:55:46 +0800 Subject: [PATCH 5/5] Check ApiKeys set from RequestStats --- .../handlers/kop/KafkaProtocolHandler.java | 1 + .../pulsar/handlers/kop/RequestStats.java | 8 ++++++++ .../handlers/kop/MetricsProviderTest.java | 20 +++++++++++++++++++ 3 files changed, 29 insertions(+) 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 fd63c115ef..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,6 +74,7 @@ public class KafkaProtocolHandler implements ProtocolHandler, TenantContextManag public static final String TLS_HANDLER = "tls"; private static final Map LOOKUP_CLIENT_MAP = new ConcurrentHashMap<>(); + @Getter private RequestStats requestStats; private PrometheusMetricsProvider statsProvider; @Getter 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 07e12ce5fe..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,9 +30,12 @@ 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; @@ -205,4 +208,9 @@ public OpStatsLogger getRequestStatsLogger(final ApiKeys apiKey, final String st __ -> statsLogger.scopeLabel(KopServerStats.REQUEST_SCOPE, apiKey.name) ).getOpStatsLogger(statsName); } + + @VisibleForTesting + public Set getApiKeysSet() { + return new TreeSet<>(apiKeysToStatsLogger.keySet()); + } } 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);