diff --git a/src/main/java/io/streamnative/kop/InternalProducer.java b/src/main/java/io/streamnative/kop/InternalProducer.java new file mode 100644 index 0000000000..d9a0450acc --- /dev/null +++ b/src/main/java/io/streamnative/kop/InternalProducer.java @@ -0,0 +1,53 @@ +/** + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.streamnative.kop; + +import java.util.concurrent.CompletableFuture; +import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.broker.service.Producer; +import org.apache.pulsar.broker.service.ServerCnx; +import org.apache.pulsar.broker.service.Topic; + +/** + * InternalServerCnx, this only used to construct internalProducer / internalConsumer. + * So when topic is unload, we could disconnect the connection between kafkaRequestHandler and client, + * by internalProducer / internalConsumer.close(); + * which means when topic unload happens, we should close the connection. + */ +@Slf4j +public class InternalProducer extends Producer { + public InternalProducer(Topic topic, ServerCnx cnx, + long producerId, String producerName) { + super(topic, cnx, producerId, producerName, null, + false, null, null); + } + + // this will call back by bundle unload + @Override + public CompletableFuture disconnect() { + InternalServerCnx cnx = (InternalServerCnx) getCnx(); + CompletableFuture future = new CompletableFuture<>(); + + cnx.getBrokerService().executor().execute(() -> { + log.info("Disconnecting producer: {}", this); + getTopic().removeProducer(this); + cnx.closeProducer(this); + future.complete(null); + }); + + return future; + } + + +} diff --git a/src/main/java/io/streamnative/kop/InternalServerCnx.java b/src/main/java/io/streamnative/kop/InternalServerCnx.java new file mode 100644 index 0000000000..7894f8a727 --- /dev/null +++ b/src/main/java/io/streamnative/kop/InternalServerCnx.java @@ -0,0 +1,63 @@ +/** + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.streamnative.kop; + +import java.net.InetSocketAddress; +import lombok.Getter; +import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.broker.service.Producer; +import org.apache.pulsar.broker.service.ServerCnx; + +/** + * InternalServerCnx, this only used to construct internalProducer / internalConsumer. + * So when topic is unload, we could disconnect the connection between kafkaRequestHandler and client, + * by internalProducer / internalConsumer.close(); + * which means when topic unload happens, we should close the connection. + */ +@Slf4j +public class InternalServerCnx extends ServerCnx { + @Getter + KafkaRequestHandler kafkaRequestHandler; + + public InternalServerCnx(KafkaRequestHandler kafkaRequestHandler) { + super(kafkaRequestHandler.getPulsarService()); + this.kafkaRequestHandler = kafkaRequestHandler; + // this is the client address that connect to this server. + this.remoteAddress = kafkaRequestHandler.remoteAddress; + + // mock some values, or Producer create will meet NPE. + // used in test, which will not call channel.active, and not call updateCtx. + if (this.remoteAddress == null) { + this.remoteAddress = new InetSocketAddress("localhost", 9999); + } + } + + // this will call back by bundle unload + @Override + public void closeProducer(Producer producer) { + // removes producer-connection from map and send close command to producer + if (log.isDebugEnabled()) { + log.debug("[{}] Removed topic: {}'s producer: {}.", + remoteAddress, producer.getTopic().getName(), producer); + } + + kafkaRequestHandler.close(); + } + + // called after channel active + public void updateCtx() { + this.remoteAddress = kafkaRequestHandler.remoteAddress; + } + +} diff --git a/src/main/java/io/streamnative/kop/KafkaChannelInitializer.java b/src/main/java/io/streamnative/kop/KafkaChannelInitializer.java index 647bfe7291..3858ff102d 100644 --- a/src/main/java/io/streamnative/kop/KafkaChannelInitializer.java +++ b/src/main/java/io/streamnative/kop/KafkaChannelInitializer.java @@ -38,8 +38,6 @@ public class KafkaChannelInitializer extends ChannelInitializer { @Getter private final KafkaServiceConfiguration kafkaConfig; @Getter - private final KafkaTopicManager kafkaTopicManager; - @Getter private final GroupCoordinator groupCoordinator; @Getter private final boolean enableTls; @@ -48,13 +46,11 @@ public class KafkaChannelInitializer extends ChannelInitializer { public KafkaChannelInitializer(PulsarService pulsarService, KafkaServiceConfiguration kafkaConfig, - KafkaTopicManager kafkaTopicManager, GroupCoordinator groupCoordinator, boolean enableTLS) throws Exception { super(); this.pulsarService = pulsarService; this.kafkaConfig = kafkaConfig; - this.kafkaTopicManager = kafkaTopicManager; this.groupCoordinator = groupCoordinator; this.enableTls = enableTLS; @@ -74,7 +70,7 @@ protected void initChannel(SocketChannel ch) throws Exception { ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder(MAX_FRAME_LENGTH, 0, 4, 0, 4)); ch.pipeline().addLast("handler", - new KafkaRequestHandler(pulsarService, kafkaConfig, kafkaTopicManager, groupCoordinator, enableTls)); + new KafkaRequestHandler(pulsarService, kafkaConfig, groupCoordinator, enableTls)); } } diff --git a/src/main/java/io/streamnative/kop/KafkaCommandDecoder.java b/src/main/java/io/streamnative/kop/KafkaCommandDecoder.java index 0cc7270cba..9c30e6bda6 100644 --- a/src/main/java/io/streamnative/kop/KafkaCommandDecoder.java +++ b/src/main/java/io/streamnative/kop/KafkaCommandDecoder.java @@ -15,6 +15,7 @@ import static com.google.common.base.Preconditions.checkArgument; import static org.apache.kafka.common.protocol.ApiKeys.API_VERSIONS; +import static org.apache.kafka.common.protocol.ApiKeys.METADATA; import com.google.common.collect.Queues; import io.netty.buffer.ByteBuf; @@ -28,8 +29,10 @@ import java.util.Queue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicBoolean; import lombok.Getter; import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.common.errors.LeaderNotAvailableException; import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.types.Struct; import org.apache.kafka.common.requests.AbstractRequest; @@ -46,31 +49,22 @@ public abstract class KafkaCommandDecoder extends ChannelInboundHandlerAdapter { protected ChannelHandlerContext ctx; protected SocketAddress remoteAddress; + @Getter + protected AtomicBoolean isActive = new AtomicBoolean(false); - // Queue to make request get response in order. - protected final ConcurrentHashMap>> responsesQueue = + // Queue to make request get responseFuture in order. + protected final ConcurrentHashMap> responsesQueue = new ConcurrentHashMap(); - // TODO: do we need keep alive? if need, messageReceived() before every command is need? - public KafkaCommandDecoder() { } @Override public void channelActive(ChannelHandlerContext ctx) throws Exception { + super.channelActive(ctx); this.remoteAddress = ctx.channel().remoteAddress(); this.ctx = ctx; - - if (log.isDebugEnabled()) { - log.debug("[{}] channel active {}", ctx.channel()); - } - } - - @Override - public void channelInactive(ChannelHandlerContext ctx) throws Exception { - super.channelInactive(ctx); - log.info("Closed connection from {}", remoteAddress); - close(); + isActive.set(true); } @Override @@ -137,8 +131,7 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception remoteAddress = channel.remoteAddress(); } - CompletableFuture responseFuture; - + CompletableFuture responseFuture; try (KafkaHeaderAndRequest kafkaHeaderAndRequest = byteBufToRequest(buffer, remoteAddress)){ if (log.isDebugEnabled()) { @@ -147,69 +140,75 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception kafkaHeaderAndRequest.getHeader(), kafkaHeaderAndRequest); } - switch (kafkaHeaderAndRequest.getHeader().apiKey()) { - case API_VERSIONS: - responseFuture = handleApiVersionsRequest(kafkaHeaderAndRequest); - break; - case METADATA: - responseFuture = handleTopicMetadataRequest(kafkaHeaderAndRequest); - break; - case PRODUCE: - responseFuture = handleProduceRequest(kafkaHeaderAndRequest); - break; - case FIND_COORDINATOR: - responseFuture = handleFindCoordinatorRequest(kafkaHeaderAndRequest); - break; - case LIST_OFFSETS: - responseFuture = handleListOffsetRequest(kafkaHeaderAndRequest); - break; - case OFFSET_FETCH: - responseFuture = handleOffsetFetchRequest(kafkaHeaderAndRequest); - break; - case OFFSET_COMMIT: - responseFuture = handleOffsetCommitRequest(kafkaHeaderAndRequest); - break; - case FETCH: - responseFuture = handleFetchRequest(kafkaHeaderAndRequest); - break; - case JOIN_GROUP: - responseFuture = handleJoinGroupRequest(kafkaHeaderAndRequest); - break; - case SYNC_GROUP: - responseFuture = handleSyncGroupRequest(kafkaHeaderAndRequest); - break; - case HEARTBEAT: - responseFuture = handleHeartbeatRequest(kafkaHeaderAndRequest); - break; - case LEAVE_GROUP: - responseFuture = handleLeaveGroupRequest(kafkaHeaderAndRequest); - break; - case DESCRIBE_GROUPS: - responseFuture = handleDescribeGroupRequest(kafkaHeaderAndRequest); - break; - case LIST_GROUPS: - responseFuture = handleListGroupsRequest(kafkaHeaderAndRequest); - break; - case DELETE_GROUPS: - responseFuture = handleDeleteGroupsRequest(kafkaHeaderAndRequest); - break; - case SASL_HANDSHAKE: - responseFuture = handleSaslHandshake(kafkaHeaderAndRequest); - break; - case SASL_AUTHENTICATE: - responseFuture = handleSaslAuthenticate(kafkaHeaderAndRequest); - break; - default: - responseFuture = handleError(kafkaHeaderAndRequest); + if (!isActive.get()) { + responseFuture = handleInactive(kafkaHeaderAndRequest); + } else { + switch (kafkaHeaderAndRequest.getHeader().apiKey()) { + case API_VERSIONS: + responseFuture = handleApiVersionsRequest(kafkaHeaderAndRequest); + break; + case METADATA: + responseFuture = handleTopicMetadataRequest(kafkaHeaderAndRequest); + // this is special, wait Metadata command return, before execute other command? + // responseFuture.get(); + break; + case PRODUCE: + responseFuture = handleProduceRequest(kafkaHeaderAndRequest); + break; + case FIND_COORDINATOR: + responseFuture = handleFindCoordinatorRequest(kafkaHeaderAndRequest); + break; + case LIST_OFFSETS: + responseFuture = handleListOffsetRequest(kafkaHeaderAndRequest); + break; + case OFFSET_FETCH: + responseFuture = handleOffsetFetchRequest(kafkaHeaderAndRequest); + break; + case OFFSET_COMMIT: + responseFuture = handleOffsetCommitRequest(kafkaHeaderAndRequest); + break; + case FETCH: + responseFuture = handleFetchRequest(kafkaHeaderAndRequest); + break; + case JOIN_GROUP: + responseFuture = handleJoinGroupRequest(kafkaHeaderAndRequest); + break; + case SYNC_GROUP: + responseFuture = handleSyncGroupRequest(kafkaHeaderAndRequest); + break; + case HEARTBEAT: + responseFuture = handleHeartbeatRequest(kafkaHeaderAndRequest); + break; + case LEAVE_GROUP: + responseFuture = handleLeaveGroupRequest(kafkaHeaderAndRequest); + break; + case DESCRIBE_GROUPS: + responseFuture = handleDescribeGroupRequest(kafkaHeaderAndRequest); + break; + case LIST_GROUPS: + responseFuture = handleListGroupsRequest(kafkaHeaderAndRequest); + break; + case DELETE_GROUPS: + responseFuture = handleDeleteGroupsRequest(kafkaHeaderAndRequest); + break; + case SASL_HANDSHAKE: + responseFuture = handleSaslHandshake(kafkaHeaderAndRequest); + break; + case SASL_AUTHENTICATE: + responseFuture = handleSaslAuthenticate(kafkaHeaderAndRequest); + break; + default: + responseFuture = handleError(kafkaHeaderAndRequest); + } } responsesQueue.compute(channel, (key, queue) -> { if (queue == null) { - Queue> newQueue = Queues.newConcurrentLinkedQueue(); - newQueue.add(responseFuture); + Queue newQueue = Queues.newConcurrentLinkedQueue(); + newQueue.add(ResponseAndRequest.of(responseFuture, kafkaHeaderAndRequest)); return newQueue; } else { - queue.add(responseFuture); + queue.add(ResponseAndRequest.of(responseFuture, kafkaHeaderAndRequest)); return queue; } }); @@ -221,86 +220,118 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception } // Write and flush continuously completed request back through channel. - // This is to make sure request get response in the same order. + // This is to make sure request get responseFuture in the same order. protected void writeAndFlushResponseToClient(Channel channel) { - Queue> responseQueue = + Queue responseQueue = responsesQueue.get(channel); - // loop from first response. - while (responseQueue != null && responseQueue.peek() != null && responseQueue.peek().isDone()) { - CompletableFuture response = responseQueue.remove(); + // loop from first responseFuture. + while (responseQueue != null && responseQueue.peek() != null + && responseQueue.peek().getResponseFuture().isDone()) { + ResponseAndRequest response = responseQueue.remove(); try { - ResponseAndRequest pair = response.join(); if (log.isDebugEnabled()) { log.debug("Write kafka cmd response back to client. \n" + "\trequest content: {} \n" + "\tresponse content: {}", - pair.getRequest().toString(), - pair.getResponse().toString(pair.getRequest().getRequest().version())); + response.getRequest().toString(), + response.getResponseFuture().join().toString(response.getRequest().getRequest().version())); + log.debug("Write kafka cmd responseFuture back to client. request: {}", + response.getRequest().getHeader()); + } + + ByteBuf result = responseToByteBuf(response.getResponseFuture().get(), response.getRequest()); + channel.writeAndFlush(result); + } catch (Exception e) { + // should not comes here. + log.error("error to get Response ByteBuf:", e); + } + } + } + + // return all the current command before a channel close. return Error response for all pending request. + protected void writeAndFlushWhenInactiveChannel(Channel channel) { + Queue responseQueue = + responsesQueue.get(channel); + + // loop from first responseFuture, and return them all + while (responseQueue != null && responseQueue.peek() != null) { + try { + ResponseAndRequest pair = responseQueue.remove(); + + if (log.isDebugEnabled()) { + log.debug("Write kafka cmd responseFuture back to client. request: {}", + pair.getRequest().getHeader()); } + AbstractRequest request = pair.getRequest().getRequest(); + AbstractResponse apiResponse = request + .getErrorResponse(new LeaderNotAvailableException("Channel is closing!")); + pair.getResponseFuture().complete(apiResponse); - ByteBuf result = responseToByteBuf(pair.getResponse(), pair.getRequest()); + ByteBuf result = responseToByteBuf(pair.getResponseFuture().get(), pair.getRequest()); channel.writeAndFlush(result); } catch (Exception e) { // should not comes here. log.error("error to get Response ByteBuf:", e); - throw e; } } } - protected abstract CompletableFuture + protected abstract CompletableFuture handleError(KafkaHeaderAndRequest kafkaHeaderAndRequest); - protected abstract CompletableFuture + protected abstract CompletableFuture + handleInactive(KafkaHeaderAndRequest kafkaHeaderAndRequest); + + protected abstract CompletableFuture handleApiVersionsRequest(KafkaHeaderAndRequest apiVersion); - protected abstract CompletableFuture + protected abstract CompletableFuture handleTopicMetadataRequest(KafkaHeaderAndRequest metadata); - protected abstract CompletableFuture + protected abstract CompletableFuture handleProduceRequest(KafkaHeaderAndRequest produce); - protected abstract CompletableFuture + protected abstract CompletableFuture handleFindCoordinatorRequest(KafkaHeaderAndRequest findCoordinator); - protected abstract CompletableFuture + protected abstract CompletableFuture handleListOffsetRequest(KafkaHeaderAndRequest listOffset); - protected abstract CompletableFuture + protected abstract CompletableFuture handleOffsetFetchRequest(KafkaHeaderAndRequest offsetFetch); - protected abstract CompletableFuture + protected abstract CompletableFuture handleOffsetCommitRequest(KafkaHeaderAndRequest offsetCommit); - protected abstract CompletableFuture + protected abstract CompletableFuture handleFetchRequest(KafkaHeaderAndRequest fetch); - protected abstract CompletableFuture + protected abstract CompletableFuture handleJoinGroupRequest(KafkaHeaderAndRequest joinGroup); - protected abstract CompletableFuture + protected abstract CompletableFuture handleSyncGroupRequest(KafkaHeaderAndRequest syncGroup); - protected abstract CompletableFuture + protected abstract CompletableFuture handleHeartbeatRequest(KafkaHeaderAndRequest heartbeat); - protected abstract CompletableFuture + protected abstract CompletableFuture handleLeaveGroupRequest(KafkaHeaderAndRequest leaveGroup); - protected abstract CompletableFuture + protected abstract CompletableFuture handleDescribeGroupRequest(KafkaHeaderAndRequest describeGroup); - protected abstract CompletableFuture + protected abstract CompletableFuture handleListGroupsRequest(KafkaHeaderAndRequest listGroups); - protected abstract CompletableFuture + protected abstract CompletableFuture handleDeleteGroupsRequest(KafkaHeaderAndRequest deleteGroups); - protected abstract CompletableFuture + protected abstract CompletableFuture handleSaslAuthenticate(KafkaHeaderAndRequest kafkaHeaderAndRequest); - protected abstract CompletableFuture + protected abstract CompletableFuture handleSaslHandshake(KafkaHeaderAndRequest kafkaHeaderAndRequest); static class KafkaHeaderAndRequest implements Closeable { @@ -398,7 +429,7 @@ static KafkaHeaderAndResponse responseForRequest(KafkaHeaderAndRequest request, } public String toString() { - return String.format("KafkaHeaderAndResponse(header=%s,response=%s)", + return String.format("KafkaHeaderAndResponse(header=%s,responseFuture=%s)", this.header.toStruct().toString(), this.response.toString(this.getApiVersion())); } @@ -409,20 +440,21 @@ public void close() { } /** - * A class that stores Kafka request and its related response. + * A class that stores Kafka request and its related responseFuture. */ static class ResponseAndRequest { @Getter - private AbstractResponse response; + private CompletableFuture responseFuture; @Getter private KafkaHeaderAndRequest request; - public static ResponseAndRequest of(AbstractResponse response, KafkaHeaderAndRequest request) { + public static ResponseAndRequest of(CompletableFuture response, + KafkaHeaderAndRequest request) { return new ResponseAndRequest(response, request); } - ResponseAndRequest(AbstractResponse response, KafkaHeaderAndRequest request) { - this.response = response; + ResponseAndRequest(CompletableFuture response, KafkaHeaderAndRequest request) { + this.responseFuture = response; this.request = request; } } diff --git a/src/main/java/io/streamnative/kop/KafkaProtocolHandler.java b/src/main/java/io/streamnative/kop/KafkaProtocolHandler.java index cebdc9344f..3b82975271 100644 --- a/src/main/java/io/streamnative/kop/KafkaProtocolHandler.java +++ b/src/main/java/io/streamnative/kop/KafkaProtocolHandler.java @@ -170,8 +170,6 @@ public enum ListenerType { @Getter private BrokerService brokerService; @Getter - private KafkaTopicManager kafkaTopicManager; - @Getter private GroupCoordinator groupCoordinator; @Getter private String bindAddress; @@ -220,9 +218,6 @@ public void start(BrokerService service) { KopVersion.getBuildHost(), KopVersion.getBuildTime()); - // a topic Manager - kafkaTopicManager = new KafkaTopicManager(service); - // init and start group coordinator if (kafkaConfig.isEnableGroupCoordinator()) { try { @@ -239,13 +234,12 @@ public void start(BrokerService service) { } } - // this is called after initialize, and with kafkaTopicManager, kafkaConfig, brokerService all set. + // this is called after initialize, and with kafkaConfig, brokerService all set. @Override public Map> newChannelInitializers() { checkState(kafkaConfig != null); checkState(kafkaConfig.getListeners() != null); checkState(brokerService != null); - checkState(kafkaTopicManager != null); if (kafkaConfig.isEnableGroupCoordinator()) { checkState(groupCoordinator != null); } @@ -265,7 +259,6 @@ public Map> newChannelIniti new InetSocketAddress(brokerService.pulsar().getBindAddress(), getListenerPort(listener)), new KafkaChannelInitializer(brokerService.pulsar(), kafkaConfig, - kafkaTopicManager, groupCoordinator, false)); } else if (listener.startsWith(SSL_PREFIX)) { @@ -273,7 +266,6 @@ public Map> newChannelIniti new InetSocketAddress(brokerService.pulsar().getBindAddress(), getListenerPort(listener)), new KafkaChannelInitializer(brokerService.pulsar(), kafkaConfig, - kafkaTopicManager, groupCoordinator, true)); } else { @@ -455,7 +447,7 @@ public static int getListenerPort(String listeners, ListenerType type) { return -1; } - public static String getBrokerUrl(String listeners, Boolean tlsEnabled) { + public static String getKopBrokerUrl(String listeners, Boolean tlsEnabled) { String[] parts = listeners.split(LISTENER_DEL); for (String listener: parts) { diff --git a/src/main/java/io/streamnative/kop/KafkaRequestHandler.java b/src/main/java/io/streamnative/kop/KafkaRequestHandler.java index 80bb47aea9..d82280c2c3 100644 --- a/src/main/java/io/streamnative/kop/KafkaRequestHandler.java +++ b/src/main/java/io/streamnative/kop/KafkaRequestHandler.java @@ -17,7 +17,7 @@ import static com.google.common.base.Preconditions.checkState; import static io.streamnative.kop.KafkaProtocolHandler.ListenerType.PLAINTEXT; import static io.streamnative.kop.KafkaProtocolHandler.ListenerType.SSL; -import static io.streamnative.kop.KafkaProtocolHandler.getBrokerUrl; +import static io.streamnative.kop.KafkaProtocolHandler.getKopBrokerUrl; import static io.streamnative.kop.KafkaProtocolHandler.getListenerPort; import static io.streamnative.kop.MessagePublishContext.publishMessages; import static io.streamnative.kop.utils.TopicNameUtils.getKafkaTopicNameFromPulsarTopicname; @@ -39,6 +39,7 @@ import java.io.IOException; import java.net.InetSocketAddress; import java.net.URI; +import java.net.URISyntaxException; import java.nio.ByteBuffer; import java.util.Collections; import java.util.HashMap; @@ -65,10 +66,13 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.LeaderNotAvailableException; +import org.apache.kafka.common.internals.Topic; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.record.MemoryRecords; import org.apache.kafka.common.record.RecordBatch; import org.apache.kafka.common.record.Records; +import org.apache.kafka.common.requests.AbstractRequest; import org.apache.kafka.common.requests.AbstractResponse; import org.apache.kafka.common.requests.ApiVersionsResponse; import org.apache.kafka.common.requests.DeleteGroupsRequest; @@ -106,6 +110,7 @@ import org.apache.kafka.common.requests.SyncGroupRequest; import org.apache.kafka.common.requests.SyncGroupResponse; import org.apache.kafka.common.utils.Utils; +import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfigurationUtils; import org.apache.pulsar.broker.authentication.AuthenticationProvider; @@ -124,9 +129,11 @@ import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.schema.KeyValue; +import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.common.util.Murmur3_32Hash; import org.apache.pulsar.policies.data.loadbalancer.ServiceLookupData; import org.apache.pulsar.zookeeper.ZooKeeperCache; +import org.apache.pulsar.zookeeper.ZooKeeperCache.Deserializer; /** * This class contains all the request handling methods. @@ -152,15 +159,12 @@ public class KafkaRequestHandler extends KafkaCommandDecoder { public KafkaRequestHandler(PulsarService pulsarService, KafkaServiceConfiguration kafkaConfig, - KafkaTopicManager kafkaTopicManager, GroupCoordinator groupCoordinator, Boolean tlsEnabled) throws Exception { super(); this.pulsarService = pulsarService; this.kafkaConfig = kafkaConfig; - this.topicManager = kafkaTopicManager; this.groupCoordinator = groupCoordinator; - this.clusterName = kafkaConfig.getClusterName(); this.executor = pulsarService.getExecutor(); this.admin = pulsarService.getAdminClient(); @@ -170,39 +174,83 @@ public KafkaRequestHandler(PulsarService pulsarService, this.namespace = NamespaceName.get( kafkaConfig.getKafkaTenant(), kafkaConfig.getKafkaNamespace()); + this.topicManager = new KafkaTopicManager(this); + } + + @Override + public void channelActive(ChannelHandlerContext ctx) throws Exception { + super.channelActive(ctx); + getTopicManager().updateCtx(); + log.info("channel active: {}", ctx.channel()); + } + + @Override + public void channelInactive(ChannelHandlerContext ctx) throws Exception { + super.channelInactive(ctx); + log.info("channel inactive {}", ctx.channel()); + + close(); + isActive.set(false); + } + + @Override + protected void close() { + if (isActive.getAndSet(false)) { + log.info("close channel {}", ctx.channel()); + writeAndFlushWhenInactiveChannel(ctx.channel()); + ctx.close(); + topicManager.close(); + } } - protected CompletableFuture handleApiVersionsRequest(KafkaHeaderAndRequest apiVersionRequest) { - AbstractResponse apiResponse = ApiVersionsResponse.defaultApiVersionsResponse(); - CompletableFuture resultFuture = new CompletableFuture<>(); + protected CompletableFuture handleApiVersionsRequest(KafkaHeaderAndRequest apiVersionRequest) { + ApiVersionsResponse apiResponse = ApiVersionsResponse.defaultApiVersionsResponse(); + CompletableFuture resultFuture = new CompletableFuture<>(); - resultFuture.complete(ResponseAndRequest.of(apiResponse, apiVersionRequest)); + resultFuture.complete(apiResponse); return resultFuture; } - protected CompletableFuture handleError(KafkaHeaderAndRequest kafkaHeaderAndRequest) { - CompletableFuture resultFuture = new CompletableFuture<>(); + protected CompletableFuture handleError(KafkaHeaderAndRequest kafkaHeaderAndRequest) { + CompletableFuture resultFuture = new CompletableFuture<>(); String err = String.format("Kafka API (%s) Not supported by kop server.", kafkaHeaderAndRequest.getHeader().apiKey()); log.error(err); AbstractResponse apiResponse = kafkaHeaderAndRequest.getRequest() .getErrorResponse(new UnsupportedOperationException(err)); - resultFuture.complete(ResponseAndRequest.of(apiResponse, kafkaHeaderAndRequest)); + resultFuture.complete(apiResponse); return resultFuture; } + protected CompletableFuture handleInactive(KafkaHeaderAndRequest kafkaHeaderAndRequest) { + CompletableFuture resultFuture = new CompletableFuture<>(); + + AbstractRequest request = kafkaHeaderAndRequest.getRequest(); + AbstractResponse apiResponse = request.getErrorResponse(new LeaderNotAvailableException("Channel is closing!")); + + log.error("Kafka API {} is send to a closing channel", kafkaHeaderAndRequest.getHeader().apiKey()); + + resultFuture.complete(apiResponse); + return resultFuture; + } + // Leverage pulsar admin to get partitioned topic metadata private CompletableFuture getPartitionedTopicMetadataAsync(String topicName) { return admin.topics().getPartitionedTopicMetadataAsync(topicName); } - protected CompletableFuture handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar) { + protected CompletableFuture handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar) { checkArgument(metadataHar.getRequest() instanceof MetadataRequest); MetadataRequest metadataRequest = (MetadataRequest) metadataHar.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); + + if (log.isDebugEnabled()) { + log.debug("[{}] Request {}: for topic {} ", + ctx.channel(), metadataHar.getHeader(), metadataRequest.topics()); + } // Command response for all topics List allTopicMetadata = Collections.synchronizedList(Lists.newArrayList()); @@ -323,7 +371,7 @@ protected CompletableFuture handleTopicMetadataRequest(Kafka clusterName, MetadataResponse.NO_CONTROLLER_ID, Collections.emptyList()); - resultFuture.complete(ResponseAndRequest.of(finalResponse, metadataHar)); + resultFuture.complete(finalResponse); return; } @@ -338,7 +386,7 @@ protected CompletableFuture handleTopicMetadataRequest(Kafka clusterName, MetadataResponse.NO_CONTROLLER_ID, allTopicMetadata); - resultFuture.complete(ResponseAndRequest.of(finalResponse, metadataHar)); + resultFuture.complete(finalResponse); return; } @@ -357,8 +405,10 @@ protected CompletableFuture handleTopicMetadataRequest(Kafka partitionMetadatas.add(newFailedPartitionMetadata(topicName)); } else { Node newNode = partitionMetadata.leader(); - if (!allNodes.stream().anyMatch(node1 -> node1.equals(newNode))) { - allNodes.add(newNode); + synchronized (allNodes) { + if (!allNodes.stream().anyMatch(node1 -> node1.equals(newNode))) { + allNodes.add(newNode); + } } partitionMetadatas.add(partitionMetadata); } @@ -381,19 +431,14 @@ protected CompletableFuture handleTopicMetadataRequest(Kafka // whether completed all the topics requests. int finishedTopics = topicsCompleted.incrementAndGet(); - if (log.isTraceEnabled()) { - log.trace("[{}] Request {}: Completed findBroker for topic {}, " + if (log.isDebugEnabled()) { + log.debug("[{}] Request {}: Completed findBroker for topic {}, " + "partitions found/all: {}/{}. \n dump All Metadata:", ctx.channel(), metadataHar.getHeader(), topic, finishedTopics, topicsNumber); allTopicMetadata.stream() - .forEach(data -> { - log.trace("topicMetadata: {}", data.toString()); - data.partitionMetadata() - .forEach(partitionData -> - log.trace(" partitionMetadata: {}", data.toString())); - }); + .forEach(data -> log.debug("TopicMetadata response: {}", data.toString())); } if (finishedTopics == topicsNumber) { // TODO: confirm right value for controller_id @@ -403,7 +448,7 @@ protected CompletableFuture handleTopicMetadataRequest(Kafka clusterName, MetadataResponse.NO_CONTROLLER_ID, allTopicMetadata); - resultFuture.complete(ResponseAndRequest.of(finalResponse, metadataHar)); + resultFuture.complete(finalResponse); } } }))); @@ -413,17 +458,16 @@ protected CompletableFuture handleTopicMetadataRequest(Kafka return resultFuture; } - protected CompletableFuture handleProduceRequest(KafkaHeaderAndRequest produceHar) { + protected CompletableFuture handleProduceRequest(KafkaHeaderAndRequest produceHar) { checkArgument(produceHar.getRequest() instanceof ProduceRequest); ProduceRequest produceRequest = (ProduceRequest) produceHar.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); if (produceRequest.transactionalId() != null) { log.warn("[{}] Transactions not supported", ctx.channel()); - resultFuture.complete(ResponseAndRequest.of( - failedResponse(produceHar, new UnsupportedOperationException("No transaction support")), - produceHar)); + resultFuture.complete( + failedResponse(produceHar, new UnsupportedOperationException("No transaction support"))); return resultFuture; } @@ -449,24 +493,17 @@ protected CompletableFuture handleProduceRequest(KafkaHeader TopicName topicName = pulsarTopicName(topicPartition, namespace); - pulsarService.getBrokerService().getTopic(topicName.toString(), true) - .whenComplete((topicOpt, exception) -> { - if (exception != null) { - log.error("[{}] Request {}: Failed to getOrCreateTopic {}. exception:", - ctx.channel(), produceHar.getHeader(), topicName, exception); - partitionResponse.complete(new PartitionResponse(Errors.KAFKA_STORAGE_ERROR)); - } else { - if (topicOpt.isPresent()) { - // TODO: need to add a produce Manager - topicManager.addTopic(topicName.toString(), (PersistentTopic) topicOpt.get()); - publishMessages((MemoryRecords) entry.getValue(), topicOpt.get(), partitionResponse); - } else { - log.error("[{}] Request {}: getOrCreateTopic get empty topic for name {}", - ctx.channel(), produceHar.getHeader(), topicName); - partitionResponse.complete(new PartitionResponse(Errors.KAFKA_STORAGE_ERROR)); - } - } - }); + topicManager.getTopic(topicName.toString(), true).whenComplete((persistentTopic, exception) -> { + if (exception != null || persistentTopic == null) { + log.error("[{}] Request {}: Failed to getOrCreateTopic {}. exception:", + ctx.channel(), produceHar.getHeader(), topicName, exception); + partitionResponse.complete(new PartitionResponse(Errors.KAFKA_STORAGE_ERROR)); + } else { + CompletableFuture topicFuture = new CompletableFuture<>(); + topicFuture.complete(persistentTopic); + publishMessages((MemoryRecords) entry.getValue(), persistentTopic, partitionResponse); + } + }); } CompletableFuture.allOf(responsesFutures.values().toArray(new CompletableFuture[responsesSize])) @@ -482,16 +519,16 @@ protected CompletableFuture handleProduceRequest(KafkaHeader log.debug("[{}] Request {}: Complete handle produce.", ctx.channel(), produceHar.toString()); } - resultFuture.complete(ResponseAndRequest.of(new ProduceResponse(responses), produceHar)); + resultFuture.complete(new ProduceResponse(responses)); }); return resultFuture; } - protected CompletableFuture + protected CompletableFuture handleFindCoordinatorRequest(KafkaHeaderAndRequest findCoordinator) { checkArgument(findCoordinator.getRequest() instanceof FindCoordinatorRequest); FindCoordinatorRequest request = (FindCoordinatorRequest) findCoordinator.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); if (request.coordinatorType() == FindCoordinatorRequest.CoordinatorType.GROUP) { int partition = groupCoordinator.partitionFor(request.coordinatorKey()); @@ -514,7 +551,7 @@ protected CompletableFuture handleProduceRequest(KafkaHeader AbstractResponse response = new FindCoordinatorResponse( Errors.NONE, node); - resultFuture.complete(ResponseAndRequest.of(response, findCoordinator)); + resultFuture.complete(response); }); } else { throw new NotImplementedException("FindCoordinatorRequest not support TRANSACTION type " @@ -524,12 +561,12 @@ protected CompletableFuture handleProduceRequest(KafkaHeader return resultFuture; } - protected CompletableFuture handleOffsetFetchRequest(KafkaHeaderAndRequest offsetFetch) { + protected CompletableFuture handleOffsetFetchRequest(KafkaHeaderAndRequest offsetFetch) { checkArgument(offsetFetch.getRequest() instanceof OffsetFetchRequest); OffsetFetchRequest request = (OffsetFetchRequest) offsetFetch.getRequest(); checkState(groupCoordinator != null, "Group Coordinator not started"); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); KeyValue> keyValue = groupCoordinator.handleFetchOffsets( @@ -537,25 +574,33 @@ protected CompletableFuture handleOffsetFetchRequest(KafkaHe Optional.of(request.partitions()) ); - resultFuture.complete(ResponseAndRequest - .of(new OffsetFetchResponse(keyValue.getKey(), keyValue.getValue()), offsetFetch)); + resultFuture.complete(new OffsetFetchResponse(keyValue.getKey(), keyValue.getValue())); return resultFuture; } private CompletableFuture - fetchOffsetForTimestamp(PersistentTopic persistentTopic, Long timestamp) { - ManagedLedgerImpl managedLedger = - (ManagedLedgerImpl) persistentTopic.getManagedLedger(); - + fetchOffsetForTimestamp(CompletableFuture persistentTopic, Long timestamp) { CompletableFuture partitionData = new CompletableFuture<>(); - try { + persistentTopic.whenComplete((perTopic, t) -> { + if (t != null || perTopic == null) { + log.error("Failed while get persistentTopic topic: {} ts: {}. ", + perTopic == null ? "null" : perTopic.getName(), timestamp, t); + + partitionData.complete(new ListOffsetResponse.PartitionData( + Errors.UNKNOWN_SERVER_ERROR, + ListOffsetResponse.UNKNOWN_TIMESTAMP, + ListOffsetResponse.UNKNOWN_OFFSET)); + return; + } + + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) perTopic.getManagedLedger(); if (timestamp == ListOffsetRequest.LATEST_TIMESTAMP) { PositionImpl position = (PositionImpl) managedLedger.getLastConfirmedEntry(); if (log.isDebugEnabled()) { log.debug("Get latest position for topic {} time {}. result: {}", - persistentTopic.getName(), timestamp, position); + perTopic.getName(), timestamp, position); } // no entry in ledger, then entry id could be -1 @@ -571,7 +616,7 @@ protected CompletableFuture handleOffsetFetchRequest(KafkaHe if (log.isDebugEnabled()) { log.debug("Get earliest position for topic {} time {}. result: {}", - persistentTopic.getName(), timestamp, position); + perTopic.getName(), timestamp, position); } partitionData.complete(new ListOffsetResponse.PartitionData( @@ -590,7 +635,7 @@ public void findEntryComplete(Position position, Object ctx) { finalPosition = OffsetFinder.getFirstValidPosition(managedLedger); if (finalPosition == null) { log.warn("Unable to find position for topic {} time {}. get NULL position", - persistentTopic.getName(), timestamp); + perTopic.getName(), timestamp); partitionData.complete(new ListOffsetResponse .PartitionData( @@ -605,7 +650,7 @@ public void findEntryComplete(Position position, Object ctx) { if (log.isDebugEnabled()) { log.debug("Find position for topic {} time {}. position: {}", - persistentTopic.getName(), timestamp, finalPosition); + perTopic.getName(), timestamp, finalPosition); } partitionData.complete(new ListOffsetResponse.PartitionData( Errors.NONE, @@ -617,7 +662,7 @@ public void findEntryComplete(Position position, Object ctx) { public void findEntryFailed(ManagedLedgerException exception, Optional position, Object ctx) { log.warn("Unable to find position for topic {} time {}. Exception:", - persistentTopic.getName(), timestamp, exception); + perTopic.getName(), timestamp, exception); partitionData.complete(new ListOffsetResponse .PartitionData( Errors.UNKNOWN_SERVER_ERROR, @@ -627,23 +672,15 @@ public void findEntryFailed(ManagedLedgerException exception, } }); } - } catch (Exception e) { - log.error("Failed while get position for topic: {} ts: {}.", - persistentTopic.getName(), timestamp, e); - - partitionData.complete(new ListOffsetResponse.PartitionData( - Errors.UNKNOWN_SERVER_ERROR, - ListOffsetResponse.UNKNOWN_TIMESTAMP, - ListOffsetResponse.UNKNOWN_OFFSET)); - } + }); return partitionData; } - private CompletableFuture handleListOffsetRequestV1AndAbove(KafkaHeaderAndRequest listOffset) { + private CompletableFuture handleListOffsetRequestV1AndAbove(KafkaHeaderAndRequest listOffset) { ListOffsetRequest request = (ListOffsetRequest) listOffset.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); Map> responseData = Maps.newHashMap(); request.partitionTimestamps().entrySet().stream().forEach(tms -> { @@ -652,19 +689,8 @@ private CompletableFuture handleListOffsetRequestV1AndAbove( Long times = tms.getValue(); CompletableFuture partitionData; - // topic not exist, return UNKNOWN_TOPIC_OR_PARTITION - if (!topicManager.topicExists(pulsarTopic.toString())) { - log.warn("Topic {} not exist in topic manager while list offset.", pulsarTopic.toString()); - partitionData = new CompletableFuture<>(); - partitionData.complete(new ListOffsetResponse - .PartitionData( - Errors.UNKNOWN_TOPIC_OR_PARTITION, - ListOffsetResponse.UNKNOWN_TIMESTAMP, - ListOffsetResponse.UNKNOWN_OFFSET)); - } else { - PersistentTopic persistentTopic = topicManager.getTopic(pulsarTopic.toString()); - partitionData = fetchOffsetForTimestamp(persistentTopic, times); - } + CompletableFuture persistentTopic = topicManager.getTopic(pulsarTopic.toString()); + partitionData = fetchOffsetForTimestamp(persistentTopic, times); responseData.put(topic, partitionData); }); @@ -675,20 +701,19 @@ private CompletableFuture handleListOffsetRequestV1AndAbove( ListOffsetResponse response = new ListOffsetResponse(CoreUtils.mapValue(responseData, future -> future.join())); - resultFuture.complete(ResponseAndRequest - .of(response, listOffset)); + resultFuture.complete(response); }); return resultFuture; } // get offset from underline managedLedger - protected CompletableFuture handleListOffsetRequest(KafkaHeaderAndRequest listOffset) { + protected CompletableFuture handleListOffsetRequest(KafkaHeaderAndRequest listOffset) { checkArgument(listOffset.getRequest() instanceof ListOffsetRequest); // not support version 0 if (listOffset.getHeader().apiVersion() == 0) { - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); ListOffsetRequest request = (ListOffsetRequest) listOffset.getRequest(); log.error("ListOffset not support V0 format request"); @@ -697,8 +722,7 @@ protected CompletableFuture handleListOffsetRequest(KafkaHea ignored -> new ListOffsetResponse .PartitionData(Errors.UNSUPPORTED_FOR_MESSAGE_FORMAT, Lists.newArrayList()))); - resultFuture.complete(ResponseAndRequest - .of(response, listOffset)); + resultFuture.complete(response); return resultFuture; } @@ -722,13 +746,13 @@ private Map nonExistingTopicErrors(OffsetCommitRequest r // )); } - protected CompletableFuture handleOffsetCommitRequest(KafkaHeaderAndRequest offsetCommit) { + protected CompletableFuture handleOffsetCommitRequest(KafkaHeaderAndRequest offsetCommit) { checkArgument(offsetCommit.getRequest() instanceof OffsetCommitRequest); checkState(groupCoordinator != null, "Group Coordinator not started"); OffsetCommitRequest request = (OffsetCommitRequest) offsetCommit.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); Map nonExistingTopic = nonExistingTopicErrors(request); @@ -748,16 +772,16 @@ protected CompletableFuture handleOffsetCommitRequest(KafkaH offsetCommitResult.putAll(nonExistingTopic); } OffsetCommitResponse response = new OffsetCommitResponse(offsetCommitResult); - resultFuture.complete(ResponseAndRequest.of(response, offsetCommit)); + resultFuture.complete(response); }); return resultFuture; } - protected CompletableFuture handleFetchRequest(KafkaHeaderAndRequest fetch) { + protected CompletableFuture handleFetchRequest(KafkaHeaderAndRequest fetch) { checkArgument(fetch.getRequest() instanceof FetchRequest); FetchRequest request = (FetchRequest) fetch.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); if (log.isDebugEnabled()) { log.debug("[{}] Request {} Fetch request. Size: {}. Each item: ", @@ -773,13 +797,13 @@ protected CompletableFuture handleFetchRequest(KafkaHeaderAn return fetchContext.handleFetch(resultFuture); } - protected CompletableFuture handleJoinGroupRequest(KafkaHeaderAndRequest joinGroup) { + protected CompletableFuture handleJoinGroupRequest(KafkaHeaderAndRequest joinGroup) { checkArgument(joinGroup.getRequest() instanceof JoinGroupRequest); checkState(groupCoordinator != null, "Group Coordinator not started"); JoinGroupRequest request = (JoinGroupRequest) joinGroup.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); Map protocols = new HashMap<>(); request.groupProtocols() @@ -814,16 +838,16 @@ protected CompletableFuture handleJoinGroupRequest(KafkaHead response, joinGroup.getHeader().correlationId(), joinGroup.getHeader().clientId()); } - resultFuture.complete(ResponseAndRequest.of(response, joinGroup)); + resultFuture.complete(response); }); return resultFuture; } - protected CompletableFuture handleSyncGroupRequest(KafkaHeaderAndRequest syncGroup) { + protected CompletableFuture handleSyncGroupRequest(KafkaHeaderAndRequest syncGroup) { checkArgument(syncGroup.getRequest() instanceof SyncGroupRequest); SyncGroupRequest request = (SyncGroupRequest) syncGroup.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); groupCoordinator.handleSyncGroup( request.groupId(), @@ -838,17 +862,17 @@ protected CompletableFuture handleSyncGroupRequest(KafkaHead ByteBuffer.wrap(syncGroupResult.getValue()) ); - resultFuture.complete(ResponseAndRequest.of(response, syncGroup)); + resultFuture.complete(response); }); return resultFuture; } - protected CompletableFuture handleHeartbeatRequest(KafkaHeaderAndRequest heartbeat) { + protected CompletableFuture handleHeartbeatRequest(KafkaHeaderAndRequest heartbeat) { checkArgument(heartbeat.getRequest() instanceof HeartbeatRequest); HeartbeatRequest request = (HeartbeatRequest) heartbeat.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); // let the coordinator to handle heartbeat groupCoordinator.handleHeartbeat( @@ -863,16 +887,16 @@ protected CompletableFuture handleHeartbeatRequest(KafkaHead response, heartbeat.getHeader().correlationId(), heartbeat.getHeader().clientId()); } - resultFuture.complete(ResponseAndRequest.of(response, heartbeat)); + resultFuture.complete(response); }); return resultFuture; } @Override - protected CompletableFuture handleLeaveGroupRequest(KafkaHeaderAndRequest leaveGroup) { + protected CompletableFuture handleLeaveGroupRequest(KafkaHeaderAndRequest leaveGroup) { checkArgument(leaveGroup.getRequest() instanceof LeaveGroupRequest); LeaveGroupRequest request = (LeaveGroupRequest) leaveGroup.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); // let the coordinator to handle heartbeat groupCoordinator.handleLeaveGroup( @@ -881,17 +905,17 @@ protected CompletableFuture handleLeaveGroupRequest(KafkaHea ).thenAccept(errors -> { LeaveGroupResponse response = new LeaveGroupResponse(errors); - resultFuture.complete(ResponseAndRequest.of(response, leaveGroup)); + resultFuture.complete(response); }); return resultFuture; } @Override - protected CompletableFuture handleDescribeGroupRequest(KafkaHeaderAndRequest describeGroup) { + protected CompletableFuture handleDescribeGroupRequest(KafkaHeaderAndRequest describeGroup) { checkArgument(describeGroup.getRequest() instanceof DescribeGroupsRequest); DescribeGroupsRequest request = (DescribeGroupsRequest) describeGroup.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); // let the coordinator to handle heartbeat Map groups = request.groupIds().stream() @@ -930,35 +954,35 @@ protected CompletableFuture handleDescribeGroupRequest(Kafka DescribeGroupsResponse response = new DescribeGroupsResponse( groups ); - resultFuture.complete(ResponseAndRequest.of(response, describeGroup)); + resultFuture.complete(response); return resultFuture; } @Override - protected CompletableFuture handleListGroupsRequest(KafkaHeaderAndRequest listGroups) { + protected CompletableFuture handleListGroupsRequest(KafkaHeaderAndRequest listGroups) { throw new NotImplementedException("Not implemented yet"); } @Override - protected CompletableFuture handleDeleteGroupsRequest(KafkaHeaderAndRequest deleteGroups) { + protected CompletableFuture handleDeleteGroupsRequest(KafkaHeaderAndRequest deleteGroups) { checkArgument(deleteGroups.getRequest() instanceof DescribeGroupsRequest); DeleteGroupsRequest request = (DeleteGroupsRequest) deleteGroups.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); Map deleteResult = groupCoordinator.handleDeleteGroups(request.groups()); DeleteGroupsResponse response = new DeleteGroupsResponse( deleteResult ); - resultFuture.complete(ResponseAndRequest.of(response, deleteGroups)); + resultFuture.complete(response); return resultFuture; } @Override - protected CompletableFuture handleSaslAuthenticate(KafkaHeaderAndRequest saslAuthenticate) { + protected CompletableFuture handleSaslAuthenticate(KafkaHeaderAndRequest saslAuthenticate) { checkArgument(saslAuthenticate.getRequest() instanceof SaslAuthenticateRequest); SaslAuthenticateRequest request = (SaslAuthenticateRequest) saslAuthenticate.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); SaslAuth saslAuth; try { @@ -990,24 +1014,24 @@ protected CompletableFuture handleSaslAuthenticate(KafkaHead // TODO: what should be answered? SaslAuthenticateResponse response = new SaslAuthenticateResponse( Errors.NONE, "", request.saslAuthBytes()); - resultFuture.complete(ResponseAndRequest.of(response, saslAuthenticate)); + resultFuture.complete(response); } catch (IOException | AuthenticationException | PulsarAdminException e) { SaslAuthenticateResponse response = new SaslAuthenticateResponse( Errors.SASL_AUTHENTICATION_FAILED, e.getMessage(), request.saslAuthBytes()); - resultFuture.complete(ResponseAndRequest.of(response, saslAuthenticate)); + resultFuture.complete(response); } return resultFuture; } @Override - protected CompletableFuture handleSaslHandshake(KafkaHeaderAndRequest saslHandshake) { + protected CompletableFuture handleSaslHandshake(KafkaHeaderAndRequest saslHandshake) { checkArgument(saslHandshake.getRequest() instanceof SaslHandshakeRequest); SaslHandshakeRequest request = (SaslHandshakeRequest) saslHandshake.getRequest(); - CompletableFuture resultFuture = new CompletableFuture<>(); + CompletableFuture resultFuture = new CompletableFuture<>(); SaslHandshakeResponse response = checkSaslMechanism(request.mechanism()); - resultFuture.complete(ResponseAndRequest.of(response, saslHandshake)); + resultFuture.complete(response); return resultFuture; } @@ -1062,114 +1086,135 @@ public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { return; } - AtomicInteger atomicInteger = new AtomicInteger(matchBrokers.size()); - matchBrokers.stream().forEach(matchBroker -> { - String path = String.format("%s/%s", LoadManager.LOADBALANCE_BROKERS_ROOT, - matchBroker); - zkCache.getDataAsync(path, pulsarService.getLoadManager().get().getLoadReportDeserializer()) - .whenComplete((serviceLookupData, th) -> { - int wait = atomicInteger.decrementAndGet(); + // Get a list of ServiceLookupData for each matchBroker. + List>> list = matchBrokers.stream() + .map(matchBroker -> + zkCache.getDataAsync( + String.format("%s/%s", LoadManager.LOADBALANCE_BROKERS_ROOT, matchBroker), + (Deserializer) + pulsarService.getLoadManager().get().getLoadReportDeserializer())) + .collect(Collectors.toList()); + + FutureUtil.waitForAll(list) + .whenComplete((ignore, th) -> { if (th != null) { - log.error("Error in getDataAsync({}) for {}", path, brokerAddress, th); + log.error("Error in getDataAsync() for {}", brokerAddress, th); returnFuture.complete(Optional.empty()); return; } - if (log.isDebugEnabled()) { - log.debug("Handle getProtocolDataToAdvertise for {}, pulsarUrl: {}, kafka: {}", - topic, serviceLookupData.get().getPulsarServiceUrl(), - serviceLookupData.get().getProtocol(KafkaProtocolHandler.PROTOCOL_NAME)); - } + try { + for (CompletableFuture> lookupData : list) { + ServiceLookupData data = lookupData.get().get(); + if (log.isDebugEnabled()) { + log.debug("Handle getProtocolDataToAdvertise for {}, pulsarUrl: {}, " + + "pulsarUrlTls: {}, webUrl: {}, webUrlTls: {} kafka: {}", + topic, data.getPulsarServiceUrl(), data.getPulsarServiceUrlTls(), + data.getWebServiceUrl(), data.getWebServiceUrlTls(), + data.getProtocol(KafkaProtocolHandler.PROTOCOL_NAME)); + } - ServiceLookupData data = serviceLookupData.get(); - if (data.getPulsarServiceUrl().contains(hostAndPort) - || data.getPulsarServiceUrlTls().contains(hostAndPort) - || data.getWebServiceUrl().contains(hostAndPort) - || data.getWebServiceUrlTls().contains(hostAndPort)) { - returnFuture.complete(data.getProtocol(KafkaProtocolHandler.PROTOCOL_NAME)); + if (lookupDataContainsAddress(data, hostAndPort)) { + returnFuture.complete(data.getProtocol(KafkaProtocolHandler.PROTOCOL_NAME)); + return; + } + } + } catch (Exception e) { + log.error("Error in {} lookupFuture get: ", brokerAddress, e); + returnFuture.complete(Optional.empty()); return; } - if (wait == 0) { - log.error("Error to search {} in all child of zk://loadbalance", brokerAddress); - returnFuture.complete(Optional.empty()); - } - }); - }); + // no matching lookup data in all matchBrokers. + log.error("Not able to search {} in all child of zk://loadbalance", brokerAddress); + returnFuture.complete(Optional.empty()); + } + ); }); return returnFuture; } + private boolean isOffsetTopic(String topic) { + String offsetsTopic = kafkaConfig.getKafkaMetadataTenant() + "/" + + kafkaConfig.getKafkaMetadataNamespace() + + "/" + Topic.GROUP_METADATA_TOPIC_NAME; + + return topic.contains(offsetsTopic); + } + private CompletableFuture findBroker(PulsarService pulsarService, TopicName topic) { if (log.isDebugEnabled()) { - log.debug("Handle Lookup for {}", topic); + log.debug("[{}] Handle Lookup for {}", ctx.channel(), topic); } + CompletableFuture returnFuture = new CompletableFuture<>(); + PulsarClientImpl pulsarClient; try { - PulsarClientImpl pulsarClient = (PulsarClientImpl) pulsarService.getClient(); - return pulsarClient.getLookup() - .getBroker(topic) - .thenCompose(pair -> getProtocolDataToAdvertise(pair, topic)) - .thenApply(stringOptional -> { - if (!stringOptional.isPresent()) { - log.error("Not get advertise data for Kafka topic:{} ", topic); - return null; - } + pulsarClient = (PulsarClientImpl) pulsarService.getClient(); + } catch (PulsarServerException e) { + log.error("[{}] findBroker for Kafka topic {} error get pulsar client. throwable: ", + topic, topic, e); + returnFuture.complete(null); + return returnFuture; + } - try { - String listeners = stringOptional.get(); - String brokerUrl = getBrokerUrl(listeners, tlsEnabled); + pulsarClient.getLookup() + .getBroker(topic) + .thenCompose(pair -> getProtocolDataToAdvertise(pair, topic)) + .whenComplete((stringOptional, throwable) -> { + if (!stringOptional.isPresent() || throwable != null) { + log.error("Not get advertise data for Kafka topic:{}. throwable", + topic, throwable); + returnFuture.complete(null); + return; + } - // get local listeners. - String localListeners = kafkaConfig.getListeners(); + String listeners = stringOptional.get(); + String kopBrokerUrl = getKopBrokerUrl(listeners, tlsEnabled); + URI kopUri; + try { + kopUri = new URI(kopBrokerUrl); + } catch (URISyntaxException e) { + log.error("[{}] findBroker for topic {}: Failed to translate URI {}. exception:", + ctx.channel(), topic.toString(), kopBrokerUrl, e); + returnFuture.complete(null); + return; + } - if (log.isDebugEnabled()) { - log.debug("Found broker listeners: {} for topicName: {}, " - + "localListeners: {}, found Listeners: {}", - listeners, topic, localListeners, listeners); - } + Node node = newNode(new InetSocketAddress( + kopUri.getHost(), + kopUri.getPort())); - if (!topicManager.topicExists(topic.toString()) && localListeners.contains(brokerUrl)) { - pulsarService.getBrokerService().getTopic(topic.toString(), true) - .whenComplete((topicOpt, exception) -> { - if (exception != null) { - log.error("[{}] findBroker: Failed to getOrCreateTopic {}. exception:", - ctx.channel(), topic.toString(), exception); - } else { - if (topicOpt.isPresent()) { - if (log.isDebugEnabled()) { - log.debug("Add topic: {} into TopicManager while findBroker.", - topic.toString()); - } - topicManager.addTopic(topic.toString(), (PersistentTopic) topicOpt.get()); - } else { - log.error("[{}] findBroker: getOrCreateTopic get empty topic for name {}", - ctx.channel(), topic.toString()); - } - } - }); - } + // get local listeners. + String localListeners = kafkaConfig.getListeners(); - URI uri = new URI(brokerUrl); - Node node = newNode(new InetSocketAddress( - uri.getHost(), - uri.getPort())); + if (log.isDebugEnabled()) { + log.debug("Found broker listeners: {} for topicName: {}, " + + "localListeners: {}, found Listeners: {}", + listeners, topic, localListeners, listeners); + } - return newPartitionMetadata(topic, node); - } catch (Exception e) { - log.error("Caught error while find Broker for topic:{} ", topic, e); - return null; - } - }).exceptionally(ex -> { - log.error("Exceptionally while find Broker for topic:{} ", topic, ex); - return null; - }); - } catch (Exception e) { - log.error("Exceptionally while get pulsar client from Pulsar Broker for topic:{} ", topic, e); - CompletableFuture completableFuture = new CompletableFuture(); - completableFuture.complete(null); - return completableFuture; - } + if (!topicManager.topicExists(topic.toString()) + && !isOffsetTopic(topic.toString()) + && localListeners.contains(kopBrokerUrl)) { + topicManager.getTopic(topic.toString()).whenComplete((persistentTopic, exception) -> { + if (exception != null || persistentTopic == null) { + log.error("[{}] findBroker: Failed to getOrCreateTopic {}. exception:", + ctx.channel(), topic.toString(), exception); + returnFuture.complete(null); + } else { + if (log.isDebugEnabled()) { + log.debug("Add topic: {} into TopicManager while findBroker.", + topic.toString()); + } + returnFuture.complete(newPartitionMetadata(topic, node)); + } + }); + } else { + returnFuture.complete(newPartitionMetadata(topic, node)); + } + }); + return returnFuture; } static Node newNode(InetSocketAddress address) { @@ -1239,4 +1284,12 @@ static AbstractResponse failedResponse(KafkaHeaderAndRequest requestHar, Throwab } return requestHar.getRequest().getErrorResponse(((Integer) THROTTLE_TIME_MS.defaultValue), e); } + + // whether a ServiceLookupData contains wanted address. + static boolean lookupDataContainsAddress(ServiceLookupData data, String hostAndPort) { + return (data.getPulsarServiceUrl() != null && data.getPulsarServiceUrl().contains(hostAndPort)) + || (data.getPulsarServiceUrlTls() != null && data.getPulsarServiceUrlTls().contains(hostAndPort)) + || (data.getWebServiceUrl() != null && data.getWebServiceUrl().contains(hostAndPort)) + || (data.getWebServiceUrlTls() != null && data.getWebServiceUrlTls().contains(hostAndPort)); + } } diff --git a/src/main/java/io/streamnative/kop/KafkaService.java b/src/main/java/io/streamnative/kop/KafkaService.java index a615316301..9fee2cd176 100644 --- a/src/main/java/io/streamnative/kop/KafkaService.java +++ b/src/main/java/io/streamnative/kop/KafkaService.java @@ -50,8 +50,6 @@ public class KafkaService extends PulsarService { @Getter private final KafkaServiceConfiguration kafkaConfig; @Getter - private KafkaTopicManager kafkaTopicManager; - @Getter @Setter private GroupCoordinator groupCoordinator; @@ -210,7 +208,6 @@ public Boolean get() { .build(); getBrokerService().startProtocolHandlers(protocolHandlers); - this.kafkaTopicManager = kafkaProtocolHandler.getKafkaTopicManager(); this.groupCoordinator = kafkaProtocolHandler.getGroupCoordinator(); setState(State.Started); @@ -231,9 +228,6 @@ public void close() throws PulsarServerException { if (groupCoordinator != null) { this.groupCoordinator.shutdown(); } - if (kafkaTopicManager != null) { - this.kafkaTopicManager.close(); - } super.close(); } diff --git a/src/main/java/io/streamnative/kop/KafkaTopicConsumerManager.java b/src/main/java/io/streamnative/kop/KafkaTopicConsumerManager.java index 97d7b241a0..550124a006 100644 --- a/src/main/java/io/streamnative/kop/KafkaTopicConsumerManager.java +++ b/src/main/java/io/streamnative/kop/KafkaTopicConsumerManager.java @@ -54,7 +54,8 @@ public CompletableFuture> remove(long offset) { CompletableFuture> cursor = consumers.remove(offset); if (cursor != null) { if (log.isDebugEnabled()) { - log.debug("Get cursor for offset: {} in cache", offset); + log.debug("Get cursor for offset: {} in cache. cache size: {}", + offset, consumers.size()); } return cursor; } diff --git a/src/main/java/io/streamnative/kop/KafkaTopicManager.java b/src/main/java/io/streamnative/kop/KafkaTopicManager.java index 54c9f0ade5..14091f75ad 100644 --- a/src/main/java/io/streamnative/kop/KafkaTopicManager.java +++ b/src/main/java/io/streamnative/kop/KafkaTopicManager.java @@ -13,81 +13,234 @@ */ package io.streamnative.kop; +import java.net.InetSocketAddress; +import java.util.Map; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; import lombok.Getter; import lombok.extern.slf4j.Slf4j; -import org.apache.bookkeeper.util.collections.ConcurrentOpenHashMap; +import org.apache.commons.lang3.tuple.Pair; +import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.service.BrokerService; +import org.apache.pulsar.broker.service.Producer; import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.client.impl.PulsarClientImpl; +import org.apache.pulsar.common.naming.TopicName; /** * KafkaTopicManager manages a Map of topic to KafkaTopicConsumerManager. * For each topic, there is a KafkaTopicConsumerManager, which manages a topic and its related offset cursor. + * This is mainly used to cache the produce/consume topic, not include offsetTopic. */ @Slf4j public class KafkaTopicManager { - private final BrokerService service; + private final KafkaRequestHandler requestHandler; + private final PulsarService pulsarService; + private final BrokerService brokerService; - // consumerTopics for consumers cache. + // consumerTopicManagers for consumers cache. @Getter - private final ConcurrentOpenHashMap> consumerTopics; + private final ConcurrentHashMap> consumerTopicManagers; - // cache for topics - private final ConcurrentOpenHashMap topics; + // cache for topics: + private final ConcurrentHashMap> topics; + // cache for references in PersistentTopic: + private final ConcurrentHashMap references; - KafkaTopicManager(BrokerService service) { - this.service = service; - consumerTopics = new ConcurrentOpenHashMap<>(); - topics = new ConcurrentOpenHashMap<>(); + private InternalServerCnx internalServerCnx; + + KafkaTopicManager(KafkaRequestHandler kafkaRequestHandler) { + this.requestHandler = kafkaRequestHandler; + this.pulsarService = kafkaRequestHandler.getPulsarService(); + this.brokerService = pulsarService.getBrokerService(); + this.internalServerCnx = new InternalServerCnx(requestHandler); + + consumerTopicManagers = new ConcurrentHashMap<>(); + topics = new ConcurrentHashMap<>(); + references = new ConcurrentHashMap<>(); + } + + // update Ctx information, since at create time there is no ctx passed into kafkaRequestHandler. + public void updateCtx() { + internalServerCnx.updateCtx(); + if (log.isDebugEnabled()) { + log.debug("internalServerCnx.remoteAddress: {}", internalServerCnx.getRemoteAddress()); + } } // topicName is in pulsar format. e.g. persistent://public/default/topic-partition-0 + // return null if topic not owned by this broker public CompletableFuture getTopicConsumerManager(String topicName) { - return consumerTopics.computeIfAbsent( + return consumerTopicManagers.computeIfAbsent( topicName, - t -> service - .getTopic(topicName, true) - .thenApply(t2 -> { + t -> { + CompletableFuture topic = getTopic(t, true); + if (topic == null) { + log.warn("Failed to getTopicConsumerManager for topic {}. return null", t); + return null; + } + + return topic.thenApply(t2 -> { if (log.isDebugEnabled()) { - log.debug("Call getTopicConsumerManager for {}, and create KafkaTopicConsumerManager.", - topicName); + log.debug("Call getTopicConsumerManager for {}, and create KafkaTopicConsumerManager.", t); } - topics.putIfAbsent(topicName, (PersistentTopic) t2.get()); - return new KafkaTopicConsumerManager((PersistentTopic) t2.get()); - }) - .exceptionally(ex -> { + // return consumer manager + return new KafkaTopicConsumerManager(t2); + }).exceptionally(ex -> { log.error("Failed to getTopicConsumerManager {}. exception:", - topicName, ex); + t, ex); return null; - }) + }); + } ); } - // whether topic exists or not + // whether topic exists in cache. public boolean topicExists(String topicName) { return topics.containsKey(topicName); } - public PersistentTopic addTopic(String topicName, PersistentTopic persistentTopic) { - return topics.putIfAbsent(topicName, persistentTopic); + private Producer registerInPersistentTopic(PersistentTopic persistentTopic) throws Exception { + Producer producer = new InternalProducer(persistentTopic, internalServerCnx, + ((PulsarClientImpl) (pulsarService.getClient())).newRequestId(), + brokerService.generateUniqueProducerName()); + + if (log.isDebugEnabled()) { + log.debug("Register Mock Producer {} into PersistentTopic {}", + producer, persistentTopic.getName()); + } + + // this will register and add USAGE_COUNT_UPDATER. + persistentTopic.addProducer(producer); + return producer; } - public PersistentTopic getTopic(String topicName) { - return topics.get(topicName); + // this should be the only entrance for getTopic, since we need register topic into PersistentTopic. + // return null if not owned by this broker. + public CompletableFuture getTopic(String topicName) { + return getTopic(topicName, false); } - public void close() { - consumerTopics.values() - .forEach(manager -> manager.join().getConsumers().values() - .forEach(pair -> { - try { - pair.join().getLeft().close(); - } catch (Exception e) { - log.error("Failed to close cursor for topic {}. exception:", - pair.join().getLeft().getName(), e); + // For Produce/Consume we need to lookup, to make sure topic served by brokerService, + // or will meet error: "Service unit is not ready when loading the topic". + // If getTopic is called after lookup, then no needLookup. + public synchronized CompletableFuture getTopic(String topicName, boolean needLookup) { + return topics.computeIfAbsent(topicName, + t -> { + try { + CompletableFuture> lookupBroker; + if (needLookup) { + lookupBroker = ((PulsarClientImpl) pulsarService.getClient()).getLookup() + .getBroker(TopicName.get(t)); + } else { + lookupBroker = new CompletableFuture<>(); + lookupBroker.complete(null); } - })); + + final CompletableFuture topicCompletableFuture = new CompletableFuture<>(); + + lookupBroker.whenCompleteAsync((ignore, th) -> { + brokerService + .getTopic(t, true) + .thenApply(t2 -> { + if (log.isDebugEnabled()) { + log.debug("GetTopic for {} in KafkaTopicManager", t); + } + + try { + if (t2.isPresent()) { + PersistentTopic persistentTopic = (PersistentTopic) t2.get(); + references.putIfAbsent(t, registerInPersistentTopic(persistentTopic)); + topicCompletableFuture.complete(persistentTopic); + } else { + log.error("Get empty topic for name {}", t); + topicCompletableFuture.complete(null); + } + } catch (Exception e) { + log.error("Failed to registerInPersistentTopic {}. exception:", + t, e); + topicCompletableFuture.complete(null); + } + + return null; + }) + .exceptionally(ex -> { + log.error("Failed to getTopic {}. exception:", + t, ex); + topicCompletableFuture.complete(null); + return null; + }); + + + }); + return topicCompletableFuture; + } catch (Exception e) { + log.error("Caught error while getclient for topic:{} ", t, e); + return null; + } + }); + } + + // when channel close, release all the topics reference in persistentTopic + public synchronized void close() { + try { + for (CompletableFuture manager : consumerTopicManagers.values()) { + manager.get().getConsumers().values() + .forEach(pair -> { + try { + pair.get().getLeft().close(); + } catch (Exception e) { + log.error("Failed to close cursor for topic {}. exception:", + pair.join().getLeft().getName(), e); + } + }); + } + consumerTopicManagers.clear(); + + for (Map.Entry> entry : topics.entrySet()) { + String topicName = entry.getKey(); + CompletableFuture topicFuture = entry.getValue(); + if (log.isDebugEnabled()) { + log.debug("remove producer {} for topic {} at close()", + references.get(topicName), topicName); + } + topicFuture.get().removeProducer(references.get(topicName)); + references.remove(topicName); + topics.remove(topicName); + } + topics.clear(); + } catch (Exception e) { + log.error("Failed to close KafkaTopicManager. exception:", e); + } + } + + public void deReference(String topicName) { + try { + if (!consumerTopicManagers.containsKey(topicName)) { + return; + } + + consumerTopicManagers.get(topicName).get().getConsumers().values().forEach(pair -> { + try { + pair.join().getLeft().close(); + consumerTopicManagers.remove(topicName); + } catch (Exception e) { + log.error("Failed to close cursor for individual topic {}. exception:", + topicName, e); + } + }); + + if (!topics.containsKey(topicName)) { + return; + } + + topics.get(topicName).get().removeProducer(references.get(topicName)); + topics.remove(topicName); + } catch (Exception e) { + log.error("Failed to close reference for individual topic {}. exception:", + topicName, e); + } } } diff --git a/src/main/java/io/streamnative/kop/MessageFetchContext.java b/src/main/java/io/streamnative/kop/MessageFetchContext.java index bd4eaeb7ab..0353a16d54 100644 --- a/src/main/java/io/streamnative/kop/MessageFetchContext.java +++ b/src/main/java/io/streamnative/kop/MessageFetchContext.java @@ -21,7 +21,6 @@ import io.netty.util.Recycler; import io.netty.util.Recycler.Handle; import io.streamnative.kop.KafkaCommandDecoder.KafkaHeaderAndRequest; -import io.streamnative.kop.KafkaCommandDecoder.ResponseAndRequest; import io.streamnative.kop.utils.MessageIdUtils; import java.util.Date; import java.util.LinkedHashMap; @@ -29,6 +28,7 @@ import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.Collectors; @@ -43,6 +43,7 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.record.MemoryRecords; +import org.apache.kafka.common.requests.AbstractResponse; import org.apache.kafka.common.requests.FetchRequest; import org.apache.kafka.common.requests.FetchResponse; import org.apache.kafka.common.requests.FetchResponse.PartitionData; @@ -86,7 +87,9 @@ public void recycle() { // handle request - public CompletableFuture handleFetch(CompletableFuture fetchResponse) { + public CompletableFuture handleFetch(CompletableFuture fetchResponse) { + LinkedHashMap> responseData = new LinkedHashMap<>(); + // Map of partition and related cursor Map>> topicsAndCursor = ((FetchRequest) fetchRequest.getRequest()) @@ -100,17 +103,37 @@ public CompletableFuture handleFetch(CompletableFuture consumerManager = + requestHandler.getTopicManager().getTopicConsumerManager(topicName.toString()); + + // topic not owned by broker + if (consumerManager == null) { + // return UNKNOWN_TOPIC_OR_PARTITION? + responseData.put(entry.getKey(), + new FetchResponse.PartitionData( + Errors.UNKNOWN_TOPIC_OR_PARTITION, + FetchResponse.INVALID_HIGHWATERMARK, + FetchResponse.INVALID_LAST_STABLE_OFFSET, + FetchResponse.INVALID_LOG_START_OFFSET, + null, + MemoryRecords.EMPTY)); + + log.warn("Partition {} not owned by this broker, will not trigger read for this partition", + entry.getKey()); + return null; + } + return Pair.of( entry.getKey(), - requestHandler.getTopicManager().getTopicConsumerManager(topicName.toString()) - .thenCompose(cm -> cm.remove(offset))); + consumerManager.thenCompose(cm -> cm.remove(offset))); }) + .filter(x -> x != null) .collect(Collectors.toMap(Pair::getKey, Pair::getValue)); // wait to get all the cursor, then readMessages CompletableFuture .allOf(topicsAndCursor.entrySet().stream().map(Map.Entry::getValue).toArray(CompletableFuture[]::new)) - .whenComplete((ignore, ex) -> readMessages(fetchRequest, topicsAndCursor, fetchResponse)); + .whenComplete((ignore, ex) -> readMessages(fetchRequest, topicsAndCursor, fetchResponse, responseData)); return fetchResponse; } @@ -118,23 +141,25 @@ public CompletableFuture handleFetch(CompletableFuture>> cursors, - CompletableFuture resultFuture) { + CompletableFuture resultFuture, + LinkedHashMap> responseData) { AtomicInteger bytesRead = new AtomicInteger(0); - Map> responseValues = new ConcurrentHashMap<>(); + Map> entryValues = new ConcurrentHashMap<>(); if (log.isDebugEnabled()) { log.debug("Request {}: Read Messages for request.", fetch.getHeader()); } - readMessagesInternal(fetch, cursors, bytesRead, responseValues, resultFuture); + readMessagesInternal(fetch, cursors, bytesRead, entryValues, resultFuture, responseData); } private void readMessagesInternal(KafkaHeaderAndRequest fetch, Map>> cursors, AtomicInteger bytesRead, Map> responseValues, - CompletableFuture resultFuture) { + CompletableFuture resultFuture, + LinkedHashMap> responseData) { AtomicInteger entriesRead = new AtomicInteger(0); Map> readFutures = readAllCursorOnce(cursors); CompletableFuture.allOf(readFutures.values().stream().toArray(CompletableFuture[]::new)) @@ -142,7 +167,7 @@ private void readMessagesInternal(KafkaHeaderAndRequest fetch, // keep entries since all read completed. currently only read 1 entry each time. readFutures.forEach((topic, readEntry) -> { try { - Entry entry = readEntry.join(); + Entry entry = readEntry.get(); List entryList = responseValues.computeIfAbsent(topic, l -> Lists.newArrayList()); if (entry != null) { @@ -156,11 +181,19 @@ private void readMessagesInternal(KafkaHeaderAndRequest fetch, } } } catch (Exception e) { - // readEntry.join failed. ignore this partition - log.error("Request {}: Failed readEntry.join for topic: {}. ", + // readEntry.get failed. ignore this partition + log.error("Request {}: Failed readEntry.get for topic: {}. ", fetch.getHeader(), topic, e); cursors.remove(topic); - responseValues.putIfAbsent(topic, Lists.newArrayList()); + + responseData.put(topic, + new FetchResponse.PartitionData( + Errors.NONE, + FetchResponse.INVALID_HIGHWATERMARK, + FetchResponse.INVALID_LAST_STABLE_OFFSET, + FetchResponse.INVALID_LOG_START_OFFSET, + null, + MemoryRecords.EMPTY)); } }); @@ -192,8 +225,6 @@ private void readMessagesInternal(KafkaHeaderAndRequest fetch, fetch.getHeader(), entriesRead.get(), allSize); } - LinkedHashMap> responseData = new LinkedHashMap<>(); - AtomicBoolean allPartitionsNoEntry = new AtomicBoolean(true); responseValues.forEach((topicPartition, entries) -> { final FetchResponse.PartitionData partitionData; @@ -231,34 +262,26 @@ private void readMessagesInternal(KafkaHeaderAndRequest fetch, log.debug("Request {}: All partitions for request read 0 entry", fetch.getHeader()); - // returned earlier, sleep for waitTime - try { - Thread.sleep(waitTime); - } catch (Exception e) { - log.info("Request {}: error while sleep, this is OK.", - fetch.getHeader(), e); - } - - resultFuture.complete(ResponseAndRequest.of( - new FetchResponse(Errors.NONE, - responseData, - ((Integer) THROTTLE_TIME_MS.defaultValue), - ((FetchRequest) fetch.getRequest()).metadata().sessionId()), - fetch)); - this.recycle(); + requestHandler.getPulsarService().getExecutor().schedule(() -> { + resultFuture.complete( + new FetchResponse(Errors.NONE, + responseData, + ((Integer) THROTTLE_TIME_MS.defaultValue), + ((FetchRequest) fetch.getRequest()).metadata().sessionId())); + this.recycle(); + }, waitTime, TimeUnit.MILLISECONDS); } else { - resultFuture.complete(ResponseAndRequest.of( + resultFuture.complete( new FetchResponse( Errors.NONE, responseData, ((Integer) THROTTLE_TIME_MS.defaultValue), - ((FetchRequest) fetch.getRequest()).metadata().sessionId()), - fetch)); + ((FetchRequest) fetch.getRequest()).metadata().sessionId())); this.recycle(); } } else { //need do another round read - readMessagesInternal(fetch, cursors, bytesRead, responseValues, resultFuture); + readMessagesInternal(fetch, cursors, bytesRead, responseValues, resultFuture, responseData); } }); } @@ -273,7 +296,11 @@ private Map> readAllCursorOnce( CompletableFuture readFuture = new CompletableFuture<>(); try { - Pair cursorOffsetPair = pair.getValue().join(); + if (pair.getValue() == null) { + throw new Exception("Topic not owned " + pair.getKey()); + } + + Pair cursorOffsetPair = pair.getValue().get(); cursor = cursorOffsetPair.getLeft(); long keptOffset = cursorOffsetPair.getRight(); diff --git a/src/test/java/io/streamnative/kop/coordinator/group/DistributedGroupCoordinatorTest.java b/src/test/java/io/streamnative/kop/DistributedClusterTest.java similarity index 76% rename from src/test/java/io/streamnative/kop/coordinator/group/DistributedGroupCoordinatorTest.java rename to src/test/java/io/streamnative/kop/DistributedClusterTest.java index fde7ca8462..6648a3bc9a 100644 --- a/src/test/java/io/streamnative/kop/coordinator/group/DistributedGroupCoordinatorTest.java +++ b/src/test/java/io/streamnative/kop/DistributedClusterTest.java @@ -11,7 +11,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package io.streamnative.kop.coordinator.group; +package io.streamnative.kop; import static io.streamnative.kop.KafkaProtocolHandler.PLAINTEXT_PREFIX; import static org.apache.kafka.common.internals.Topic.GROUP_METADATA_TOPIC_NAME; @@ -23,9 +23,6 @@ import com.google.common.collect.Maps; import com.google.common.collect.Sets; import com.google.gson.Gson; -import io.streamnative.kop.KafkaService; -import io.streamnative.kop.KafkaServiceConfiguration; -import io.streamnative.kop.MockKafkaServiceBaseTest; import java.time.Duration; import java.util.List; import java.util.Map; @@ -52,11 +49,11 @@ import org.testng.annotations.Test; - /** - * Unit test {@link GroupCoordinator}. + * Test KoP cluster mode. + * Will setup 2 brokers and do the tests. */ -public class DistributedGroupCoordinatorTest extends MockKafkaServiceBaseTest { +public class DistributedClusterTest extends MockKafkaServiceBaseTest { protected KafkaServiceConfiguration conf1; protected KafkaServiceConfiguration conf2; @@ -72,7 +69,7 @@ public class DistributedGroupCoordinatorTest extends MockKafkaServiceBaseTest { protected int offsetsTopicNumPartitions; - private static final Logger log = LoggerFactory.getLogger(DistributedGroupCoordinatorTest.class); + private static final Logger log = LoggerFactory.getLogger(DistributedClusterTest.class); protected KafkaServiceConfiguration resetConfig(int brokerPort, int webPort, int kafkaPort) { KafkaServiceConfiguration kConfig = new KafkaServiceConfiguration(); @@ -240,6 +237,7 @@ protected void kafkaConsumeCommitMessage(KConsumer kConsumer, assertEquals(i, numMessages); } + // Unit test {@link GroupCoordinator}. @Test(timeOut = 30000) public void testMutiBrokerAndCoordinator() throws Exception { int partitionNumber = 10; @@ -256,6 +254,8 @@ public void testMutiBrokerAndCoordinator() throws Exception { // In setting, each ns has 2 bundles. unload the first part, and this part will be served by broker2. kafkaService1.getAdminClient().namespaces().unloadNamespaceBundle(offsetNs, "0x00000000_0x80000000"); + log.info("unloaded offset namespace, will call lookup to force reload"); + // Offsets partitions should be served by 2 brokers now. Map> offsetTopicMap = Maps.newHashMap(); for (int ii = 0; ii < offsetsTopicNumPartitions; ii++) { @@ -398,4 +398,102 @@ public void testMutiBrokerAndCoordinator() throws Exception { records = kConsumer4.getConsumer().poll(Duration.ofMillis(200)); assertTrue(records.isEmpty()); } + + // Unit test for unload / reload user topic bundle, verify it works well. + @Test(timeOut = 30000) + public void testMutiBrokerUnloadReload() throws Exception { + int partitionNumber = 10; + String kafkaTopicName = "kopMutiBrokerUnloadReload" + partitionNumber; + String pulsarTopicName = "persistent://public/default/" + kafkaTopicName; + String kopNamespace = "public/default"; + + // 0. Preparing: create partitioned topic. + kafkaService1.getAdminClient().topics().createPartitionedTopic(kafkaTopicName, partitionNumber); + + // 1. use a map for serving broker and topics , verify both broker has messages served. + Map> topicMap = Maps.newHashMap(); + for (int ii = 0; ii < partitionNumber; ii++) { + String topicName = pulsarTopicName + PARTITIONED_TOPIC_SUFFIX + ii; + String result = admin.lookups().lookupTopic(topicName); + topicMap.putIfAbsent(result, Lists.newArrayList()); + topicMap.get(result).add(topicName); + log.info("serving broker for topic {} is {}", topicName, result); + } + assertTrue(topicMap.size() == 2); + + // 2. produce consume message with Kafka producer. + int totalMsgs = 50; + String messageStrPrefix = "Message_" + kafkaTopicName + "_"; + @Cleanup + KProducer kProducer = new KProducer(kafkaTopicName, false, getKafkaBrokerPort()); + kafkaPublishMessage(kProducer, totalMsgs, messageStrPrefix); + + List topicPartitions = IntStream.range(0, partitionNumber) + .mapToObj(i -> new TopicPartition(kafkaTopicName, i)).collect(Collectors.toList()); + @Cleanup + KConsumer kConsumer1 = new KConsumer(kafkaTopicName, getKafkaBrokerPort(), "consumer-group-1"); + @Cleanup + KConsumer kConsumer2 = new KConsumer(kafkaTopicName, getKafkaBrokerPort(), "consumer-group-2"); + log.info("Partition size: {}, will consume and commitOffset for 2 consumers", + topicPartitions.size()); + kafkaConsumeCommitMessage(kConsumer1, totalMsgs, messageStrPrefix, topicPartitions); + kafkaConsumeCommitMessage(kConsumer2, totalMsgs, messageStrPrefix, topicPartitions); + + // 3. unload + log.info("Unload namespace, lookup will trigger another reload."); + kafkaService1.getAdminClient().namespaces().unload(kopNamespace); + + // 4. publish consume again + log.info("Re Publish / Consume again."); + kafkaPublishMessage(kProducer, totalMsgs, messageStrPrefix); + kafkaConsumeCommitMessage(kConsumer1, totalMsgs, messageStrPrefix, topicPartitions); + kafkaConsumeCommitMessage(kConsumer2, totalMsgs, messageStrPrefix, topicPartitions); + } + + @Test(timeOut = 30000) + public void testOneBrokerShutdown() throws Exception { + int partitionNumber = 10; + String kafkaTopicName = "kopOneBrokerShutdown" + partitionNumber; + String pulsarTopicName = "persistent://public/default/" + kafkaTopicName; + + // 0. Preparing: create partitioned topic. + kafkaService1.getAdminClient().topics().createPartitionedTopic(kafkaTopicName, partitionNumber); + + // 1. use a map for serving broker and topics , verify both broker has messages served. + Map> topicMap = Maps.newHashMap(); + for (int ii = 0; ii < partitionNumber; ii++) { + String topicName = pulsarTopicName + PARTITIONED_TOPIC_SUFFIX + ii; + String result = admin.lookups().lookupTopic(topicName); + topicMap.putIfAbsent(result, Lists.newArrayList()); + topicMap.get(result).add(topicName); + log.info("serving broker for topic {} is {}", topicName, result); + } + assertTrue(topicMap.size() == 2); + + // 2. produce consume message with Kafka producer. + int totalMsgs = 50; + String messageStrPrefix = "Message_" + kafkaTopicName + "_"; + @Cleanup + KProducer kProducer = new KProducer(kafkaTopicName, false, getKafkaBrokerPort()); + kafkaPublishMessage(kProducer, totalMsgs, messageStrPrefix); + + List topicPartitions = IntStream.range(0, partitionNumber) + .mapToObj(i -> new TopicPartition(kafkaTopicName, i)).collect(Collectors.toList()); + @Cleanup + KConsumer kConsumer1 = new KConsumer(kafkaTopicName, getKafkaBrokerPort(), "consumer-group-1"); + @Cleanup + KConsumer kConsumer2 = new KConsumer(kafkaTopicName, getKafkaBrokerPort(), "consumer-group-2"); + log.info("Partition size: {}, will consume and commitOffset for 2 consumers", + topicPartitions.size()); + kafkaConsumeCommitMessage(kConsumer1, totalMsgs, messageStrPrefix, topicPartitions); + kafkaConsumeCommitMessage(kConsumer2, totalMsgs, messageStrPrefix, topicPartitions); + + kafkaService1.close(); + + // 4. publish consume again + log.info("Re Publish / Consume again."); + kafkaPublishMessage(kProducer, totalMsgs, messageStrPrefix); + kafkaConsumeCommitMessage(kConsumer1, totalMsgs, messageStrPrefix, topicPartitions); + kafkaConsumeCommitMessage(kConsumer2, totalMsgs, messageStrPrefix, topicPartitions); + } } diff --git a/src/test/java/io/streamnative/kop/KafkaApisTest.java b/src/test/java/io/streamnative/kop/KafkaApisTest.java index eb8a2eb466..5275e91898 100644 --- a/src/test/java/io/streamnative/kop/KafkaApisTest.java +++ b/src/test/java/io/streamnative/kop/KafkaApisTest.java @@ -30,7 +30,6 @@ import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; import io.streamnative.kop.KafkaCommandDecoder.KafkaHeaderAndRequest; -import io.streamnative.kop.KafkaCommandDecoder.ResponseAndRequest; import io.streamnative.kop.utils.MessageIdUtils; import java.net.InetSocketAddress; import java.net.SocketAddress; @@ -58,6 +57,7 @@ import org.apache.kafka.common.record.MemoryRecords; import org.apache.kafka.common.record.MutableRecordBatch; import org.apache.kafka.common.requests.AbstractRequest; +import org.apache.kafka.common.requests.AbstractResponse; import org.apache.kafka.common.requests.FetchRequest; import org.apache.kafka.common.requests.FetchResponse; import org.apache.kafka.common.requests.IsolationLevel; @@ -132,7 +132,6 @@ protected void setup() throws Exception { kafkaRequestHandler = new KafkaRequestHandler( kafkaService, kafkaService.getKafkaConfig(), - kafkaService.getKafkaTopicManager(), kafkaService.getGroupCoordinator(), false); ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class); Channel mockChannel = mock(Channel.class); @@ -166,8 +165,8 @@ KafkaHeaderAndRequest buildRequest(AbstractRequest.Builder builder) { return new KafkaHeaderAndRequest(header, body, byteBuf, serviceAddress); } - CompletableFuture checkInvalidPartition(String topic, - int invalidPartitionId) { + CompletableFuture checkInvalidPartition(String topic, + int invalidPartitionId) { TopicPartition invalidTopicPartition = new TopicPartition(topic, invalidPartitionId); PartitionData partitionOffsetCommitData = new OffsetCommitRequest.PartitionData(15L, ""); Map offsetData = Maps.newHashMap(); @@ -182,19 +181,17 @@ public void testOffsetCommitWithInvalidPartition() throws Exception { String topicName = "kopOffsetCommitWithInvalidPartition"; // invalid partition id -1; - CompletableFuture invalidResponse1 = checkInvalidPartition(topicName, -1); - ResponseAndRequest response1 = invalidResponse1.get(); - assertEquals(response1.getRequest().getHeader().apiKey(), ApiKeys.OFFSET_COMMIT); + CompletableFuture invalidResponse1 = checkInvalidPartition(topicName, -1); + AbstractResponse response1 = invalidResponse1.get(); TopicPartition topicPartition1 = new TopicPartition(topicName, -1); - assertEquals(((OffsetCommitResponse) response1.getResponse()).responseData().get(topicPartition1), + assertEquals(((OffsetCommitResponse) response1).responseData().get(topicPartition1), Errors.UNKNOWN_TOPIC_OR_PARTITION); // invalid partition id 1. - CompletableFuture invalidResponse2 = checkInvalidPartition(topicName, 1); + CompletableFuture invalidResponse2 = checkInvalidPartition(topicName, 1); TopicPartition topicPartition2 = new TopicPartition(topicName, 1); - ResponseAndRequest response2 = invalidResponse2.get(); - assertEquals(response2.getRequest().getHeader().apiKey(), ApiKeys.OFFSET_COMMIT); - assertEquals(((OffsetCommitResponse) response2.getResponse()).responseData().get(topicPartition2), + AbstractResponse response2 = invalidResponse2.get(); + assertEquals(((OffsetCommitResponse) response2).responseData().get(topicPartition2), Errors.UNKNOWN_TOPIC_OR_PARTITION); } @@ -270,12 +267,11 @@ public void testReadUncommittedConsumerListOffsetEarliestOffsetEquals() throws E .setTargetTimes(targetTimes); KafkaHeaderAndRequest request = buildRequest(builder); - CompletableFuture responseFuture = kafkaRequestHandler + CompletableFuture responseFuture = kafkaRequestHandler .handleListOffsetRequest(request); - ResponseAndRequest response = responseFuture.get(); - ListOffsetResponse listOffsetResponse = (ListOffsetResponse) response.getResponse(); - assertEquals(response.getRequest().getHeader().apiKey(), ApiKeys.LIST_OFFSETS); + AbstractResponse response = responseFuture.get(); + ListOffsetResponse listOffsetResponse = (ListOffsetResponse) response; assertEquals(listOffsetResponse.responseData().get(tp).error, Errors.NONE); assertEquals(listOffsetResponse.responseData().get(tp).offset, Long.valueOf(limitOffset)); assertEquals(listOffsetResponse.responseData().get(tp).timestamp, Long.valueOf(NO_TIMESTAMP)); @@ -339,12 +335,11 @@ public void testConsumerListOffsetLatest() throws Exception { .setTargetTimes(targetTimes); KafkaHeaderAndRequest request = buildRequest(builder); - CompletableFuture responseFuture = kafkaRequestHandler + CompletableFuture responseFuture = kafkaRequestHandler .handleListOffsetRequest(request); - ResponseAndRequest response = responseFuture.get(); - ListOffsetResponse listOffsetResponse = (ListOffsetResponse) response.getResponse(); - assertEquals(response.getRequest().getHeader().apiKey(), ApiKeys.LIST_OFFSETS); + AbstractResponse response = responseFuture.get(); + ListOffsetResponse listOffsetResponse = (ListOffsetResponse) response; assertEquals(listOffsetResponse.responseData().get(tp).error, Errors.NONE); assertEquals(listOffsetResponse.responseData().get(tp).offset, Long.valueOf(limitOffset)); assertEquals(listOffsetResponse.responseData().get(tp).timestamp, Long.valueOf(NO_TIMESTAMP)); @@ -538,9 +533,9 @@ public void testBrokerRespectsPartitionsOrderAndSizeLimits() throws Exception { maxPartitionBytes, shuffledTopicPartitions1, Collections.EMPTY_MAP); - CompletableFuture responseFuture1 = kafkaRequestHandler.handleFetchRequest(fetchRequest1); + CompletableFuture responseFuture1 = kafkaRequestHandler.handleFetchRequest(fetchRequest1); FetchResponse fetchResponse1 = - (FetchResponse) responseFuture1.get().getResponse(); + (FetchResponse) responseFuture1.get(); checkFetchResponse(shuffledTopicPartitions1, fetchResponse1, maxPartitionBytes, maxResponseBytes, messagesPerPartition); @@ -556,9 +551,9 @@ public void testBrokerRespectsPartitionsOrderAndSizeLimits() throws Exception { maxPartitionBytes, shuffledTopicPartitions2, Collections.EMPTY_MAP); - CompletableFuture responseFuture2 = kafkaRequestHandler.handleFetchRequest(fetchRequest2); + CompletableFuture responseFuture2 = kafkaRequestHandler.handleFetchRequest(fetchRequest2); FetchResponse fetchResponse2 = - (FetchResponse) responseFuture2.get().getResponse(); + (FetchResponse) responseFuture2.get(); checkFetchResponse(shuffledTopicPartitions2, fetchResponse2, maxPartitionBytes, maxResponseBytes, messagesPerPartition); @@ -577,15 +572,14 @@ public void testBrokerRespectsPartitionsOrderAndSizeLimits() throws Exception { maxPartitionBytes, shuffledTopicPartitions3, offsetMaps); - CompletableFuture responseFuture3 = kafkaRequestHandler.handleFetchRequest(fetchRequest3); + CompletableFuture responseFuture3 = kafkaRequestHandler.handleFetchRequest(fetchRequest3); FetchResponse fetchResponse3 = - (FetchResponse) responseFuture3.get().getResponse(); + (FetchResponse) responseFuture3.get(); checkFetchResponse(shuffledTopicPartitions3, fetchResponse3, maxPartitionBytes, maxResponseBytes, messagesPerPartition); } - // verify Metadata request handling. @Test(timeOut = 20000) public void testBrokerHandleTopicMetadataRequest() throws Exception { @@ -596,10 +590,10 @@ public void testBrokerHandleTopicMetadataRequest() throws Exception { List topicPartitions = createTopics(topicName, numberTopics, numberPartitions); List kafkaTopics = getCreatedTopics(topicName, numberTopics); KafkaHeaderAndRequest metadataRequest = createTopicMetadataRequest(kafkaTopics); - CompletableFuture responseFuture = + CompletableFuture responseFuture = kafkaRequestHandler.handleTopicMetadataRequest(metadataRequest); - MetadataResponse metadataResponse = (MetadataResponse) responseFuture.get().getResponse(); + MetadataResponse metadataResponse = (MetadataResponse) responseFuture.get(); // verify all served by same broker : localhost:port assertEquals(metadataResponse.brokers().size(), 1); @@ -624,4 +618,27 @@ public void testBrokerHandleTopicMetadataRequest() throws Exception { assertEquals(topicMetadata.partitionMetadata().size(), numberPartitions); }); } + + @Test(timeOut = 20000, enabled = false) + // https://github.com/streamnative/kop/issues/51 + public void testGetOffsetsForUnknownTopic() throws Exception { + String topicName = "kopTestGetOffsetsForUnknownTopic"; + + TopicPartition tp = new TopicPartition(topicName, 0); + Map targetTimes = Maps.newHashMap(); + targetTimes.put(tp, ListOffsetRequest.LATEST_TIMESTAMP); + + ListOffsetRequest.Builder builder = ListOffsetRequest.Builder + .forConsumer(false, IsolationLevel.READ_UNCOMMITTED) + .setTargetTimes(targetTimes); + + KafkaHeaderAndRequest request = buildRequest(builder); + CompletableFuture responseFuture = kafkaRequestHandler + .handleListOffsetRequest(request); + + AbstractResponse response = responseFuture.get(); + ListOffsetResponse listOffsetResponse = (ListOffsetResponse) response; + assertEquals(listOffsetResponse.responseData().get(tp).error, + Errors.UNKNOWN_TOPIC_OR_PARTITION); + } } diff --git a/src/test/java/io/streamnative/kop/KafkaRequestHandlerTest.java b/src/test/java/io/streamnative/kop/KafkaRequestHandlerTest.java index 5623e9393b..06f124af17 100644 --- a/src/test/java/io/streamnative/kop/KafkaRequestHandlerTest.java +++ b/src/test/java/io/streamnative/kop/KafkaRequestHandlerTest.java @@ -17,24 +17,15 @@ import static io.streamnative.kop.utils.TopicNameUtils.getKafkaTopicNameFromPulsarTopicname; import static io.streamnative.kop.utils.TopicNameUtils.getPartitionedTopicNameWithoutPartitions; import static org.apache.pulsar.common.naming.TopicName.PARTITIONED_TOPIC_SUFFIX; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.Mockito.CALLS_REAL_METHODS; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.times; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; import com.google.common.collect.Sets; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; -import io.netty.channel.Channel; -import io.netty.channel.ChannelHandlerContext; import io.streamnative.kop.KafkaCommandDecoder.KafkaHeaderAndRequest; import io.streamnative.kop.KafkaCommandDecoder.KafkaHeaderAndResponse; import java.net.InetSocketAddress; -import java.net.SocketAddress; import java.nio.ByteBuffer; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.common.Node; @@ -101,7 +92,6 @@ protected void setup() throws Exception { handler = new KafkaRequestHandler( kafkaService, kafkaService.getKafkaConfig(), - kafkaService.getKafkaTopicManager(), kafkaService.getGroupCoordinator(), false); } @@ -175,41 +165,6 @@ public void testResponseToByteBuf() throws Exception { assertEquals(parsedResponse.apiVersions().size(), apiVersionsResponse.apiVersions().size()); } - @Test - public void testChannelRead() throws Exception { - int correlationId = 7777; - String clientId = "KopClientId"; - - ApiVersionsRequest apiVersionsRequest = new ApiVersionsRequest.Builder().build(); - RequestHeader header = new RequestHeader( - ApiKeys.API_VERSIONS, - ApiKeys.API_VERSIONS.latestVersion(), - clientId, - correlationId); - - ByteBuffer serializedRequest = apiVersionsRequest.serialize(header); - int size = serializedRequest.remaining(); - ByteBuf inputBuf = Unpooled.buffer(size); - inputBuf.writeBytes(serializedRequest); - - ChannelHandlerContext ctx = mock(ChannelHandlerContext.class); - Channel channel = mock(Channel.class); - SocketAddress address = mock(SocketAddress.class); - - when(ctx.channel()).thenReturn(channel); - when(channel.remoteAddress()).thenReturn(address); - - try { - handler = mock(KafkaRequestHandler.class, CALLS_REAL_METHODS); - handler.channelActive(ctx); - handler.channelRead(mock(ChannelHandlerContext.class), inputBuf); - } catch (Exception e) { - // not mock other module, expect meet exception. - } - - verify(handler, times(1)).handleApiVersionsRequest(any()); - } - @Test public void testNewNode() { String host = "192.168.168.168"; diff --git a/src/test/java/io/streamnative/kop/KafkaRequestTypeTest.java b/src/test/java/io/streamnative/kop/KafkaRequestTypeTest.java index dec0b09e2d..60bf4414b1 100644 --- a/src/test/java/io/streamnative/kop/KafkaRequestTypeTest.java +++ b/src/test/java/io/streamnative/kop/KafkaRequestTypeTest.java @@ -15,17 +15,13 @@ import static java.nio.charset.StandardCharsets.UTF_8; -import static org.apache.pulsar.common.naming.TopicName.PARTITIONED_TOPIC_SUFFIX; import static org.testng.Assert.assertEquals; -import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; -import static org.testng.Assert.fail; import com.google.common.collect.Lists; import com.google.common.collect.Sets; -import io.streamnative.kop.utils.MessageIdUtils; import java.time.Duration; import java.util.Base64; import java.util.Collections; @@ -38,9 +34,6 @@ import java.util.stream.IntStream; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; -import org.apache.bookkeeper.mledger.ManagedCursor; -import org.apache.bookkeeper.mledger.impl.PositionImpl; -import org.apache.commons.lang3.tuple.Pair; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.producer.ProducerRecord; @@ -51,7 +44,6 @@ import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerBuilder; -import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TenantInfo; @@ -466,138 +458,4 @@ public void testPulsarProduceKafkaConsume2(int partitionNumber) throws Exception ConsumerRecords records = kConsumer2.getConsumer().poll(Duration.ofSeconds(1)); assertTrue(records.isEmpty()); } - - - // Test kafka topic consumer manager. - // 1. topic has no entry, read no entry, tm has one cursor and target to read the first entry. - // 2. produce entry, after read all entry. tm has one cursor, and target to read entry after lastEntry. - // 3. has no entry to read again. tm has one cursor, and after each empty read, cursor read offset not changed. - @Test(timeOut = 20000) - public void testTopicConsumerManager() throws Exception { - int partitionNumber = 1; - String kafkaTopicName = "testTopicConsumerManager" + partitionNumber; - String pulsarTopicName = "persistent://public/default/" + kafkaTopicName + PARTITIONED_TOPIC_SUFFIX + 0; - - // create partitioned topic. - kafkaService.getAdminClient().topics().createPartitionedTopic(kafkaTopicName, partitionNumber); - - int totalMsgs = 10; - String messageStrPrefix = "Message_Kop_testTopicConsumerManager_" + partitionNumber + "_"; - - ProducerBuilder producerBuilder = pulsarClient.newProducer() - .topic(pulsarTopicName) - .enableBatching(false); - @Cleanup - Producer producer = producerBuilder.create(); - - // above producer only created a topic, but with no data. consumer retry read but read no entry. - @Cleanup - KConsumer kConsumer = new KConsumer(kafkaTopicName, getKafkaBrokerPort(), true); - kConsumer.getConsumer().subscribe(Collections.singletonList(kafkaTopicName)); - - KafkaTopicConsumerManager tm = kafkaService - .getKafkaTopicManager() - .getTopicConsumerManager(pulsarTopicName) - .get(); - - // read empty entry will remove and add cursor each time. - int i = 0; - while (i < 7) { - if (log.isDebugEnabled()) { - log.debug("start poll empty entry: {}", i); - } - ConsumerRecords records = kConsumer.getConsumer().poll(Duration.ofSeconds(1)); - for (ConsumerRecord record : records) { - Integer key = record.key(); - assertEquals(messageStrPrefix + key.toString(), record.value()); - if (log.isDebugEnabled()) { - log.debug("Kafka Consumer Received message: {}, {} at offset {}", - record.key(), record.value(), record.offset()); - } - } - i++; - } - - // expected tm only have one item. and entryId should be 0 - long size = tm.getConsumers().size(); - assertEquals(size, 1); - tm.getConsumers().forEach((offset, cursor) -> { - try { - PositionImpl position = MessageIdUtils.getPosition(offset); - long ledgerId = position.getLedgerId(); - assertNotEquals(ledgerId, 0); - assertEquals(position.getEntryId(), 0); - assertEquals(cursor.get().getRight(), Long.valueOf(offset)); - } catch (Exception e) { - fail("should not throw exception"); - } - }); - - MessageId messageId = null; - // produce some message - for (i = 0; i < totalMsgs; i++) { - String message = messageStrPrefix + i; - messageId = producer.newMessage() - .keyBytes(kafkaIntSerialize(Integer.valueOf(i))) - .value(message.getBytes()) - .send(); - } - - i = 0; - // receive all message. - while (i < totalMsgs) { - if (log.isDebugEnabled()) { - log.debug("start poll message: {}", i); - } - ConsumerRecords records = kConsumer.getConsumer().poll(Duration.ofSeconds(1)); - for (ConsumerRecord record : records) { - Integer key = record.key(); - assertEquals(messageStrPrefix + key.toString(), record.value()); - if (log.isDebugEnabled()) { - log.debug("Kafka Consumer Received message: {}, {} at offset {}", - record.key(), record.value(), record.offset()); - } - i++; - } - } - - // expect have one item, and offset equals to lastmessageId + 1 - size = tm.getConsumers().size(); - assertEquals(size, 1); - - MessageIdImpl lastMessageId = (MessageIdImpl) messageId; - long ledgerId = lastMessageId.getLedgerId(); - long entryId = lastMessageId.getEntryId(); - long lastOffset = MessageIdUtils.getOffset(ledgerId, entryId + 1); - CompletableFuture> cursor = tm.getConsumers().get(lastOffset); - assertNotNull(cursor); - assertEquals(cursor.get().getRight(), Long.valueOf(lastOffset)); - - - // After read all entry, read no entry again, this will remove and add cursor each time. - i = 0; - while (i < 7) { - if (log.isDebugEnabled()) { - log.debug("start poll empty entry again: {}", i); - } - ConsumerRecords records = kConsumer.getConsumer().poll(Duration.ofSeconds(1)); - for (ConsumerRecord record : records) { - Integer key = record.key(); - assertEquals(messageStrPrefix + key.toString(), record.value()); - if (log.isDebugEnabled()) { - log.debug("Kafka Consumer Received message: {}, {} at offset {}", - record.key(), record.value(), record.offset()); - } - } - i++; - } - - // expect have one item, and offset equals to lastmessageId + 1 - size = tm.getConsumers().size(); - assertEquals(size, 1); - cursor = tm.getConsumers().get(lastOffset); - assertNotNull(cursor); - assertEquals(cursor.get().getRight(), Long.valueOf(lastOffset)); - } - } diff --git a/src/test/java/io/streamnative/kop/KafkaTopicConsumerManagerTest.java b/src/test/java/io/streamnative/kop/KafkaTopicConsumerManagerTest.java index a1ad900226..170dc24d37 100644 --- a/src/test/java/io/streamnative/kop/KafkaTopicConsumerManagerTest.java +++ b/src/test/java/io/streamnative/kop/KafkaTopicConsumerManagerTest.java @@ -13,13 +13,20 @@ */ package io.streamnative.kop; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertTrue; import com.google.common.collect.Sets; +import io.netty.channel.Channel; +import io.netty.channel.ChannelHandlerContext; import io.streamnative.kop.utils.MessageIdUtils; +import java.net.InetSocketAddress; +import java.net.SocketAddress; import java.util.concurrent.CompletableFuture; +import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.client.api.Producer; @@ -33,9 +40,12 @@ /** * Pulsar service configuration object. */ +@Slf4j public class KafkaTopicConsumerManagerTest extends MockKafkaServiceBaseTest { private KafkaTopicManager kafkaTopicManager; + private KafkaRequestHandler kafkaRequestHandler; + private SocketAddress serviceAddress; @BeforeMethod @Override @@ -47,7 +57,18 @@ protected void setup() throws Exception { admin.namespaces().setRetention("public/default", new RetentionPolicies(20, 100)); - kafkaTopicManager = new KafkaTopicManager(kafkaService.getBrokerService()); + kafkaRequestHandler = new KafkaRequestHandler( + kafkaService, + kafkaService.getKafkaConfig(), + kafkaService.getGroupCoordinator(), false); + ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class); + Channel mockChannel = mock(Channel.class); + doReturn(mockChannel).when(mockCtx).channel(); + kafkaRequestHandler.ctx = mockCtx; + + serviceAddress = new InetSocketAddress(kafkaService.getBindAddress(), kafkaBrokerPort); + + kafkaTopicManager = new KafkaTopicManager(kafkaRequestHandler); } @AfterMethod @@ -68,7 +89,7 @@ public void testGetTopicConsumerManager() throws Exception { KafkaTopicConsumerManager topicConsumerManager2 = tcm.get(); assertTrue(topicConsumerManager == topicConsumerManager2); - assertEquals(kafkaTopicManager.getConsumerTopics().size(), 1); + assertEquals(kafkaTopicManager.getConsumerTopicManagers().size(), 1); // 2. verify another get with different topic will return different tcm String topicName2 = "persistent://public/default/testGetTopicConsumerManager2"; @@ -76,7 +97,7 @@ public void testGetTopicConsumerManager() throws Exception { tcm = kafkaTopicManager.getTopicConsumerManager(topicName2); topicConsumerManager2 = tcm.get(); assertTrue(topicConsumerManager != topicConsumerManager2); - assertEquals(kafkaTopicManager.getConsumerTopics().size(), 2); + assertEquals(kafkaTopicManager.getConsumerTopicManagers().size(), 2); } @@ -155,4 +176,5 @@ public void testTopicConsumerManagerRemoveAndAdd() throws Exception { assertNotEquals(cursor2.getName(), cursor.getName()); assertEquals(cursorCompletableFuture.get().getRight(), Long.valueOf(offset)); } + } diff --git a/src/test/java/io/streamnative/kop/LogOffsetTest.java b/src/test/java/io/streamnative/kop/LogOffsetTest.java deleted file mode 100644 index 233e80b31b..0000000000 --- a/src/test/java/io/streamnative/kop/LogOffsetTest.java +++ /dev/null @@ -1,61 +0,0 @@ -/** - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package io.streamnative.kop; - -import static org.testng.Assert.assertEquals; - -import com.google.common.collect.Maps; -import io.streamnative.kop.KafkaCommandDecoder.KafkaHeaderAndRequest; -import io.streamnative.kop.KafkaCommandDecoder.ResponseAndRequest; -import java.util.Map; -import java.util.concurrent.CompletableFuture; -import lombok.extern.slf4j.Slf4j; -import org.apache.kafka.common.TopicPartition; -import org.apache.kafka.common.protocol.ApiKeys; -import org.apache.kafka.common.protocol.Errors; -import org.apache.kafka.common.requests.IsolationLevel; -import org.apache.kafka.common.requests.ListOffsetRequest; -import org.apache.kafka.common.requests.ListOffsetResponse; -import org.testng.annotations.Test; - -/** - * Validate LogOffset. - */ -@Slf4j -public class LogOffsetTest extends KafkaApisTest { - - @Test(timeOut = 20000, enabled = false) - // https://github.com/streamnative/kop/issues/51 - public void testGetOffsetsForUnknownTopic() throws Exception { - String topicName = "kopTestGetOffsetsForUnknownTopic"; - - TopicPartition tp = new TopicPartition(topicName, 0); - Map targetTimes = Maps.newHashMap(); - targetTimes.put(tp, ListOffsetRequest.LATEST_TIMESTAMP); - - ListOffsetRequest.Builder builder = ListOffsetRequest.Builder - .forConsumer(false, IsolationLevel.READ_UNCOMMITTED) - .setTargetTimes(targetTimes); - - KafkaHeaderAndRequest request = buildRequest(builder); - CompletableFuture responseFuture = kafkaRequestHandler - .handleListOffsetRequest(request); - - ResponseAndRequest response = responseFuture.get(); - ListOffsetResponse listOffsetResponse = (ListOffsetResponse) response.getResponse(); - assertEquals(response.getRequest().getHeader().apiKey(), ApiKeys.LIST_OFFSETS); - assertEquals(listOffsetResponse.responseData().get(tp).error, - Errors.UNKNOWN_TOPIC_OR_PARTITION); - } -}