From e81a043530a675b7d466ede2c0136d6948e66e81 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Tue, 23 Nov 2021 18:30:03 +0800 Subject: [PATCH 01/15] Split store message into independent class --- .../handlers/kop/DelayedProduceAndFetch.java | 4 +- .../handlers/kop/KafkaProtocolHandler.java | 14 + .../handlers/kop/KafkaRequestHandler.java | 110 +++--- .../handlers/kop/TenantContextManager.java | 11 + .../handlers/kop/storage/PartitionLog.java | 356 ++++++++++++++++++ .../kop/storage/PartitionLogManager.java | 49 +++ .../handlers/kop/storage/ReplicaManager.java | 152 ++++++++ .../handlers/kop/storage/package-info.java | 14 + .../pulsar/handlers/kop/utils/KopTopic.java | 9 + .../kop/KopProtocolHandlerTestBase.java | 8 + 10 files changed, 666 insertions(+), 61 deletions(-) create mode 100644 kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java create mode 100644 kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java create mode 100644 kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java create mode 100644 kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/package-info.java diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java index a1070e0377..be68b7b4b2 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java @@ -20,12 +20,12 @@ /** * A delayed create topic operation that is stored in the topic purgatory. */ -class DelayedProduceAndFetch extends DelayedOperation { +public class DelayedProduceAndFetch extends DelayedOperation { private final AtomicInteger topicPartitionNum; private final Runnable callback; - DelayedProduceAndFetch(long delayMs, AtomicInteger topicPartitionNum, Runnable callback) { + public DelayedProduceAndFetch(long delayMs, AtomicInteger topicPartitionNum, Runnable callback) { super(delayMs, Optional.empty()); this.topicPartitionNum = topicPartitionNum; this.callback = callback; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index 73463ea946..2bdd29698a 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -28,6 +28,7 @@ import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.stats.PrometheusMetricsProvider; import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; +import io.streamnative.pulsar.handlers.kop.storage.ReplicaManager; import io.streamnative.pulsar.handlers.kop.utils.ConfigurationUtils; import io.streamnative.pulsar.handlers.kop.utils.MetadataUtils; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; @@ -97,6 +98,7 @@ public class KafkaProtocolHandler implements ProtocolHandler, TenantContextManag private final Map groupCoordinatorsByTenant = new ConcurrentHashMap<>(); private final Map transactionCoordinatorByTenant = new ConcurrentHashMap<>(); + private final Map replicaManagerByTenant = new ConcurrentHashMap<>(); @Override public GroupCoordinator getGroupCoordinator(String tenant) { @@ -113,6 +115,18 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) { return transactionCoordinatorByTenant.computeIfAbsent(tenant, this::createAndBootTransactionCoordinator); } + @Override + public ReplicaManager getReplicaManager(String tenant) { + return replicaManagerByTenant.computeIfAbsent(tenant, s -> { + try { + return new ReplicaManager(kafkaConfig, Time.SYSTEM, producePurgatory, fetchPurgatory); + } catch (Exception e) { + log.error("Failed to init ReplicaManager for tenant {}", tenant, e); + throw new IllegalStateException(e); + } + }); + } + /** * Listener for invalidating the global Broker ownership cache. */ diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index d411550786..a683e11e0f 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -49,6 +49,7 @@ import io.streamnative.pulsar.handlers.kop.security.auth.ResourceType; import io.streamnative.pulsar.handlers.kop.security.auth.SimpleAclAuthorizer; import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; +import io.streamnative.pulsar.handlers.kop.storage.ReplicaManager; import io.streamnative.pulsar.handlers.kop.utils.CoreUtils; import io.streamnative.pulsar.handlers.kop.utils.GroupIdUtils; import io.streamnative.pulsar.handlers.kop.utils.KafkaRequestUtils; @@ -58,7 +59,6 @@ import io.streamnative.pulsar.handlers.kop.utils.OffsetFinder; import io.streamnative.pulsar.handlers.kop.utils.TopicNameUtils; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; -import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationKey; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory; import java.net.InetSocketAddress; import java.nio.ByteBuffer; @@ -279,6 +279,10 @@ public TransactionCoordinator getTransactionCoordinator() { return tenantContextManager.getTransactionCoordinator(getCurrentTenant()); } + public ReplicaManager getReplicaManager() { + return tenantContextManager.getReplicaManager(getCurrentTenant()); + } + public KafkaRequestHandler(PulsarService pulsarService, KafkaServiceConfiguration kafkaConfig, TenantContextManager tenantContextManager, @@ -956,82 +960,70 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, final int numPartitions = produceRequest.partitionRecordsOrFail().size(); - final Map responseMap = new ConcurrentHashMap<>(); - // delay produce - final AtomicInteger topicPartitionNum = new AtomicInteger(produceRequest.partitionRecordsOrFail().size()); + final Map unauthorizedTopicResponsesMap = new ConcurrentHashMap<>(); + final Map authorizedRequestInfo = new ConcurrentHashMap<>(); int timeoutMs = produceRequest.timeout(); - Runnable complete = () -> { - topicPartitionNum.set(0); - if (resultFuture.isDone()) { - // It may be triggered again in DelayedProduceAndFetch - return; - } - // add the topicPartition with timeout error if it's not existed in responseMap - produceRequest.partitionRecordsOrFail().keySet().forEach(topicPartition -> { - if (!responseMap.containsKey(topicPartition)) { - responseMap.put(topicPartition, new PartitionResponse(Errors.REQUEST_TIMED_OUT)); - } - }); - if (log.isDebugEnabled()) { - log.debug("[{}] Request {}: Complete handle produce.", ctx.channel(), produceHar.toString()); - } - resultFuture.complete(new ProduceResponse(responseMap)); - }; - BiConsumer addPartitionResponse = (topicPartition, response) -> { - responseMap.put(topicPartition, response); - // reset topicPartitionNum - int restTopicPartitionNum = topicPartitionNum.decrementAndGet(); - if (restTopicPartitionNum < 0) { - return; - } - if (restTopicPartitionNum == 0) { - complete.run(); + String namespacePrefix = currentNamespacePrefix(); + final AtomicInteger unfinishedAuthorizationCount = new AtomicInteger(numPartitions); + Consumer completeOne = (action) -> { + // When complete one authorization or failed, will do the action first. + action.run(); + if (unfinishedAuthorizationCount.decrementAndGet() == 0) { + CompletableFuture> responseCallback = + new CompletableFuture<>(); + getReplicaManager().appendRecords( + timeoutMs, + false, + produceHar.getRequest().version(), + topicManager, + namespacePrefix, + authorizedRequestInfo, + requestStats, + getTransactionCoordinator(), + this::startSendOperationForThrottling, + this::completeSendOperationForThrottling, + pendingTopicFuturesMap, + responseCallback + ); + responseCallback.thenAccept(response -> { + Map mergedResponse = Maps.newHashMap(); + mergedResponse.putAll(response); + mergedResponse.putAll(unauthorizedTopicResponsesMap); + resultFuture.complete(new ProduceResponse(mergedResponse)); + }).exceptionally(ex -> { + resultFuture.completeExceptionally(ex); + return null; + }); } }; - String namespacePrefix = currentNamespacePrefix(); produceRequest.partitionRecordsOrFail().forEach((topicPartition, records) -> { - final Consumer offsetConsumer = offset -> addPartitionResponse.accept( - topicPartition, new PartitionResponse(Errors.NONE, offset, -1L, -1L)); - final Consumer errorsConsumer = - errors -> addPartitionResponse.accept(topicPartition, new PartitionResponse(errors)); - final Consumer exceptionConsumer = - e -> addPartitionResponse.accept(topicPartition, new PartitionResponse(Errors.forException(e))); final String fullPartitionName = KopTopic.toString(topicPartition, namespacePrefix); - authorize(AclOperation.WRITE, Resource.of(ResourceType.TOPIC, fullPartitionName)) .whenComplete((isAuthorized, ex) -> { if (ex != null) { log.error("Write topic authorize failed, topic - {}. {}", fullPartitionName, ex.getMessage()); - errorsConsumer.accept(Errors.TOPIC_AUTHORIZATION_FAILED); + completeOne.accept(() -> { + unauthorizedTopicResponsesMap.put(topicPartition, + new ProduceResponse.PartitionResponse(Errors.TOPIC_AUTHORIZATION_FAILED)); + }); return; } if (!isAuthorized) { - errorsConsumer.accept(Errors.TOPIC_AUTHORIZATION_FAILED); + completeOne.accept(() -> { + unauthorizedTopicResponsesMap.put(topicPartition, + new ProduceResponse.PartitionResponse(Errors.TOPIC_AUTHORIZATION_FAILED)); + }); return; } - - handlePartitionRecords(produceHar, - topicPartition, - records, - numPartitions, - fullPartitionName, - offsetConsumer, - errorsConsumer, - exceptionConsumer); + completeOne.accept(() -> { + authorizedRequestInfo.put(topicPartition, records); + }); }); }); - // delay produce - if (timeoutMs <= 0) { - complete.run(); - } else { - List delayedCreateKeys = - produceRequest.partitionRecordsOrFail().keySet().stream() - .map(DelayedOperationKey.TopicPartitionOperationKey::new).collect(Collectors.toList()); - DelayedProduceAndFetch delayedProduce = new DelayedProduceAndFetch(timeoutMs, topicPartitionNum, complete); - producePurgatory.tryCompleteElseWatch(delayedProduce, delayedCreateKeys); - } + + } private void handlePartitionRecords(final KafkaHeaderAndRequest produceHar, diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/TenantContextManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/TenantContextManager.java index f5071bc382..0b63726dc5 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/TenantContextManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/TenantContextManager.java @@ -15,11 +15,13 @@ import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; +import io.streamnative.pulsar.handlers.kop.storage.ReplicaManager; /** * Access Tenant level coordinators. */ public interface TenantContextManager { + /** * Access the GroupCoordinator for the current Tenant. * This method bootstraps a new GroupCoordinator if it is not started @@ -27,6 +29,7 @@ public interface TenantContextManager { * @return the GroupCoordinator */ GroupCoordinator getGroupCoordinator(String tenant); + /** * Access the TransactionCoordinator for the current Tenant. * This method bootstraps a new TransactionCoordinator if it is not started @@ -34,4 +37,12 @@ public interface TenantContextManager { * @return the TransactionCoordinator */ TransactionCoordinator getTransactionCoordinator(String tenant); + + /** + * Access the ReplicaManager for the current Tenant. + * This method bootstraps a new ReplicaManager if it is not started + * @param tenant + * @return the ReplicaManager + */ + ReplicaManager getReplicaManager(String tenant); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java new file mode 100644 index 0000000000..b58da36158 --- /dev/null +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -0,0 +1,356 @@ +/** + * 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.pulsar.handlers.kop.storage; + +import io.netty.buffer.ByteBuf; +import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; +import io.streamnative.pulsar.handlers.kop.KafkaTopicManager; +import io.streamnative.pulsar.handlers.kop.MessagePublishContext; +import io.streamnative.pulsar.handlers.kop.PendingTopicFutures; +import io.streamnative.pulsar.handlers.kop.RequestStats; +import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; +import io.streamnative.pulsar.handlers.kop.format.EncodeRequest; +import io.streamnative.pulsar.handlers.kop.format.EncodeResult; +import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; +import io.streamnative.pulsar.handlers.kop.format.KafkaMixedEntryFormatter; +import io.streamnative.pulsar.handlers.kop.utils.MessageMetadataUtils; +import java.nio.ByteBuffer; +import java.util.Iterator; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.function.Consumer; +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.extern.slf4j.Slf4j; +import org.apache.bookkeeper.common.util.MathUtils; +import org.apache.bookkeeper.mledger.ManagedLedger; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.CorruptRecordException; +import org.apache.kafka.common.errors.RecordTooLargeException; +import org.apache.kafka.common.protocol.Errors; +import org.apache.kafka.common.record.InvalidRecordException; +import org.apache.kafka.common.record.MemoryRecords; +import org.apache.kafka.common.record.MutableRecordBatch; +import org.apache.kafka.common.record.RecordBatch; +import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.common.naming.TopicName; + +@AllArgsConstructor +@Slf4j +public class PartitionLog { + private KafkaServiceConfiguration kafkaConfig; + private TopicPartition topicPartition; + private String namespacePrefix; + private String fullPartitionName; + private EntryFormatter entryFormatter; + + + // A lock that guards all modifications to the log + private final Object lock = new Object(); + + @Data + @AllArgsConstructor + public static final class LogAppendInfo { + private Optional firstOffset; + private Long lastOffset; + private Integer shallowCount; + private Boolean offsetsMonotonic; + private Long lastOffsetOfFirstBatch; + private Integer validBytes; + + public Long numMessages() { + return firstOffset.map(firstOffsetVal -> { + if (firstOffsetVal >= 0 && lastOffset >= 0) { + return lastOffset - firstOffsetVal + 1; + } + return 0L; + }).orElse(0L); + } + } + + public void appendRecords(final MemoryRecords records, + final short version, + final KafkaTopicManager topicManager, + final RequestStats requestStats, + final TransactionCoordinator coordinator, + final Consumer offsetConsumer, + final Consumer errorsConsumer, + final Consumer exceptionConsumer, + final Consumer startSendOperationForThrottlingConsumer, + final Consumer completeSendOperationForThrottlingConsumer, + final Map pendingTopicFuturesMap) { + append(records, + version, + topicManager, + requestStats, + coordinator, + false, + offsetConsumer, + errorsConsumer, + exceptionConsumer, + startSendOperationForThrottlingConsumer, + completeSendOperationForThrottlingConsumer, + pendingTopicFuturesMap); + } + + /** + * Append this message set to the active segment of the log, rolling over to a fresh segment if necessary. + * + * This method will generally be responsible for assigning offsets to the messages, + * however if the assignOffsets=false flag is passed we will only check that the existing offsets are valid. + * + * @param records The log records to append + * @param version Inter-broker message protocol version + * @param ignoreRecordSize true to skip validation of record size. + */ + private void append(final MemoryRecords records, + final short version, + final KafkaTopicManager topicManager, + final RequestStats requestStats, + final TransactionCoordinator coordinator, + final boolean ignoreRecordSize, + final Consumer offsetConsumer, + final Consumer errorsConsumer, + final Consumer exceptionConsumer, + final Consumer startSendOperationForThrottlingConsumer, + final Consumer completeSendOperationForThrottlingConsumer, + final Map pendingTopicFuturesMap + ) { + final long beforeRecordsProcess = MathUtils.nowInNano(); + final LogAppendInfo appendInfo = analyzeAndValidateRecords(records, version, topicPartition, ignoreRecordSize); + + // return if we have no valid messages or if this is a duplicate of the last appended entry + if (appendInfo.getShallowCount() == 0) { + return; + } + // trim any invalid bytes or partial messages before appending it to the on-disk log + MemoryRecords validRecords = trimInvalidBytes(records, appendInfo); + synchronized (lock) { + long offset = -1; + appendInfo.setFirstOffset(Optional.of(offset)); + // TODO: validateMessagesAndAssignOffsets + + // Append Message into pulsar + final CompletableFuture> topicFuture = + topicManager.getTopic(fullPartitionName); + if (topicFuture.isCompletedExceptionally()) { + topicFuture.exceptionally(e -> { + exceptionConsumer.accept(e); + return Optional.empty(); + }); + return; + } + if (topicFuture.isDone() && !topicFuture.getNow(Optional.empty()).isPresent()) { + errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); + return; + } + final Consumer> persistentTopicConsumer = persistentTopicOpt -> { + if (!persistentTopicOpt.isPresent()) { + errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); + return; + } + + final EncodeRequest encodeRequest = EncodeRequest.get(validRecords); + if (entryFormatter instanceof KafkaMixedEntryFormatter) { + final ManagedLedger managedLedger = persistentTopicOpt.get().getManagedLedger(); + final long logEndOffset = MessageMetadataUtils.getLogEndOffset(managedLedger); + encodeRequest.setBaseOffset(logEndOffset); + } + + final EncodeResult encodeResult = entryFormatter.encode(encodeRequest); + encodeRequest.recycle(); + requestStats.getProduceEncodeStats().registerSuccessfulEvent( + MathUtils.elapsedNanos(beforeRecordsProcess), TimeUnit.NANOSECONDS); +// startSendOperationForThrottling(encodeResult.getEncodedByteBuf().readableBytes()); + startSendOperationForThrottlingConsumer.accept(encodeResult.getEncodedByteBuf().readableBytes()); + if (log.isDebugEnabled()) { + log.debug("Produce messages for topic {} partition {}", + topicPartition.topic(), topicPartition.partition()); + } + + publishMessages(persistentTopicOpt, + topicManager, + coordinator, + requestStats, + encodeResult, + topicPartition, + offsetConsumer, + errorsConsumer, + completeSendOperationForThrottlingConsumer); + }; + + if (topicFuture.isDone()) { + persistentTopicConsumer.accept(topicFuture.getNow(Optional.empty())); + } else { + // topic is not available now + pendingTopicFuturesMap + .computeIfAbsent(topicPartition, ignored -> + new PendingTopicFutures(requestStats)) + .addListener(topicFuture, persistentTopicConsumer, exceptionConsumer); + } + } + } + + private void publishMessages(final Optional persistentTopicOpt, + final KafkaTopicManager topicManager, + final TransactionCoordinator coordinator, + final RequestStats requestStats, + final EncodeResult encodeResult, + final TopicPartition topicPartition, + final Consumer offsetConsumer, + final Consumer errorsConsumer, + final Consumer completeSendOperationForThrottlingConsumer) { + final MemoryRecords records = encodeResult.getRecords(); + final int numMessages = encodeResult.getNumMessages(); + final ByteBuf byteBuf = encodeResult.getEncodedByteBuf(); + if (!persistentTopicOpt.isPresent()) { + encodeResult.recycle(); + // It will trigger a retry send of Kafka client + errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); + return; + } + PersistentTopic persistentTopic = persistentTopicOpt.get(); + if (persistentTopic.isSystemTopic()) { + encodeResult.recycle(); + log.error("Not support producing message to system topic: {}", persistentTopic); + errorsConsumer.accept(Errors.INVALID_TOPIC_EXCEPTION); + return; + } + + topicManager.registerProducerInPersistentTopic(fullPartitionName, persistentTopic); + + // collect metrics + encodeResult.updateProducerStats(topicPartition, requestStats, namespacePrefix); + + // publish + final CompletableFuture offsetFuture = new CompletableFuture<>(); + final long beforePublish = MathUtils.nowInNano(); + persistentTopic.publishMessage(byteBuf, + MessagePublishContext.get(offsetFuture, persistentTopic, numMessages, System.nanoTime())); + final RecordBatch batch = records.batchIterator().next(); + offsetFuture.whenComplete((offset, e) -> { +// completeSendOperationForThrottling(byteBuf.readableBytes()); + completeSendOperationForThrottlingConsumer.accept(byteBuf.readableBytes()); + encodeResult.recycle(); + if (e == null) { + if (batch.isTransactional()) { + coordinator.addActivePidOffset(TopicName.get(fullPartitionName), batch.producerId(), + offset); + } + requestStats.getMessagePublishStats().registerSuccessfulEvent( + MathUtils.elapsedNanos(beforePublish), TimeUnit.NANOSECONDS); + offsetConsumer.accept(offset); + } else { + log.error("publishMessages for topic partition: {} failed when write.", fullPartitionName, e); + requestStats.getMessagePublishStats().registerFailedEvent( + MathUtils.elapsedNanos(beforePublish), TimeUnit.NANOSECONDS); + errorsConsumer.accept(Errors.KAFKA_STORAGE_ERROR); + } + }); + } + + private LogAppendInfo analyzeAndValidateRecords(MemoryRecords records, + short version, + TopicPartition topicPartition, + boolean ignoreRecordSize) { + int shallowMessageCount = 0; + long lastOffset = -1L; + Optional firstOffset = Optional.empty(); + long lastOffsetOfFirstBatch = -1L; + boolean readFirstMessage = false; + boolean monotonic = true; + + if (version >= 3) { + Iterator iterator = records.batches().iterator(); + if (!iterator.hasNext()) { + throw new InvalidRecordException("Produce requests with version " + version + " must have at least " + + "one record batch"); + } + + MutableRecordBatch entry = iterator.next(); + if (entry.magic() != RecordBatch.MAGIC_VALUE_V2) { + throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " + + "contain record batches with magic version 2"); + } + + if (iterator.hasNext()) { + throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " + + "contain exactly one record batch"); + } + } + + int validBytesCount = 0; + for (RecordBatch batch : records.batches()) { + if (batch.magic() >= RecordBatch.MAGIC_VALUE_V2 && batch.baseOffset() != 0) { + throw new InvalidRecordException("The baseOffset of the record batch in the append to " + + topicPartition + " should be 0, but it is " + batch.baseOffset()); + } + if (!readFirstMessage) { + if (batch.magic() >= RecordBatch.MAGIC_VALUE_V2) { + firstOffset = Optional.of(batch.baseOffset()); + } + lastOffsetOfFirstBatch = batch.lastOffset(); + readFirstMessage = true; + } + // check that offsets are monotonically increasing + if (lastOffset >= batch.lastOffset()){ + monotonic = false; + } + + // update the last offset seen + lastOffset = batch.lastOffset(); + + int batchSize = batch.sizeInBytes(); + if (!ignoreRecordSize && batchSize > kafkaConfig.getMaxMessageSize()) { + throw new RecordTooLargeException(String.format("Message batch size is %s " + + "in append to partition %s which exceeds the maximum configured size of %s .", + batchSize, topicPartition, kafkaConfig.getMaxMessageSize())); + } + + batch.ensureValid(); + shallowMessageCount += 1; + validBytesCount += batchSize; + } + + if (validBytesCount < 0) { + throw new CorruptRecordException("Cannot append record batch with illegal length " + + validBytesCount + " to log for " + topicPartition + + ". A possible cause is corrupted produce request."); + } + + return new LogAppendInfo( + firstOffset, + lastOffset, + shallowMessageCount, + monotonic, + lastOffsetOfFirstBatch, + validBytesCount); + } + + private MemoryRecords trimInvalidBytes(MemoryRecords records, LogAppendInfo info) { + Integer validBytes = info.getValidBytes(); + if (validBytes < 0){ + throw new CorruptRecordException(String.format("Cannot append record batch with illegal length %s to " + + "log for %s. A possible cause is a corrupted produce request.", validBytes, topicPartition)); + } else if (validBytes == records.sizeInBytes()) { + return records; + } else { + ByteBuffer validByteBuffer = records.buffer().duplicate(); + validByteBuffer.limit(validBytes); + return MemoryRecords.readableRecords(validByteBuffer); + } + } +} diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java new file mode 100644 index 0000000000..7e50c592bc --- /dev/null +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java @@ -0,0 +1,49 @@ +/** + * 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.pulsar.handlers.kop.storage; + +import com.google.common.collect.Maps; +import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; +import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; +import io.streamnative.pulsar.handlers.kop.format.EntryFormatterFactory; +import io.streamnative.pulsar.handlers.kop.utils.KopTopic; +import java.util.Map; +import lombok.AllArgsConstructor; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.utils.Time; + +@AllArgsConstructor +public class PartitionLogManager { + + private final Map logMap; + private final KafkaServiceConfiguration config; + private final EntryFormatter formatter; + private final Time time; + + public PartitionLogManager(KafkaServiceConfiguration config, + Time time) { + this.logMap = Maps.newConcurrentMap(); + this.formatter = EntryFormatterFactory.create(config); + this.config = config; + this.time = time; + } + + public PartitionLog getLog(TopicPartition topicPartition, String namespacePrefix) { + String kopTopic = KopTopic.toString(topicPartition, namespacePrefix); + return logMap.computeIfAbsent(kopTopic, key -> + new PartitionLog(config, topicPartition, namespacePrefix, kopTopic, formatter) + ); + } +} + diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java new file mode 100644 index 0000000000..a0dc795333 --- /dev/null +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -0,0 +1,152 @@ +/** + * 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.pulsar.handlers.kop.storage; + +import io.streamnative.pulsar.handlers.kop.DelayedProduceAndFetch; +import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; +import io.streamnative.pulsar.handlers.kop.KafkaTopicManager; +import io.streamnative.pulsar.handlers.kop.PendingTopicFutures; +import io.streamnative.pulsar.handlers.kop.RequestStats; +import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; +import io.streamnative.pulsar.handlers.kop.utils.KopTopic; +import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; +import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationKey; +import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.BiConsumer; +import java.util.function.Consumer; +import java.util.stream.Collectors; +import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.InvalidTopicException; +import org.apache.kafka.common.protocol.Errors; +import org.apache.kafka.common.record.MemoryRecords; +import org.apache.kafka.common.requests.ProduceResponse; +import org.apache.kafka.common.utils.Time; + +@Slf4j +public class ReplicaManager { + private final PartitionLogManager logManager; + private final DelayedOperationPurgatory producePurgatory; + private final DelayedOperationPurgatory fetchPurgatory; + + public ReplicaManager(KafkaServiceConfiguration config, + Time time, + DelayedOperationPurgatory producePurgatory, + DelayedOperationPurgatory fetchPurgatory) { + this.logManager = new PartitionLogManager(config, time); + this.producePurgatory = producePurgatory; + this.fetchPurgatory = fetchPurgatory; + } + + public PartitionLog getPartitionLog(TopicPartition topicPartition, String namespacePrefix) { + return logManager.getLog(topicPartition, namespacePrefix); + } + + public void appendRecords( + final long timeout, + final boolean internalTopicsAllowed, + final short version, + final KafkaTopicManager topicManager, + final String namespacePrefix, + final Map entriesPerPartition, + final RequestStats requestStats, + final TransactionCoordinator coordinator, + final Consumer startSendOperationForThrottlingConsumer, + final Consumer completeSendOperationForThrottlingConsumer, + final Map pendingTopicFuturesMap, + final CompletableFuture> responseCallback) { + + final AtomicInteger topicPartitionNum = new AtomicInteger(entriesPerPartition.size()); + final Map responseMap = new ConcurrentHashMap<>(); + + Runnable complete = () -> { + topicPartitionNum.set(0); + if (responseCallback.isDone()) { + // It may be triggered again in DelayedProduceAndFetch + return; + } + // add the topicPartition with timeout error if it's not existed in responseMap + entriesPerPartition.keySet().forEach(topicPartition -> { + if (!responseMap.containsKey(topicPartition)) { + responseMap.put(topicPartition, new ProduceResponse.PartitionResponse(Errors.REQUEST_TIMED_OUT)); + } + }); + if (log.isDebugEnabled()) { + log.debug("Complete handle appendRecords."); + } + responseCallback.complete(responseMap); + }; + BiConsumer addPartitionResponse = + (topicPartition, response) -> { + responseMap.put(topicPartition, response); + // reset topicPartitionNum + int restTopicPartitionNum = topicPartitionNum.decrementAndGet(); + if (restTopicPartitionNum < 0) { + return; + } + if (restTopicPartitionNum == 0) { + complete.run(); + } + }; + entriesPerPartition.forEach((topicPartition, memoryRecords) -> { + final Consumer offsetConsumer = offset -> addPartitionResponse.accept( + topicPartition, + new ProduceResponse.PartitionResponse(Errors.NONE, offset, -1L, -1L)); + final Consumer errorsConsumer = + errors -> addPartitionResponse + .accept(topicPartition, new ProduceResponse.PartitionResponse(errors)); + final Consumer exceptionConsumer = + e -> addPartitionResponse + .accept(topicPartition, new ProduceResponse.PartitionResponse(Errors.forException(e))); + + String fullPartitionName = KopTopic.toString(topicPartition, namespacePrefix); + // reject appending to internal topics if it is not allowed + if (!internalTopicsAllowed && KopTopic.isInternalTopic(fullPartitionName)) { + exceptionConsumer.accept(new InvalidTopicException( + String.format("Cannot append to internal topic %s", topicPartition.topic()))); + } else { + PartitionLog partitionLog = getPartitionLog(topicPartition, namespacePrefix); + partitionLog.appendRecords(memoryRecords, + version, + topicManager, + requestStats, + coordinator, + offsetConsumer, + errorsConsumer, + exceptionConsumer, + startSendOperationForThrottlingConsumer, + completeSendOperationForThrottlingConsumer, + pendingTopicFuturesMap); + } + + }); + // delay produce + if (timeout <= 0) { + complete.run(); + } else { + List delayedCreateKeys = + entriesPerPartition.keySet().stream() + .map(DelayedOperationKey.TopicPartitionOperationKey::new).collect(Collectors.toList()); + DelayedProduceAndFetch delayedProduce = new DelayedProduceAndFetch(timeout, topicPartitionNum, complete); + producePurgatory.tryCompleteElseWatch(delayedProduce, delayedCreateKeys); + } + + } + +} diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/package-info.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/package-info.java new file mode 100644 index 0000000000..ad5d79b79c --- /dev/null +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/package-info.java @@ -0,0 +1,14 @@ +/** + * 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.pulsar.handlers.kop.storage; \ No newline at end of file diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/KopTopic.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/KopTopic.java index 1f5f4a71dd..fe75a7f654 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/KopTopic.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/KopTopic.java @@ -13,11 +13,14 @@ */ package io.streamnative.pulsar.handlers.kop.utils; +import static org.apache.kafka.common.internals.Topic.GROUP_METADATA_TOPIC_NAME; +import static org.apache.kafka.common.internals.Topic.TRANSACTION_STATE_TOPIC_NAME; import static org.apache.pulsar.common.naming.TopicName.PARTITIONED_TOPIC_SUFFIX; import io.streamnative.pulsar.handlers.kop.exceptions.KoPTopicException; import lombok.Getter; import org.apache.kafka.common.TopicPartition; +import org.apache.pulsar.common.naming.TopicName; /** * KopTopic maintains two topic name, one is the original topic name, the other is the full topic name used in Pulsar. @@ -103,4 +106,10 @@ public static String toString(String topic, int partition, String namespacePrefi return (new KopTopic(topic, namespacePrefix)).getPartitionName(partition); } + public static boolean isInternalTopic(final String fullTopicName) { + String partitionedTopicName = TopicName.get(fullTopicName).getPartitionedTopicName(); + return partitionedTopicName.endsWith("/" + GROUP_METADATA_TOPIC_NAME) + || partitionedTopicName.endsWith("/" + TRANSACTION_STATE_TOPIC_NAME); + } + } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java index e32de60c8d..e37f4011d1 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java @@ -26,6 +26,7 @@ import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; +import io.streamnative.pulsar.handlers.kop.storage.ReplicaManager; import io.streamnative.pulsar.handlers.kop.utils.MetadataUtils; import java.io.Closeable; import java.io.IOException; @@ -768,6 +769,8 @@ public KafkaRequestHandler newRequestHandler() throws Exception { final GroupCoordinator groupCoordinator = handler.getGroupCoordinator(conf.getKafkaMetadataTenant()); final TransactionCoordinator transactionCoordinator = handler.getTransactionCoordinator(conf.getKafkaMetadataTenant()); + final ReplicaManager replicaManager = + handler.getReplicaManager(conf.getKafkaMetadataTenant()); return ((KafkaChannelInitializer) handler.getChannelInitializerMap().entrySet().iterator().next().getValue()) .newCnx(new TenantContextManager() { @@ -780,6 +783,11 @@ public GroupCoordinator getGroupCoordinator(String tenant) { public TransactionCoordinator getTransactionCoordinator(String tenant) { return transactionCoordinator; } + + @Override + public ReplicaManager getReplicaManager(String tenant) { + return replicaManager; + } }, NullStatsLogger.INSTANCE); } } From 79ce3bacda4f0881ea710f96935bd624b935cc0f Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Tue, 23 Nov 2021 21:49:05 +0800 Subject: [PATCH 02/15] Fix units test --- .../handlers/kop/KafkaRequestHandler.java | 9 +- .../handlers/kop/storage/PartitionLog.java | 126 +++++++++--------- 2 files changed, 73 insertions(+), 62 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index a683e11e0f..e8b60d09a9 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -959,7 +959,10 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, ProduceRequest produceRequest = (ProduceRequest) produceHar.getRequest(); final int numPartitions = produceRequest.partitionRecordsOrFail().size(); - + if (numPartitions == 0) { + resultFuture.complete(new ProduceResponse(Collections.emptyMap())); + return; + } final Map unauthorizedTopicResponsesMap = new ConcurrentHashMap<>(); final Map authorizedRequestInfo = new ConcurrentHashMap<>(); int timeoutMs = produceRequest.timeout(); @@ -969,6 +972,10 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, // When complete one authorization or failed, will do the action first. action.run(); if (unfinishedAuthorizationCount.decrementAndGet() == 0) { + if (authorizedRequestInfo.isEmpty()) { + resultFuture.complete(new ProduceResponse(unauthorizedTopicResponsesMap)); + return; + } CompletableFuture> responseCallback = new CompletableFuture<>(); getReplicaManager().appendRecords( diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index b58da36158..5889a4fd39 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -130,78 +130,83 @@ private void append(final MemoryRecords records, final Map pendingTopicFuturesMap ) { final long beforeRecordsProcess = MathUtils.nowInNano(); - final LogAppendInfo appendInfo = analyzeAndValidateRecords(records, version, topicPartition, ignoreRecordSize); + try { + final LogAppendInfo appendInfo = analyzeAndValidateRecords(records, version, topicPartition, ignoreRecordSize); - // return if we have no valid messages or if this is a duplicate of the last appended entry - if (appendInfo.getShallowCount() == 0) { - return; - } - // trim any invalid bytes or partial messages before appending it to the on-disk log - MemoryRecords validRecords = trimInvalidBytes(records, appendInfo); - synchronized (lock) { - long offset = -1; - appendInfo.setFirstOffset(Optional.of(offset)); - // TODO: validateMessagesAndAssignOffsets - - // Append Message into pulsar - final CompletableFuture> topicFuture = - topicManager.getTopic(fullPartitionName); - if (topicFuture.isCompletedExceptionally()) { - topicFuture.exceptionally(e -> { - exceptionConsumer.accept(e); - return Optional.empty(); - }); + // return if we have no valid messages or if this is a duplicate of the last appended entry + if (appendInfo.getShallowCount() == 0) { return; } - if (topicFuture.isDone() && !topicFuture.getNow(Optional.empty()).isPresent()) { - errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); - return; - } - final Consumer> persistentTopicConsumer = persistentTopicOpt -> { - if (!persistentTopicOpt.isPresent()) { + // trim any invalid bytes or partial messages before appending it to the on-disk log + MemoryRecords validRecords = trimInvalidBytes(records, appendInfo); + synchronized (lock) { + long offset = -1; + appendInfo.setFirstOffset(Optional.of(offset)); + // TODO: validateMessagesAndAssignOffsets + + // Append Message into pulsar + final CompletableFuture> topicFuture = + topicManager.getTopic(fullPartitionName); + if (topicFuture.isCompletedExceptionally()) { + topicFuture.exceptionally(e -> { + exceptionConsumer.accept(e); + return Optional.empty(); + }); + return; + } + if (topicFuture.isDone() && !topicFuture.getNow(Optional.empty()).isPresent()) { errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); return; } + final Consumer> persistentTopicConsumer = persistentTopicOpt -> { + if (!persistentTopicOpt.isPresent()) { + errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); + return; + } - final EncodeRequest encodeRequest = EncodeRequest.get(validRecords); - if (entryFormatter instanceof KafkaMixedEntryFormatter) { - final ManagedLedger managedLedger = persistentTopicOpt.get().getManagedLedger(); - final long logEndOffset = MessageMetadataUtils.getLogEndOffset(managedLedger); - encodeRequest.setBaseOffset(logEndOffset); - } + final EncodeRequest encodeRequest = EncodeRequest.get(validRecords); + if (entryFormatter instanceof KafkaMixedEntryFormatter) { + final ManagedLedger managedLedger = persistentTopicOpt.get().getManagedLedger(); + final long logEndOffset = MessageMetadataUtils.getLogEndOffset(managedLedger); + encodeRequest.setBaseOffset(logEndOffset); + } - final EncodeResult encodeResult = entryFormatter.encode(encodeRequest); - encodeRequest.recycle(); - requestStats.getProduceEncodeStats().registerSuccessfulEvent( - MathUtils.elapsedNanos(beforeRecordsProcess), TimeUnit.NANOSECONDS); -// startSendOperationForThrottling(encodeResult.getEncodedByteBuf().readableBytes()); - startSendOperationForThrottlingConsumer.accept(encodeResult.getEncodedByteBuf().readableBytes()); - if (log.isDebugEnabled()) { - log.debug("Produce messages for topic {} partition {}", - topicPartition.topic(), topicPartition.partition()); - } + final EncodeResult encodeResult = entryFormatter.encode(encodeRequest); + encodeRequest.recycle(); + requestStats.getProduceEncodeStats().registerSuccessfulEvent( + MathUtils.elapsedNanos(beforeRecordsProcess), TimeUnit.NANOSECONDS); + startSendOperationForThrottlingConsumer.accept(encodeResult.getEncodedByteBuf().readableBytes()); + if (log.isDebugEnabled()) { + log.debug("Produce messages for topic {} partition {}", + topicPartition.topic(), topicPartition.partition()); + } - publishMessages(persistentTopicOpt, - topicManager, - coordinator, - requestStats, - encodeResult, - topicPartition, - offsetConsumer, - errorsConsumer, - completeSendOperationForThrottlingConsumer); - }; + publishMessages(persistentTopicOpt, + topicManager, + coordinator, + requestStats, + encodeResult, + topicPartition, + offsetConsumer, + errorsConsumer, + completeSendOperationForThrottlingConsumer); + }; - if (topicFuture.isDone()) { - persistentTopicConsumer.accept(topicFuture.getNow(Optional.empty())); - } else { - // topic is not available now - pendingTopicFuturesMap - .computeIfAbsent(topicPartition, ignored -> - new PendingTopicFutures(requestStats)) - .addListener(topicFuture, persistentTopicConsumer, exceptionConsumer); + if (topicFuture.isDone()) { + persistentTopicConsumer.accept(topicFuture.getNow(Optional.empty())); + } else { + // topic is not available now + pendingTopicFuturesMap + .computeIfAbsent(topicPartition, ignored -> + new PendingTopicFutures(requestStats)) + .addListener(topicFuture, persistentTopicConsumer, exceptionConsumer); + } } + } catch (Exception exception) { + log.error("Failed to handle produce request for {}", topicPartition, exception); + exceptionConsumer.accept(exception); } + } private void publishMessages(final Optional persistentTopicOpt, @@ -242,7 +247,6 @@ private void publishMessages(final Optional persistentTopicOpt, MessagePublishContext.get(offsetFuture, persistentTopic, numMessages, System.nanoTime())); final RecordBatch batch = records.batchIterator().next(); offsetFuture.whenComplete((offset, e) -> { -// completeSendOperationForThrottling(byteBuf.readableBytes()); completeSendOperationForThrottlingConsumer.accept(byteBuf.readableBytes()); encodeResult.recycle(); if (e == null) { From 1d6fbed8c4b3602683aa7cac274c1965c862c6db Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Tue, 23 Nov 2021 22:30:46 +0800 Subject: [PATCH 03/15] Fix code style --- .../handlers/kop/KafkaProtocolHandler.java | 10 +- .../handlers/kop/KafkaRequestHandler.java | 148 ------------------ .../handlers/kop/storage/PartitionLog.java | 34 ++-- .../kop/storage/PartitionLogManager.java | 8 +- .../handlers/kop/storage/ReplicaManager.java | 6 +- 5 files changed, 33 insertions(+), 173 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index 2bdd29698a..73c9460d23 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -118,8 +118,16 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) { @Override public ReplicaManager getReplicaManager(String tenant) { return replicaManagerByTenant.computeIfAbsent(tenant, s -> { + Optional transactionCoordinatorOptional = Optional.empty(); + if (kafkaConfig.isEnableTransactionCoordinator()) { + transactionCoordinatorOptional = Optional.of(getTransactionCoordinator(tenant)); + } try { - return new ReplicaManager(kafkaConfig, Time.SYSTEM, producePurgatory, fetchPurgatory); + return new ReplicaManager(kafkaConfig, + Time.SYSTEM, + transactionCoordinatorOptional, + producePurgatory, + fetchPurgatory); } catch (Exception e) { log.error("Failed to init ReplicaManager for tenant {}", tenant, e); throw new IllegalStateException(e); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index e8b60d09a9..271c3960b9 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -35,11 +35,8 @@ import io.streamnative.pulsar.handlers.kop.coordinator.transaction.AbortedIndexEntry; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.exceptions.KoPTopicException; -import io.streamnative.pulsar.handlers.kop.format.EncodeRequest; -import io.streamnative.pulsar.handlers.kop.format.EncodeResult; import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; import io.streamnative.pulsar.handlers.kop.format.EntryFormatterFactory; -import io.streamnative.pulsar.handlers.kop.format.KafkaMixedEntryFormatter; import io.streamnative.pulsar.handlers.kop.offset.OffsetAndMetadata; import io.streamnative.pulsar.handlers.kop.offset.OffsetMetadata; import io.streamnative.pulsar.handlers.kop.security.SaslAuthenticator; @@ -86,9 +83,7 @@ import java.util.stream.IntStream; import lombok.Getter; import lombok.extern.slf4j.Slf4j; -import org.apache.bookkeeper.common.util.MathUtils; import org.apache.bookkeeper.mledger.AsyncCallbacks; -import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; @@ -103,7 +98,6 @@ import org.apache.kafka.common.errors.AuthenticationException; import org.apache.kafka.common.errors.CorruptRecordException; import org.apache.kafka.common.errors.LeaderNotAvailableException; -import org.apache.kafka.common.errors.RecordTooLargeException; import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.record.ControlRecordType; @@ -898,60 +892,6 @@ private void completeSendOperationForThrottling(long msgSize) { } } - private void publishMessages(final Optional persistentTopicOpt, - final EncodeResult encodeResult, - final TopicPartition topicPartition, - final Consumer offsetConsumer, - final Consumer errorsConsumer) { - final MemoryRecords records = encodeResult.getRecords(); - final int numMessages = encodeResult.getNumMessages(); - final ByteBuf byteBuf = encodeResult.getEncodedByteBuf(); - if (!persistentTopicOpt.isPresent()) { - encodeResult.recycle(); - // It will trigger a retry send of Kafka client - errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); - return; - } - PersistentTopic persistentTopic = persistentTopicOpt.get(); - if (persistentTopic.isSystemTopic()) { - encodeResult.recycle(); - log.error("Not support producing message to system topic: {}", persistentTopic); - errorsConsumer.accept(Errors.INVALID_TOPIC_EXCEPTION); - return; - } - String namespacePrefix = currentNamespacePrefix(); - final String partitionName = KopTopic.toString(topicPartition, namespacePrefix); - - topicManager.registerProducerInPersistentTopic(partitionName, persistentTopic); - // collect metrics - encodeResult.updateProducerStats(topicPartition, requestStats, namespacePrefix); - - // publish - final CompletableFuture offsetFuture = new CompletableFuture<>(); - final long beforePublish = MathUtils.nowInNano(); - persistentTopic.publishMessage(byteBuf, - MessagePublishContext.get(offsetFuture, persistentTopic, numMessages, System.nanoTime())); - final RecordBatch batch = records.batchIterator().next(); - offsetFuture.whenComplete((offset, e) -> { - completeSendOperationForThrottling(byteBuf.readableBytes()); - encodeResult.recycle(); - if (e == null) { - if (batch.isTransactional()) { - getTransactionCoordinator().addActivePidOffset(TopicName.get(partitionName), batch.producerId(), - offset); - } - requestStats.getMessagePublishStats().registerSuccessfulEvent( - MathUtils.elapsedNanos(beforePublish), TimeUnit.NANOSECONDS); - offsetConsumer.accept(offset); - } else { - log.error("publishMessages for topic partition: {} failed when write.", partitionName, e); - requestStats.getMessagePublishStats().registerFailedEvent( - MathUtils.elapsedNanos(beforePublish), TimeUnit.NANOSECONDS); - errorsConsumer.accept(Errors.KAFKA_STORAGE_ERROR); - } - }); - } - @Override protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, CompletableFuture resultFuture) { @@ -986,7 +926,6 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, namespacePrefix, authorizedRequestInfo, requestStats, - getTransactionCoordinator(), this::startSendOperationForThrottling, this::completeSendOperationForThrottling, pendingTopicFuturesMap, @@ -1033,93 +972,6 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, } - private void handlePartitionRecords(final KafkaHeaderAndRequest produceHar, - final TopicPartition topicPartition, - final MemoryRecords records, - final int numPartitions, - final String fullPartitionName, - final Consumer offsetConsumer, - final Consumer errorsConsumer, - final Consumer exceptionConsumer) { - // check KOP inner topic - if (isInternalTopic(fullPartitionName)) { - log.error("[{}] Request {}: not support produce message to inner topic. topic: {}", - ctx.channel(), produceHar.getHeader(), topicPartition); - errorsConsumer.accept(Errors.INVALID_TOPIC_EXCEPTION); - return; - } - - try { - final long beforeRecordsProcess = MathUtils.nowInNano(); - final MemoryRecords validRecords = - validateRecords(produceHar.getHeader().apiVersion(), topicPartition, records); - - validRecords.batches().forEach(batch->{ - if (batch.sizeInBytes() > kafkaConfig.getMaxMessageSize()) { - throw new RecordTooLargeException(String.format("Message batch size is %s " - + "in append to partition %s which exceeds the maximum configured size of %s .", - batch.sizeInBytes(), topicPartition, kafkaConfig.getMaxMessageSize())); - } - }); - - final CompletableFuture> topicFuture = - topicManager.getTopic(fullPartitionName); - if (topicFuture.isCompletedExceptionally()) { - topicFuture.exceptionally(e -> { - exceptionConsumer.accept(e); - return Optional.empty(); - }); - return; - } - if (topicFuture.isDone() && !topicFuture.getNow(Optional.empty()).isPresent()) { - errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); - return; - } - - final Consumer> persistentTopicConsumer = persistentTopicOpt -> { - if (!persistentTopicOpt.isPresent()) { - errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); - return; - } - - final EncodeRequest encodeRequest = EncodeRequest.get(validRecords); - if (entryFormatter instanceof KafkaMixedEntryFormatter) { - final ManagedLedger managedLedger = persistentTopicOpt.get().getManagedLedger(); - final long logEndOffset = MessageMetadataUtils.getLogEndOffset(managedLedger); - encodeRequest.setBaseOffset(logEndOffset); - } - - final EncodeResult encodeResult = entryFormatter.encode(encodeRequest); - encodeRequest.recycle(); - requestStats.getProduceEncodeStats().registerSuccessfulEvent( - MathUtils.elapsedNanos(beforeRecordsProcess), TimeUnit.NANOSECONDS); - startSendOperationForThrottling(encodeResult.getEncodedByteBuf().readableBytes()); - - if (log.isDebugEnabled()) { - log.debug("[{}] Request {}: Produce messages for topic {} partition {}, " - + "request size: {} ", ctx.channel(), produceHar.getHeader(), - topicPartition.topic(), topicPartition.partition(), numPartitions); - } - - publishMessages(persistentTopicOpt, encodeResult, topicPartition, offsetConsumer, errorsConsumer); - }; - - if (topicFuture.isDone()) { - persistentTopicConsumer.accept(topicFuture.getNow(Optional.empty())); - } else { - // topic is not available now - pendingTopicFuturesMap - .computeIfAbsent(topicPartition, ignored -> - new PendingTopicFutures(requestStats)) - .addListener(topicFuture, persistentTopicConsumer, exceptionConsumer); - } - } catch (Exception e) { - log.error("[{}] Failed to handle produce request for {}", - ctx.channel(), topicPartition, e); - exceptionConsumer.accept(e); - } - } - @Override protected void handleFindCoordinatorRequest(KafkaHeaderAndRequest findCoordinator, CompletableFuture resultFuture) { diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index 5889a4fd39..b74ca22817 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -35,7 +35,6 @@ import lombok.AllArgsConstructor; import lombok.Data; import lombok.extern.slf4j.Slf4j; -import org.apache.bookkeeper.common.util.MathUtils; import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.CorruptRecordException; @@ -45,6 +44,7 @@ import org.apache.kafka.common.record.MemoryRecords; import org.apache.kafka.common.record.MutableRecordBatch; import org.apache.kafka.common.record.RecordBatch; +import org.apache.kafka.common.utils.Time; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.common.naming.TopicName; @@ -52,10 +52,12 @@ @Slf4j public class PartitionLog { private KafkaServiceConfiguration kafkaConfig; + private Time time; private TopicPartition topicPartition; private String namespacePrefix; private String fullPartitionName; private EntryFormatter entryFormatter; + private Optional transactionCoordinator; // A lock that guards all modifications to the log @@ -85,7 +87,6 @@ public void appendRecords(final MemoryRecords records, final short version, final KafkaTopicManager topicManager, final RequestStats requestStats, - final TransactionCoordinator coordinator, final Consumer offsetConsumer, final Consumer errorsConsumer, final Consumer exceptionConsumer, @@ -96,7 +97,6 @@ public void appendRecords(final MemoryRecords records, version, topicManager, requestStats, - coordinator, false, offsetConsumer, errorsConsumer, @@ -120,18 +120,17 @@ private void append(final MemoryRecords records, final short version, final KafkaTopicManager topicManager, final RequestStats requestStats, - final TransactionCoordinator coordinator, final boolean ignoreRecordSize, final Consumer offsetConsumer, final Consumer errorsConsumer, final Consumer exceptionConsumer, final Consumer startSendOperationForThrottlingConsumer, final Consumer completeSendOperationForThrottlingConsumer, - final Map pendingTopicFuturesMap - ) { - final long beforeRecordsProcess = MathUtils.nowInNano(); + final Map pendingTopicFuturesMap) { + final long beforeRecordsProcess = time.nanoseconds(); try { - final LogAppendInfo appendInfo = analyzeAndValidateRecords(records, version, topicPartition, ignoreRecordSize); + final LogAppendInfo appendInfo = + analyzeAndValidateRecords(records, version, topicPartition, ignoreRecordSize); // return if we have no valid messages or if this is a duplicate of the last appended entry if (appendInfo.getShallowCount() == 0) { @@ -140,10 +139,6 @@ private void append(final MemoryRecords records, // trim any invalid bytes or partial messages before appending it to the on-disk log MemoryRecords validRecords = trimInvalidBytes(records, appendInfo); synchronized (lock) { - long offset = -1; - appendInfo.setFirstOffset(Optional.of(offset)); - // TODO: validateMessagesAndAssignOffsets - // Append Message into pulsar final CompletableFuture> topicFuture = topicManager.getTopic(fullPartitionName); @@ -163,6 +158,7 @@ private void append(final MemoryRecords records, errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); return; } + // TODO: validateMessagesAndAssignOffsets here. final EncodeRequest encodeRequest = EncodeRequest.get(validRecords); if (entryFormatter instanceof KafkaMixedEntryFormatter) { @@ -174,7 +170,7 @@ private void append(final MemoryRecords records, final EncodeResult encodeResult = entryFormatter.encode(encodeRequest); encodeRequest.recycle(); requestStats.getProduceEncodeStats().registerSuccessfulEvent( - MathUtils.elapsedNanos(beforeRecordsProcess), TimeUnit.NANOSECONDS); + time.nanoseconds() - beforeRecordsProcess, TimeUnit.NANOSECONDS); startSendOperationForThrottlingConsumer.accept(encodeResult.getEncodedByteBuf().readableBytes()); if (log.isDebugEnabled()) { log.debug("Produce messages for topic {} partition {}", @@ -183,7 +179,6 @@ private void append(final MemoryRecords records, publishMessages(persistentTopicOpt, topicManager, - coordinator, requestStats, encodeResult, topicPartition, @@ -211,7 +206,6 @@ private void append(final MemoryRecords records, private void publishMessages(final Optional persistentTopicOpt, final KafkaTopicManager topicManager, - final TransactionCoordinator coordinator, final RequestStats requestStats, final EncodeResult encodeResult, final TopicPartition topicPartition, @@ -242,7 +236,7 @@ private void publishMessages(final Optional persistentTopicOpt, // publish final CompletableFuture offsetFuture = new CompletableFuture<>(); - final long beforePublish = MathUtils.nowInNano(); + final long beforePublish = time.nanoseconds(); persistentTopic.publishMessage(byteBuf, MessagePublishContext.get(offsetFuture, persistentTopic, numMessages, System.nanoTime())); final RecordBatch batch = records.batchIterator().next(); @@ -251,16 +245,16 @@ private void publishMessages(final Optional persistentTopicOpt, encodeResult.recycle(); if (e == null) { if (batch.isTransactional()) { - coordinator.addActivePidOffset(TopicName.get(fullPartitionName), batch.producerId(), - offset); + transactionCoordinator.ifPresent(coordinator -> coordinator.addActivePidOffset( + TopicName.get(fullPartitionName), batch.producerId(), offset)); } requestStats.getMessagePublishStats().registerSuccessfulEvent( - MathUtils.elapsedNanos(beforePublish), TimeUnit.NANOSECONDS); + time.nanoseconds() - beforePublish, TimeUnit.NANOSECONDS); offsetConsumer.accept(offset); } else { log.error("publishMessages for topic partition: {} failed when write.", fullPartitionName, e); requestStats.getMessagePublishStats().registerFailedEvent( - MathUtils.elapsedNanos(beforePublish), TimeUnit.NANOSECONDS); + time.nanoseconds() - beforePublish, TimeUnit.NANOSECONDS); errorsConsumer.accept(Errors.KAFKA_STORAGE_ERROR); } }); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java index 7e50c592bc..8712b7323b 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java @@ -15,10 +15,12 @@ import com.google.common.collect.Maps; import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; +import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; import io.streamnative.pulsar.handlers.kop.format.EntryFormatterFactory; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import java.util.Map; +import java.util.Optional; import lombok.AllArgsConstructor; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.utils.Time; @@ -28,12 +30,15 @@ public class PartitionLogManager { private final Map logMap; private final KafkaServiceConfiguration config; + private final Optional transactionCoordinator; private final EntryFormatter formatter; private final Time time; public PartitionLogManager(KafkaServiceConfiguration config, + Optional transactionCoordinator, Time time) { this.logMap = Maps.newConcurrentMap(); + this.transactionCoordinator = transactionCoordinator; this.formatter = EntryFormatterFactory.create(config); this.config = config; this.time = time; @@ -42,7 +47,8 @@ public PartitionLogManager(KafkaServiceConfiguration config, public PartitionLog getLog(TopicPartition topicPartition, String namespacePrefix) { String kopTopic = KopTopic.toString(topicPartition, namespacePrefix); return logMap.computeIfAbsent(kopTopic, key -> - new PartitionLog(config, topicPartition, namespacePrefix, kopTopic, formatter) + new PartitionLog(config, time, topicPartition, namespacePrefix, kopTopic, formatter, + this.transactionCoordinator) ); } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java index a0dc795333..e01364e15e 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -25,6 +25,7 @@ import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; @@ -47,9 +48,10 @@ public class ReplicaManager { public ReplicaManager(KafkaServiceConfiguration config, Time time, + Optional transactionCoordinator, DelayedOperationPurgatory producePurgatory, DelayedOperationPurgatory fetchPurgatory) { - this.logManager = new PartitionLogManager(config, time); + this.logManager = new PartitionLogManager(config, transactionCoordinator, time); this.producePurgatory = producePurgatory; this.fetchPurgatory = fetchPurgatory; } @@ -66,7 +68,6 @@ public void appendRecords( final String namespacePrefix, final Map entriesPerPartition, final RequestStats requestStats, - final TransactionCoordinator coordinator, final Consumer startSendOperationForThrottlingConsumer, final Consumer completeSendOperationForThrottlingConsumer, final Map pendingTopicFuturesMap, @@ -126,7 +127,6 @@ public void appendRecords( version, topicManager, requestStats, - coordinator, offsetConsumer, errorsConsumer, exceptionConsumer, From 367f471e9a652a336f7de8fec51da2807f18b78a Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 24 Nov 2021 10:06:00 +0800 Subject: [PATCH 04/15] Fix codacy check --- .../handlers/kop/KafkaProtocolHandler.java | 3 +- .../handlers/kop/storage/PartitionLog.java | 41 ++++++++++--------- .../handlers/kop/storage/ReplicaManager.java | 5 +-- 3 files changed, 24 insertions(+), 25 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index 73c9460d23..e807a9c46b 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -126,8 +126,7 @@ public ReplicaManager getReplicaManager(String tenant) { return new ReplicaManager(kafkaConfig, Time.SYSTEM, transactionCoordinatorOptional, - producePurgatory, - fetchPurgatory); + producePurgatory); } catch (Exception e) { log.error("Failed to init ReplicaManager for tenant {}", tenant, e); throw new IllegalStateException(e); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index b74ca22817..eda021ecb3 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -271,25 +271,7 @@ private LogAppendInfo analyzeAndValidateRecords(MemoryRecords records, boolean readFirstMessage = false; boolean monotonic = true; - if (version >= 3) { - Iterator iterator = records.batches().iterator(); - if (!iterator.hasNext()) { - throw new InvalidRecordException("Produce requests with version " + version + " must have at least " - + "one record batch"); - } - - MutableRecordBatch entry = iterator.next(); - if (entry.magic() != RecordBatch.MAGIC_VALUE_V2) { - throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " - + "contain record batches with magic version 2"); - } - - if (iterator.hasNext()) { - throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " - + "contain exactly one record batch"); - } - } - + validateRecords(version, records); int validBytesCount = 0; for (RecordBatch batch : records.batches()) { if (batch.magic() >= RecordBatch.MAGIC_VALUE_V2 && batch.baseOffset() != 0) { @@ -351,4 +333,25 @@ private MemoryRecords trimInvalidBytes(MemoryRecords records, LogAppendInfo info return MemoryRecords.readableRecords(validByteBuffer); } } + + private static void validateRecords(short version, MemoryRecords records) { + if (version >= 3) { + Iterator iterator = records.batches().iterator(); + if (!iterator.hasNext()) { + throw new InvalidRecordException("Produce requests with version " + version + " must have at least " + + "one record batch"); + } + + MutableRecordBatch entry = iterator.next(); + if (entry.magic() != RecordBatch.MAGIC_VALUE_V2) { + throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " + + "contain record batches with magic version 2"); + } + + if (iterator.hasNext()) { + throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " + + "contain exactly one record batch"); + } + } + } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java index e01364e15e..cce0cc4d1f 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -44,16 +44,13 @@ public class ReplicaManager { private final PartitionLogManager logManager; private final DelayedOperationPurgatory producePurgatory; - private final DelayedOperationPurgatory fetchPurgatory; public ReplicaManager(KafkaServiceConfiguration config, Time time, Optional transactionCoordinator, - DelayedOperationPurgatory producePurgatory, - DelayedOperationPurgatory fetchPurgatory) { + DelayedOperationPurgatory producePurgatory) { this.logManager = new PartitionLogManager(config, transactionCoordinator, time); this.producePurgatory = producePurgatory; - this.fetchPurgatory = fetchPurgatory; } public PartitionLog getPartitionLog(TopicPartition topicPartition, String namespacePrefix) { From ea6fcda5c6efb0450fa1d84e4a6c35947d8326aa Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 24 Nov 2021 17:46:20 +0800 Subject: [PATCH 05/15] Fix codacy check --- .../handlers/kop/PendingTopicFutures.java | 11 +++-- .../handlers/kop/storage/PartitionLog.java | 49 +++++++------------ .../handlers/kop/storage/ReplicaManager.java | 37 ++++++-------- .../handlers/kop/PendingTopicFuturesTest.java | 10 +++- 4 files changed, 47 insertions(+), 60 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/PendingTopicFutures.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/PendingTopicFutures.java index c514fa0d24..66e0a870e9 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/PendingTopicFutures.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/PendingTopicFutures.java @@ -55,7 +55,7 @@ private void registerQueueLatency(boolean success) { public void addListener(CompletableFuture> topicFuture, @NonNull Consumer> persistentTopicConsumer, - @NonNull Consumer exceptionConsumer) { + @NonNull CompletableFuture completableFuture) { if (count.compareAndSet(0, 1)) { // The first pending future comes currentTopicFuture = topicFuture.thenApply(persistentTopic -> { @@ -65,7 +65,8 @@ public void addListener(CompletableFuture> topicFuture return TopicThrowablePair.withTopic(persistentTopic); }).exceptionally(e -> { registerQueueLatency(false); - exceptionConsumer.accept(e.getCause()); +// exceptionConsumer.accept(e.getCause()); + completableFuture.completeExceptionally(e.getCause()); count.decrementAndGet(); return TopicThrowablePair.withThrowable(e.getCause()); }); @@ -77,13 +78,15 @@ public void addListener(CompletableFuture> topicFuture persistentTopicConsumer.accept(topicThrowablePair.getPersistentTopicOpt()); } else { registerQueueLatency(false); - exceptionConsumer.accept(topicThrowablePair.getThrowable()); + completableFuture.completeExceptionally(topicThrowablePair.getThrowable()); +// exceptionConsumer.accept(topicThrowablePair.getThrowable()); } count.decrementAndGet(); return topicThrowablePair; }).exceptionally(e -> { registerQueueLatency(false); - exceptionConsumer.accept(e.getCause()); +// exceptionConsumer.accept(e.getCause()); + completableFuture.completeExceptionally(e.getCause()); count.decrementAndGet(); return TopicThrowablePair.withThrowable(e.getCause()); }); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index eda021ecb3..c393473c26 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -83,24 +83,18 @@ public Long numMessages() { } } - public void appendRecords(final MemoryRecords records, + public CompletableFuture appendRecords(final MemoryRecords records, final short version, final KafkaTopicManager topicManager, final RequestStats requestStats, - final Consumer offsetConsumer, - final Consumer errorsConsumer, - final Consumer exceptionConsumer, final Consumer startSendOperationForThrottlingConsumer, final Consumer completeSendOperationForThrottlingConsumer, final Map pendingTopicFuturesMap) { - append(records, + return append(records, version, topicManager, requestStats, false, - offsetConsumer, - errorsConsumer, - exceptionConsumer, startSendOperationForThrottlingConsumer, completeSendOperationForThrottlingConsumer, pendingTopicFuturesMap); @@ -116,26 +110,20 @@ public void appendRecords(final MemoryRecords records, * @param version Inter-broker message protocol version * @param ignoreRecordSize true to skip validation of record size. */ - private void append(final MemoryRecords records, + private CompletableFuture append(final MemoryRecords records, final short version, final KafkaTopicManager topicManager, final RequestStats requestStats, final boolean ignoreRecordSize, - final Consumer offsetConsumer, - final Consumer errorsConsumer, - final Consumer exceptionConsumer, final Consumer startSendOperationForThrottlingConsumer, final Consumer completeSendOperationForThrottlingConsumer, final Map pendingTopicFuturesMap) { + CompletableFuture appendFuture = new CompletableFuture<>(); final long beforeRecordsProcess = time.nanoseconds(); try { final LogAppendInfo appendInfo = analyzeAndValidateRecords(records, version, topicPartition, ignoreRecordSize); - // return if we have no valid messages or if this is a duplicate of the last appended entry - if (appendInfo.getShallowCount() == 0) { - return; - } // trim any invalid bytes or partial messages before appending it to the on-disk log MemoryRecords validRecords = trimInvalidBytes(records, appendInfo); synchronized (lock) { @@ -144,18 +132,18 @@ private void append(final MemoryRecords records, topicManager.getTopic(fullPartitionName); if (topicFuture.isCompletedExceptionally()) { topicFuture.exceptionally(e -> { - exceptionConsumer.accept(e); + appendFuture.completeExceptionally(e); return Optional.empty(); }); - return; + return appendFuture; } if (topicFuture.isDone() && !topicFuture.getNow(Optional.empty()).isPresent()) { - errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); - return; + appendFuture.completeExceptionally(Errors.NOT_LEADER_FOR_PARTITION.exception()); + return appendFuture; } final Consumer> persistentTopicConsumer = persistentTopicOpt -> { if (!persistentTopicOpt.isPresent()) { - errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); + appendFuture.completeExceptionally(Errors.NOT_LEADER_FOR_PARTITION.exception()); return; } // TODO: validateMessagesAndAssignOffsets here. @@ -180,10 +168,9 @@ private void append(final MemoryRecords records, publishMessages(persistentTopicOpt, topicManager, requestStats, + appendFuture, encodeResult, topicPartition, - offsetConsumer, - errorsConsumer, completeSendOperationForThrottlingConsumer); }; @@ -194,23 +181,23 @@ private void append(final MemoryRecords records, pendingTopicFuturesMap .computeIfAbsent(topicPartition, ignored -> new PendingTopicFutures(requestStats)) - .addListener(topicFuture, persistentTopicConsumer, exceptionConsumer); + .addListener(topicFuture, persistentTopicConsumer, appendFuture); } } } catch (Exception exception) { log.error("Failed to handle produce request for {}", topicPartition, exception); - exceptionConsumer.accept(exception); + appendFuture.completeExceptionally(exception); } + return appendFuture; } private void publishMessages(final Optional persistentTopicOpt, final KafkaTopicManager topicManager, final RequestStats requestStats, + final CompletableFuture appendFuture, final EncodeResult encodeResult, final TopicPartition topicPartition, - final Consumer offsetConsumer, - final Consumer errorsConsumer, final Consumer completeSendOperationForThrottlingConsumer) { final MemoryRecords records = encodeResult.getRecords(); final int numMessages = encodeResult.getNumMessages(); @@ -218,14 +205,14 @@ private void publishMessages(final Optional persistentTopicOpt, if (!persistentTopicOpt.isPresent()) { encodeResult.recycle(); // It will trigger a retry send of Kafka client - errorsConsumer.accept(Errors.NOT_LEADER_FOR_PARTITION); + appendFuture.completeExceptionally(Errors.NOT_LEADER_FOR_PARTITION.exception()); return; } PersistentTopic persistentTopic = persistentTopicOpt.get(); if (persistentTopic.isSystemTopic()) { encodeResult.recycle(); log.error("Not support producing message to system topic: {}", persistentTopic); - errorsConsumer.accept(Errors.INVALID_TOPIC_EXCEPTION); + appendFuture.completeExceptionally(Errors.INVALID_TOPIC_EXCEPTION.exception()); return; } @@ -250,12 +237,12 @@ private void publishMessages(final Optional persistentTopicOpt, } requestStats.getMessagePublishStats().registerSuccessfulEvent( time.nanoseconds() - beforePublish, TimeUnit.NANOSECONDS); - offsetConsumer.accept(offset); + appendFuture.complete(offset); } else { log.error("publishMessages for topic partition: {} failed when write.", fullPartitionName, e); requestStats.getMessagePublishStats().registerFailedEvent( time.nanoseconds() - beforePublish, TimeUnit.NANOSECONDS); - errorsConsumer.accept(Errors.KAFKA_STORAGE_ERROR); + appendFuture.completeExceptionally(Errors.KAFKA_STORAGE_ERROR.exception()); } }); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java index cce0cc4d1f..383651bf4c 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -103,35 +103,26 @@ public void appendRecords( } }; entriesPerPartition.forEach((topicPartition, memoryRecords) -> { - final Consumer offsetConsumer = offset -> addPartitionResponse.accept( - topicPartition, - new ProduceResponse.PartitionResponse(Errors.NONE, offset, -1L, -1L)); - final Consumer errorsConsumer = - errors -> addPartitionResponse - .accept(topicPartition, new ProduceResponse.PartitionResponse(errors)); - final Consumer exceptionConsumer = - e -> addPartitionResponse - .accept(topicPartition, new ProduceResponse.PartitionResponse(Errors.forException(e))); - String fullPartitionName = KopTopic.toString(topicPartition, namespacePrefix); // reject appending to internal topics if it is not allowed if (!internalTopicsAllowed && KopTopic.isInternalTopic(fullPartitionName)) { - exceptionConsumer.accept(new InvalidTopicException( - String.format("Cannot append to internal topic %s", topicPartition.topic()))); + addPartitionResponse.accept(topicPartition, new ProduceResponse.PartitionResponse( + Errors.forException(new InvalidTopicException( + String.format("Cannot append to internal topic %s", topicPartition.topic()))))); } else { PartitionLog partitionLog = getPartitionLog(topicPartition, namespacePrefix); - partitionLog.appendRecords(memoryRecords, - version, - topicManager, - requestStats, - offsetConsumer, - errorsConsumer, - exceptionConsumer, - startSendOperationForThrottlingConsumer, - completeSendOperationForThrottlingConsumer, - pendingTopicFuturesMap); + partitionLog.appendRecords(memoryRecords, version, topicManager, requestStats, + startSendOperationForThrottlingConsumer, + completeSendOperationForThrottlingConsumer, + pendingTopicFuturesMap) + .thenAccept(offset -> addPartitionResponse.accept(topicPartition, + new ProduceResponse.PartitionResponse(Errors.NONE, offset, -1L, -1L))) + .exceptionally(ex -> { + addPartitionResponse.accept(topicPartition, + new ProduceResponse.PartitionResponse(Errors.forException(ex))); + return null; + }); } - }); // delay produce if (timeout <= 0) { diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/PendingTopicFuturesTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/PendingTopicFuturesTest.java index 7131599cff..053c7bd42a 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/PendingTopicFuturesTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/PendingTopicFuturesTest.java @@ -64,7 +64,8 @@ void testNormalComplete() throws ExecutionException, InterruptedException { for (int i = 0; i < 10; i++) { final int index = i; - pendingTopicFutures.addListener(topicFuture, ignored -> completedIndexes.add(index), e -> {}); + pendingTopicFutures.addListener( + topicFuture, ignored -> completedIndexes.add(index), new CompletableFuture<>()); changesOfPendingCount.add(pendingTopicFutures.size()); sleep(234); } @@ -95,7 +96,12 @@ void testExceptionalComplete() throws ExecutionException, InterruptedException { final List changesOfPendingCount = new ArrayList<>(); for (int i = 0; i < 10; i++) { - pendingTopicFutures.addListener(topicFuture, topic -> {}, e -> exceptionMessages.add(e.getMessage())); + CompletableFuture longCompletableFuture = new CompletableFuture<>(); + longCompletableFuture.exceptionally(ex -> { + exceptionMessages.add(ex.getMessage()); + return null; + }); + pendingTopicFutures.addListener(topicFuture, topic -> {}, longCompletableFuture); changesOfPendingCount.add(pendingTopicFutures.size()); sleep(200); } From 82033f771cd2c1498a43a01247ee7a233d2e2d74 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 24 Nov 2021 18:22:29 +0800 Subject: [PATCH 06/15] Fix wrong exception --- .../handlers/kop/KafkaRequestHandler.java | 21 ++++++++----------- .../handlers/kop/storage/ReplicaManager.java | 16 +++++++------- 2 files changed, 17 insertions(+), 20 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index 271c3960b9..1d3700ef2b 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -928,18 +928,15 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, requestStats, this::startSendOperationForThrottling, this::completeSendOperationForThrottling, - pendingTopicFuturesMap, - responseCallback - ); - responseCallback.thenAccept(response -> { - Map mergedResponse = Maps.newHashMap(); - mergedResponse.putAll(response); - mergedResponse.putAll(unauthorizedTopicResponsesMap); - resultFuture.complete(new ProduceResponse(mergedResponse)); - }).exceptionally(ex -> { - resultFuture.completeExceptionally(ex); - return null; - }); + pendingTopicFuturesMap).thenAccept(response -> { + Map mergedResponse = Maps.newHashMap(); + mergedResponse.putAll(response); + mergedResponse.putAll(unauthorizedTopicResponsesMap); + resultFuture.complete(new ProduceResponse(mergedResponse)); + }).exceptionally(ex -> { + resultFuture.completeExceptionally(ex.getCause()); + return null; + }); } }; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java index 383651bf4c..e3d4916935 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -57,7 +57,7 @@ public PartitionLog getPartitionLog(TopicPartition topicPartition, String namesp return logManager.getLog(topicPartition, namespacePrefix); } - public void appendRecords( + public CompletableFuture> appendRecords( final long timeout, final boolean internalTopicsAllowed, final short version, @@ -67,15 +67,15 @@ public void appendRecords( final RequestStats requestStats, final Consumer startSendOperationForThrottlingConsumer, final Consumer completeSendOperationForThrottlingConsumer, - final Map pendingTopicFuturesMap, - final CompletableFuture> responseCallback) { - + final Map pendingTopicFuturesMap) { + CompletableFuture> completableFuture = + new CompletableFuture<>(); final AtomicInteger topicPartitionNum = new AtomicInteger(entriesPerPartition.size()); final Map responseMap = new ConcurrentHashMap<>(); Runnable complete = () -> { topicPartitionNum.set(0); - if (responseCallback.isDone()) { + if (completableFuture.isDone()) { // It may be triggered again in DelayedProduceAndFetch return; } @@ -88,7 +88,7 @@ public void appendRecords( if (log.isDebugEnabled()) { log.debug("Complete handle appendRecords."); } - responseCallback.complete(responseMap); + completableFuture.complete(responseMap); }; BiConsumer addPartitionResponse = (topicPartition, response) -> { @@ -119,7 +119,7 @@ public void appendRecords( new ProduceResponse.PartitionResponse(Errors.NONE, offset, -1L, -1L))) .exceptionally(ex -> { addPartitionResponse.accept(topicPartition, - new ProduceResponse.PartitionResponse(Errors.forException(ex))); + new ProduceResponse.PartitionResponse(Errors.forException(ex.getCause()))); return null; }); } @@ -134,7 +134,7 @@ public void appendRecords( DelayedProduceAndFetch delayedProduce = new DelayedProduceAndFetch(timeout, topicPartitionNum, complete); producePurgatory.tryCompleteElseWatch(delayedProduce, delayedCreateKeys); } - + return completableFuture; } } From c4f09743e2abababc290a7169b629d07475143ea Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 24 Nov 2021 19:04:41 +0800 Subject: [PATCH 07/15] Add AppendRecordsContext to fix codacy check --- .../handlers/kop/KafkaRequestHandler.java | 34 +++++---- .../kop/storage/AppendRecordsContext.java | 70 +++++++++++++++++++ .../handlers/kop/storage/PartitionLog.java | 48 +++++-------- .../handlers/kop/storage/ReplicaManager.java | 15 +--- 4 files changed, 107 insertions(+), 60 deletions(-) create mode 100644 kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index 1d3700ef2b..db66e0a332 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -46,6 +46,7 @@ import io.streamnative.pulsar.handlers.kop.security.auth.ResourceType; import io.streamnative.pulsar.handlers.kop.security.auth.SimpleAclAuthorizer; import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; +import io.streamnative.pulsar.handlers.kop.storage.AppendRecordsContext; import io.streamnative.pulsar.handlers.kop.storage.ReplicaManager; import io.streamnative.pulsar.handlers.kop.utils.CoreUtils; import io.streamnative.pulsar.handlers.kop.utils.GroupIdUtils; @@ -916,27 +917,30 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, resultFuture.complete(new ProduceResponse(unauthorizedTopicResponsesMap)); return; } - CompletableFuture> responseCallback = - new CompletableFuture<>(); + AppendRecordsContext appendRecordsContext = AppendRecordsContext.get( + topicManager, + requestStats, + this::startSendOperationForThrottling, + this::completeSendOperationForThrottling, + pendingTopicFuturesMap); getReplicaManager().appendRecords( timeoutMs, false, produceHar.getRequest().version(), - topicManager, namespacePrefix, authorizedRequestInfo, - requestStats, - this::startSendOperationForThrottling, - this::completeSendOperationForThrottling, - pendingTopicFuturesMap).thenAccept(response -> { - Map mergedResponse = Maps.newHashMap(); - mergedResponse.putAll(response); - mergedResponse.putAll(unauthorizedTopicResponsesMap); - resultFuture.complete(new ProduceResponse(mergedResponse)); - }).exceptionally(ex -> { - resultFuture.completeExceptionally(ex.getCause()); - return null; - }); + appendRecordsContext + ).thenAccept(response -> { + appendRecordsContext.recycle(); + Map mergedResponse = new HashMap<>(); + mergedResponse.putAll(response); + mergedResponse.putAll(unauthorizedTopicResponsesMap); + resultFuture.complete(new ProduceResponse(mergedResponse)); + }).exceptionally(ex -> { + appendRecordsContext.recycle(); + resultFuture.completeExceptionally(ex.getCause()); + return null; + }); } }; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java new file mode 100644 index 0000000000..2a3e7c7c74 --- /dev/null +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java @@ -0,0 +1,70 @@ +/** + * 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.pulsar.handlers.kop.storage; + +import io.netty.util.Recycler; +import io.streamnative.pulsar.handlers.kop.KafkaTopicManager; +import io.streamnative.pulsar.handlers.kop.PendingTopicFutures; +import io.streamnative.pulsar.handlers.kop.RequestStats; +import lombok.Getter; +import org.apache.kafka.common.TopicPartition; + +import java.util.Map; +import java.util.function.Consumer; + +@Getter +public class AppendRecordsContext { + private static final Recycler RECYCLER = new Recycler() { + protected AppendRecordsContext newObject(Handle handle) { + return new AppendRecordsContext(handle); + } + }; + + private final Recycler.Handle recyclerHandle; + private KafkaTopicManager topicManager; + private RequestStats requestStats; + private Consumer startSendOperationForThrottling; + private Consumer completeSendOperationForThrottling; + private Map pendingTopicFuturesMap; + + private AppendRecordsContext(Recycler.Handle recyclerHandle) { + this.recyclerHandle = recyclerHandle; + } + + // recycler and get for this object + public static AppendRecordsContext get(final KafkaTopicManager topicManager, + final RequestStats requestStats, + final Consumer startSendOperationForThrottling, + final Consumer completeSendOperationForThrottling, + final Map pendingTopicFuturesMap) { + AppendRecordsContext context = RECYCLER.get(); + context.topicManager = topicManager; + context.requestStats = requestStats; + context.startSendOperationForThrottling = startSendOperationForThrottling; + context.completeSendOperationForThrottling = completeSendOperationForThrottling; + context.pendingTopicFuturesMap = pendingTopicFuturesMap; + + return context; + } + + public void recycle() { + topicManager = null; + requestStats = null; + startSendOperationForThrottling = null; + completeSendOperationForThrottling = null; + pendingTopicFuturesMap = null; + recyclerHandle.recycle(this); + } + +} diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index c393473c26..25de5e95b4 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -27,7 +27,6 @@ import io.streamnative.pulsar.handlers.kop.utils.MessageMetadataUtils; import java.nio.ByteBuffer; import java.util.Iterator; -import java.util.Map; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; @@ -84,20 +83,9 @@ public Long numMessages() { } public CompletableFuture appendRecords(final MemoryRecords records, - final short version, - final KafkaTopicManager topicManager, - final RequestStats requestStats, - final Consumer startSendOperationForThrottlingConsumer, - final Consumer completeSendOperationForThrottlingConsumer, - final Map pendingTopicFuturesMap) { - return append(records, - version, - topicManager, - requestStats, - false, - startSendOperationForThrottlingConsumer, - completeSendOperationForThrottlingConsumer, - pendingTopicFuturesMap); + final short version, + final AppendRecordsContext appendRecordsContext) { + return append(records, version, false, appendRecordsContext); } /** @@ -111,14 +99,12 @@ public CompletableFuture appendRecords(final MemoryRecords records, * @param ignoreRecordSize true to skip validation of record size. */ private CompletableFuture append(final MemoryRecords records, - final short version, - final KafkaTopicManager topicManager, - final RequestStats requestStats, - final boolean ignoreRecordSize, - final Consumer startSendOperationForThrottlingConsumer, - final Consumer completeSendOperationForThrottlingConsumer, - final Map pendingTopicFuturesMap) { + final short version, + final boolean ignoreRecordSize, + final AppendRecordsContext appendRecordsContext) { CompletableFuture appendFuture = new CompletableFuture<>(); + RequestStats requestStats = appendRecordsContext.getRequestStats(); + KafkaTopicManager topicManager = appendRecordsContext.getTopicManager(); final long beforeRecordsProcess = time.nanoseconds(); try { final LogAppendInfo appendInfo = @@ -159,26 +145,25 @@ private CompletableFuture append(final MemoryRecords records, encodeRequest.recycle(); requestStats.getProduceEncodeStats().registerSuccessfulEvent( time.nanoseconds() - beforeRecordsProcess, TimeUnit.NANOSECONDS); - startSendOperationForThrottlingConsumer.accept(encodeResult.getEncodedByteBuf().readableBytes()); + appendRecordsContext.getStartSendOperationForThrottling() + .accept(encodeResult.getEncodedByteBuf().readableBytes()); if (log.isDebugEnabled()) { log.debug("Produce messages for topic {} partition {}", topicPartition.topic(), topicPartition.partition()); } publishMessages(persistentTopicOpt, - topicManager, - requestStats, appendFuture, encodeResult, topicPartition, - completeSendOperationForThrottlingConsumer); + appendRecordsContext); }; if (topicFuture.isDone()) { persistentTopicConsumer.accept(topicFuture.getNow(Optional.empty())); } else { // topic is not available now - pendingTopicFuturesMap + appendRecordsContext.getPendingTopicFuturesMap() .computeIfAbsent(topicPartition, ignored -> new PendingTopicFutures(requestStats)) .addListener(topicFuture, persistentTopicConsumer, appendFuture); @@ -193,15 +178,14 @@ private CompletableFuture append(final MemoryRecords records, } private void publishMessages(final Optional persistentTopicOpt, - final KafkaTopicManager topicManager, - final RequestStats requestStats, final CompletableFuture appendFuture, final EncodeResult encodeResult, final TopicPartition topicPartition, - final Consumer completeSendOperationForThrottlingConsumer) { + final AppendRecordsContext appendRecordsContext) { final MemoryRecords records = encodeResult.getRecords(); final int numMessages = encodeResult.getNumMessages(); final ByteBuf byteBuf = encodeResult.getEncodedByteBuf(); + RequestStats requestStats = appendRecordsContext.getRequestStats(); if (!persistentTopicOpt.isPresent()) { encodeResult.recycle(); // It will trigger a retry send of Kafka client @@ -216,7 +200,7 @@ private void publishMessages(final Optional persistentTopicOpt, return; } - topicManager.registerProducerInPersistentTopic(fullPartitionName, persistentTopic); + appendRecordsContext.getTopicManager().registerProducerInPersistentTopic(fullPartitionName, persistentTopic); // collect metrics encodeResult.updateProducerStats(topicPartition, requestStats, namespacePrefix); @@ -228,7 +212,7 @@ private void publishMessages(final Optional persistentTopicOpt, MessagePublishContext.get(offsetFuture, persistentTopic, numMessages, System.nanoTime())); final RecordBatch batch = records.batchIterator().next(); offsetFuture.whenComplete((offset, e) -> { - completeSendOperationForThrottlingConsumer.accept(byteBuf.readableBytes()); + appendRecordsContext.getCompleteSendOperationForThrottling().accept(byteBuf.readableBytes()); encodeResult.recycle(); if (e == null) { if (batch.isTransactional()) { diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java index e3d4916935..14926c3743 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -15,9 +15,6 @@ import io.streamnative.pulsar.handlers.kop.DelayedProduceAndFetch; import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; -import io.streamnative.pulsar.handlers.kop.KafkaTopicManager; -import io.streamnative.pulsar.handlers.kop.PendingTopicFutures; -import io.streamnative.pulsar.handlers.kop.RequestStats; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; @@ -30,7 +27,6 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BiConsumer; -import java.util.function.Consumer; import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.common.TopicPartition; @@ -61,13 +57,9 @@ public CompletableFuture> final long timeout, final boolean internalTopicsAllowed, final short version, - final KafkaTopicManager topicManager, final String namespacePrefix, final Map entriesPerPartition, - final RequestStats requestStats, - final Consumer startSendOperationForThrottlingConsumer, - final Consumer completeSendOperationForThrottlingConsumer, - final Map pendingTopicFuturesMap) { + final AppendRecordsContext appendRecordsContext) { CompletableFuture> completableFuture = new CompletableFuture<>(); final AtomicInteger topicPartitionNum = new AtomicInteger(entriesPerPartition.size()); @@ -111,10 +103,7 @@ public CompletableFuture> String.format("Cannot append to internal topic %s", topicPartition.topic()))))); } else { PartitionLog partitionLog = getPartitionLog(topicPartition, namespacePrefix); - partitionLog.appendRecords(memoryRecords, version, topicManager, requestStats, - startSendOperationForThrottlingConsumer, - completeSendOperationForThrottlingConsumer, - pendingTopicFuturesMap) + partitionLog.appendRecords(memoryRecords, version, appendRecordsContext) .thenAccept(offset -> addPartitionResponse.accept(topicPartition, new ProduceResponse.PartitionResponse(Errors.NONE, offset, -1L, -1L))) .exceptionally(ex -> { From b273b07f1d5b1a7d0d60a48ec13ae41aa40220f6 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 24 Nov 2021 22:13:30 +0800 Subject: [PATCH 08/15] Fix code style --- .../pulsar/handlers/kop/storage/AppendRecordsContext.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java index 2a3e7c7c74..973ffbc260 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java @@ -17,11 +17,10 @@ import io.streamnative.pulsar.handlers.kop.KafkaTopicManager; import io.streamnative.pulsar.handlers.kop.PendingTopicFutures; import io.streamnative.pulsar.handlers.kop.RequestStats; -import lombok.Getter; -import org.apache.kafka.common.TopicPartition; - import java.util.Map; import java.util.function.Consumer; +import lombok.Getter; +import org.apache.kafka.common.TopicPartition; @Getter public class AppendRecordsContext { From 6445cc050b247a53a277f0a10062ce1fda3bb8f7 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Thu, 25 Nov 2021 09:54:58 +0800 Subject: [PATCH 09/15] Add more code comments --- .../handlers/kop/storage/AppendRecordsContext.java | 3 +++ .../pulsar/handlers/kop/storage/PartitionLog.java | 10 +++++++--- .../handlers/kop/storage/PartitionLogManager.java | 3 +++ .../pulsar/handlers/kop/storage/ReplicaManager.java | 3 +++ 4 files changed, 16 insertions(+), 3 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java index 973ffbc260..546f5d1b0f 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/AppendRecordsContext.java @@ -22,6 +22,9 @@ import lombok.Getter; import org.apache.kafka.common.TopicPartition; +/** + * AppendRecordsContext is use for pass parameters to ReplicaManager, to avoid long parameter lists. + */ @Getter public class AppendRecordsContext { private static final Recycler RECYCLER = new Recycler() { diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index 25de5e95b4..612bf7da0c 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -47,8 +47,11 @@ import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.common.naming.TopicName; -@AllArgsConstructor +/** + * An append-only log for storing messages. Mapping to Kafka Log.scala. + */ @Slf4j +@AllArgsConstructor public class PartitionLog { private KafkaServiceConfiguration kafkaConfig; private Time time; @@ -89,14 +92,15 @@ public CompletableFuture appendRecords(final MemoryRecords records, } /** - * Append this message set to the active segment of the log, rolling over to a fresh segment if necessary. + * Append this message to pulsar. * * This method will generally be responsible for assigning offsets to the messages, * however if the assignOffsets=false flag is passed we will only check that the existing offsets are valid. * * @param records The log records to append * @param version Inter-broker message protocol version - * @param ignoreRecordSize true to skip validation of record size. + * @param ignoreRecordSize true to skip validation of record size + * @param appendRecordsContext See {@link AppendRecordsContext} */ private CompletableFuture append(final MemoryRecords records, final short version, diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java index 8712b7323b..b25a5f1471 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java @@ -25,6 +25,9 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.utils.Time; +/** + * Manage {@link PartitionLog}. + */ @AllArgsConstructor public class PartitionLogManager { diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java index 14926c3743..ecbaab618c 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -36,6 +36,9 @@ import org.apache.kafka.common.requests.ProduceResponse; import org.apache.kafka.common.utils.Time; +/** + * Used to append records. Mapping to Kafka ReplicaManager.scala. + */ @Slf4j public class ReplicaManager { private final PartitionLogManager logManager; From 234dedca5a5561dfd99f5c59b4cc507fe0913096 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Thu, 25 Nov 2021 11:18:08 +0800 Subject: [PATCH 10/15] Remove unused code comments --- .../streamnative/pulsar/handlers/kop/PendingTopicFutures.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/PendingTopicFutures.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/PendingTopicFutures.java index 66e0a870e9..21372f8f92 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/PendingTopicFutures.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/PendingTopicFutures.java @@ -65,7 +65,6 @@ public void addListener(CompletableFuture> topicFuture return TopicThrowablePair.withTopic(persistentTopic); }).exceptionally(e -> { registerQueueLatency(false); -// exceptionConsumer.accept(e.getCause()); completableFuture.completeExceptionally(e.getCause()); count.decrementAndGet(); return TopicThrowablePair.withThrowable(e.getCause()); @@ -79,13 +78,11 @@ public void addListener(CompletableFuture> topicFuture } else { registerQueueLatency(false); completableFuture.completeExceptionally(topicThrowablePair.getThrowable()); -// exceptionConsumer.accept(topicThrowablePair.getThrowable()); } count.decrementAndGet(); return topicThrowablePair; }).exceptionally(e -> { registerQueueLatency(false); -// exceptionConsumer.accept(e.getCause()); completableFuture.completeExceptionally(e.getCause()); count.decrementAndGet(); return TopicThrowablePair.withThrowable(e.getCause()); From 678ea1f9979dac95d0b5ca6722cf198d8404efef Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Thu, 25 Nov 2021 20:45:33 +0800 Subject: [PATCH 11/15] Addressed review comments --- .../handlers/kop/KafkaProtocolHandler.java | 17 +++++++++++------ .../handlers/kop/KafkaRequestHandler.java | 10 +++++----- .../kop/storage/PartitionLogManager.java | 4 ++-- .../handlers/kop/storage/ReplicaManager.java | 4 +++- 4 files changed, 21 insertions(+), 14 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index e807a9c46b..0b3e305ad6 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -26,6 +26,8 @@ import io.streamnative.pulsar.handlers.kop.coordinator.group.OffsetConfig; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionConfig; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; +import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; +import io.streamnative.pulsar.handlers.kop.format.EntryFormatterFactory; import io.streamnative.pulsar.handlers.kop.stats.PrometheusMetricsProvider; import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; import io.streamnative.pulsar.handlers.kop.storage.ReplicaManager; @@ -122,15 +124,18 @@ public ReplicaManager getReplicaManager(String tenant) { if (kafkaConfig.isEnableTransactionCoordinator()) { transactionCoordinatorOptional = Optional.of(getTransactionCoordinator(tenant)); } + EntryFormatter entryFormatter; try { - return new ReplicaManager(kafkaConfig, - Time.SYSTEM, - transactionCoordinatorOptional, - producePurgatory); - } catch (Exception e) { - log.error("Failed to init ReplicaManager for tenant {}", tenant, e); + entryFormatter = EntryFormatterFactory.create(kafkaConfig); + } catch (IllegalArgumentException e) { + log.error("Failed to init create enter formatter {}", tenant, e); throw new IllegalStateException(e); } + return new ReplicaManager(kafkaConfig, + Time.SYSTEM, + entryFormatter, + transactionCoordinatorOptional, + producePurgatory); }); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index db66e0a332..9d7a97d624 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -930,16 +930,16 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, namespacePrefix, authorizedRequestInfo, appendRecordsContext - ).thenAccept(response -> { + ).whenComplete((response, ex) -> { appendRecordsContext.recycle(); + if (ex != null) { + resultFuture.completeExceptionally(ex.getCause()); + return; + } Map mergedResponse = new HashMap<>(); mergedResponse.putAll(response); mergedResponse.putAll(unauthorizedTopicResponsesMap); resultFuture.complete(new ProduceResponse(mergedResponse)); - }).exceptionally(ex -> { - appendRecordsContext.recycle(); - resultFuture.completeExceptionally(ex.getCause()); - return null; }); } }; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java index b25a5f1471..5c2fad5862 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java @@ -17,7 +17,6 @@ import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; -import io.streamnative.pulsar.handlers.kop.format.EntryFormatterFactory; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import java.util.Map; import java.util.Optional; @@ -38,11 +37,12 @@ public class PartitionLogManager { private final Time time; public PartitionLogManager(KafkaServiceConfiguration config, + EntryFormatter entryFormatter, Optional transactionCoordinator, Time time) { this.logMap = Maps.newConcurrentMap(); this.transactionCoordinator = transactionCoordinator; - this.formatter = EntryFormatterFactory.create(config); + this.formatter = entryFormatter; this.config = config; this.time = time; } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java index ecbaab618c..52549db785 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -16,6 +16,7 @@ import io.streamnative.pulsar.handlers.kop.DelayedProduceAndFetch; import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; +import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationKey; @@ -46,9 +47,10 @@ public class ReplicaManager { public ReplicaManager(KafkaServiceConfiguration config, Time time, + EntryFormatter entryFormatter, Optional transactionCoordinator, DelayedOperationPurgatory producePurgatory) { - this.logManager = new PartitionLogManager(config, transactionCoordinator, time); + this.logManager = new PartitionLogManager(config, entryFormatter, transactionCoordinator, time); this.producePurgatory = producePurgatory; } From 2d26708432c96e4a60b27f67a561d0fdb20b5a1a Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Fri, 26 Nov 2021 16:17:35 +0800 Subject: [PATCH 12/15] Addressed review comments --- .../handlers/kop/KafkaProtocolHandler.java | 2 +- .../handlers/kop/KafkaRequestHandler.java | 93 ++------ .../handlers/kop/storage/PartitionLog.java | 216 +++++++----------- .../kop/storage/PartitionLogManager.java | 8 +- .../handlers/kop/storage/ReplicaManager.java | 6 +- 5 files changed, 98 insertions(+), 227 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index 0b3e305ad6..297f6d5715 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -131,7 +131,7 @@ public ReplicaManager getReplicaManager(String tenant) { log.error("Failed to init create enter formatter {}", tenant, e); throw new IllegalStateException(e); } - return new ReplicaManager(kafkaConfig, + return new ReplicaManager( Time.SYSTEM, entryFormatter, transactionCoordinatorOptional, diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index 9d7a97d624..c0e2a25fa4 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -19,7 +19,6 @@ import static io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration.TENANT_PLACEHOLDER; import static java.nio.charset.StandardCharsets.UTF_8; import static org.apache.kafka.common.internals.Topic.GROUP_METADATA_TOPIC_NAME; -import static org.apache.kafka.common.internals.Topic.TRANSACTION_STATE_TOPIC_NAME; import static org.apache.kafka.common.protocol.CommonFields.THROTTLE_TIME_MS; import static org.apache.kafka.common.requests.CreateTopicsRequest.TopicDetails; @@ -65,7 +64,6 @@ import java.util.Collections; import java.util.HashMap; import java.util.HashSet; -import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Optional; @@ -97,15 +95,12 @@ import org.apache.kafka.common.acl.AclOperation; import org.apache.kafka.common.config.ConfigResource; import org.apache.kafka.common.errors.AuthenticationException; -import org.apache.kafka.common.errors.CorruptRecordException; import org.apache.kafka.common.errors.LeaderNotAvailableException; import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.record.ControlRecordType; import org.apache.kafka.common.record.EndTransactionMarker; -import org.apache.kafka.common.record.InvalidRecordException; import org.apache.kafka.common.record.MemoryRecords; -import org.apache.kafka.common.record.MutableRecordBatch; import org.apache.kafka.common.record.RecordBatch; import org.apache.kafka.common.requests.AbstractRequest; import org.apache.kafka.common.requests.AbstractResponse; @@ -471,12 +466,6 @@ private CompletableFuture getPartitionedTopicMetadataA return admin.topics().getPartitionedTopicMetadataAsync(topicName); } - private static boolean isInternalTopic(final String fullTopicName) { - String partitionedTopicName = TopicName.get(fullTopicName).getPartitionedTopicName(); - return partitionedTopicName.endsWith("/" + GROUP_METADATA_TOPIC_NAME) - || partitionedTopicName.endsWith("/" + TRANSACTION_STATE_TOPIC_NAME); - } - private CompletableFuture> expandAllowedNamespaces(Set allowedNamespaces) { String currentTenant = getCurrentTenant(kafkaConfig.getKafkaTenant()); return expandAllowedNamespaces(allowedNamespaces, currentTenant, pulsarService); @@ -598,7 +587,7 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, allTopicMetadata.add(new TopicMetadata( Errors.TOPIC_AUTHORIZATION_FAILED, topic, - isInternalTopic(topicName.toString()), + KopTopic.isInternalTopic(topicName.toString()), Collections.emptyList())); return; } @@ -643,7 +632,7 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, allTopicMetadata.add(new TopicMetadata( Errors.TOPIC_AUTHORIZATION_FAILED, topic, - isInternalTopic(fullTopicName), + KopTopic.isInternalTopic(fullTopicName), Collections.emptyList())); completeOneTopic.run(); }; @@ -698,7 +687,7 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, new TopicMetadata( Errors.UNKNOWN_TOPIC_OR_PARTITION, topic, - isInternalTopic(fullTopicName), + KopTopic.isInternalTopic(fullTopicName), Collections.emptyList())); completeOneTopic.run(); } @@ -708,7 +697,7 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, new TopicMetadata( Errors.UNKNOWN_TOPIC_OR_PARTITION, topic, - isInternalTopic(fullTopicName), + KopTopic.isInternalTopic(fullTopicName), Collections.emptyList())); log.warn("[{}] Request {}: Failed to get partitioned pulsar topic {} " + "metadata: {}", @@ -813,7 +802,8 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, // The topic returned to Kafka clients should be // the same with what it sent topic, - isInternalTopic(new KopTopic(topic, namespacePrefix).getFullName()), + KopTopic.isInternalTopic( + new KopTopic(topic, namespacePrefix).getFullName()), partitionMetadatas)); // whether completed all the topics requests. @@ -909,9 +899,8 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, int timeoutMs = produceRequest.timeout(); String namespacePrefix = currentNamespacePrefix(); final AtomicInteger unfinishedAuthorizationCount = new AtomicInteger(numPartitions); - Consumer completeOne = (action) -> { + Runnable completeOne = () -> { // When complete one authorization or failed, will do the action first. - action.run(); if (unfinishedAuthorizationCount.decrementAndGet() == 0) { if (authorizedRequestInfo.isEmpty()) { resultFuture.complete(new ProduceResponse(unauthorizedTopicResponsesMap)); @@ -951,22 +940,19 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, if (ex != null) { log.error("Write topic authorize failed, topic - {}. {}", fullPartitionName, ex.getMessage()); - completeOne.accept(() -> { - unauthorizedTopicResponsesMap.put(topicPartition, - new ProduceResponse.PartitionResponse(Errors.TOPIC_AUTHORIZATION_FAILED)); - }); + unauthorizedTopicResponsesMap.put(topicPartition, + new ProduceResponse.PartitionResponse(Errors.TOPIC_AUTHORIZATION_FAILED)); + completeOne.run(); return; } if (!isAuthorized) { - completeOne.accept(() -> { - unauthorizedTopicResponsesMap.put(topicPartition, - new ProduceResponse.PartitionResponse(Errors.TOPIC_AUTHORIZATION_FAILED)); - }); + unauthorizedTopicResponsesMap.put(topicPartition, + new ProduceResponse.PartitionResponse(Errors.TOPIC_AUTHORIZATION_FAILED)); + completeOne.run(); return; } - completeOne.accept(() -> { - authorizedRequestInfo.put(topicPartition, records); - }); + authorizedRequestInfo.put(topicPartition, records); + completeOne.run(); }); }); @@ -2533,55 +2519,6 @@ static AbstractResponse failedResponse(KafkaHeaderAndRequest requestHar, Throwab return requestHar.getRequest().getErrorResponse(((Integer) THROTTLE_TIME_MS.defaultValue), e); } - private static MemoryRecords validateRecords(short version, TopicPartition topicPartition, MemoryRecords records) { - if (version >= 3) { - Iterator iterator = records.batches().iterator(); - if (!iterator.hasNext()) { - throw new InvalidRecordException("Produce requests with version " + version + " must have at least " - + "one record batch"); - } - - MutableRecordBatch entry = iterator.next(); - if (entry.magic() != RecordBatch.MAGIC_VALUE_V2) { - throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " - + "contain record batches with magic version 2"); - } - - if (iterator.hasNext()) { - throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " - + "contain exactly one record batch"); - } - } - - int validBytesCount = 0; - for (RecordBatch batch : records.batches()) { - if (batch.magic() >= RecordBatch.MAGIC_VALUE_V2 && batch.baseOffset() != 0) { - throw new InvalidRecordException("The baseOffset of the record batch in the append to " - + topicPartition + " should be 0, but it is " + batch.baseOffset()); - } - - batch.ensureValid(); - validBytesCount += batch.sizeInBytes(); - } - - if (validBytesCount < 0) { - throw new CorruptRecordException("Cannot append record batch with illegal length " - + validBytesCount + " to log for " + topicPartition - + ". A possible cause is corrupted produce request."); - } - - MemoryRecords validRecords; - if (validBytesCount == records.sizeInBytes()) { - validRecords = records; - } else { - ByteBuffer validByteBuffer = records.buffer().duplicate(); - validByteBuffer.limit(validBytesCount); - validRecords = MemoryRecords.readableRecords(validByteBuffer); - } - - return validRecords; - } - @VisibleForTesting protected CompletableFuture authorize(AclOperation operation, Resource resource) { Session session = authenticator != null ? authenticator.session() : null; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index 612bf7da0c..9487a02c67 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -14,7 +14,6 @@ package io.streamnative.pulsar.handlers.kop.storage; import io.netty.buffer.ByteBuf; -import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.KafkaTopicManager; import io.streamnative.pulsar.handlers.kop.MessagePublishContext; import io.streamnative.pulsar.handlers.kop.PendingTopicFutures; @@ -37,7 +36,6 @@ import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.CorruptRecordException; -import org.apache.kafka.common.errors.RecordTooLargeException; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.record.InvalidRecordException; import org.apache.kafka.common.record.MemoryRecords; @@ -53,7 +51,6 @@ @Slf4j @AllArgsConstructor public class PartitionLog { - private KafkaServiceConfiguration kafkaConfig; private Time time; private TopicPartition topicPartition; private String namespacePrefix; @@ -88,90 +85,79 @@ public Long numMessages() { public CompletableFuture appendRecords(final MemoryRecords records, final short version, final AppendRecordsContext appendRecordsContext) { - return append(records, version, false, appendRecordsContext); + return append(records, version, appendRecordsContext); } /** * Append this message to pulsar. * - * This method will generally be responsible for assigning offsets to the messages, - * however if the assignOffsets=false flag is passed we will only check that the existing offsets are valid. - * * @param records The log records to append * @param version Inter-broker message protocol version - * @param ignoreRecordSize true to skip validation of record size * @param appendRecordsContext See {@link AppendRecordsContext} */ private CompletableFuture append(final MemoryRecords records, final short version, - final boolean ignoreRecordSize, final AppendRecordsContext appendRecordsContext) { CompletableFuture appendFuture = new CompletableFuture<>(); RequestStats requestStats = appendRecordsContext.getRequestStats(); KafkaTopicManager topicManager = appendRecordsContext.getTopicManager(); final long beforeRecordsProcess = time.nanoseconds(); try { - final LogAppendInfo appendInfo = - analyzeAndValidateRecords(records, version, topicPartition, ignoreRecordSize); - - // trim any invalid bytes or partial messages before appending it to the on-disk log - MemoryRecords validRecords = trimInvalidBytes(records, appendInfo); - synchronized (lock) { - // Append Message into pulsar - final CompletableFuture> topicFuture = - topicManager.getTopic(fullPartitionName); - if (topicFuture.isCompletedExceptionally()) { - topicFuture.exceptionally(e -> { - appendFuture.completeExceptionally(e); - return Optional.empty(); - }); - return appendFuture; - } - if (topicFuture.isDone() && !topicFuture.getNow(Optional.empty()).isPresent()) { + MemoryRecords validRecords = validateRecords(version, fullPartitionName, records); + // Append Message into pulsar + final CompletableFuture> topicFuture = + topicManager.getTopic(fullPartitionName); + if (topicFuture.isCompletedExceptionally()) { + topicFuture.exceptionally(e -> { + appendFuture.completeExceptionally(e); + return Optional.empty(); + }); + return appendFuture; + } + if (topicFuture.isDone() && !topicFuture.getNow(Optional.empty()).isPresent()) { + appendFuture.completeExceptionally(Errors.NOT_LEADER_FOR_PARTITION.exception()); + return appendFuture; + } + final Consumer> persistentTopicConsumer = persistentTopicOpt -> { + if (!persistentTopicOpt.isPresent()) { appendFuture.completeExceptionally(Errors.NOT_LEADER_FOR_PARTITION.exception()); - return appendFuture; + return; } - final Consumer> persistentTopicConsumer = persistentTopicOpt -> { - if (!persistentTopicOpt.isPresent()) { - appendFuture.completeExceptionally(Errors.NOT_LEADER_FOR_PARTITION.exception()); - return; - } - // TODO: validateMessagesAndAssignOffsets here. + // TODO: validateMessagesAndAssignOffsets here. - final EncodeRequest encodeRequest = EncodeRequest.get(validRecords); - if (entryFormatter instanceof KafkaMixedEntryFormatter) { - final ManagedLedger managedLedger = persistentTopicOpt.get().getManagedLedger(); - final long logEndOffset = MessageMetadataUtils.getLogEndOffset(managedLedger); - encodeRequest.setBaseOffset(logEndOffset); - } + final EncodeRequest encodeRequest = EncodeRequest.get(validRecords); + if (entryFormatter instanceof KafkaMixedEntryFormatter) { + final ManagedLedger managedLedger = persistentTopicOpt.get().getManagedLedger(); + final long logEndOffset = MessageMetadataUtils.getLogEndOffset(managedLedger); + encodeRequest.setBaseOffset(logEndOffset); + } - final EncodeResult encodeResult = entryFormatter.encode(encodeRequest); - encodeRequest.recycle(); - requestStats.getProduceEncodeStats().registerSuccessfulEvent( - time.nanoseconds() - beforeRecordsProcess, TimeUnit.NANOSECONDS); - appendRecordsContext.getStartSendOperationForThrottling() - .accept(encodeResult.getEncodedByteBuf().readableBytes()); - if (log.isDebugEnabled()) { - log.debug("Produce messages for topic {} partition {}", - topicPartition.topic(), topicPartition.partition()); - } + final EncodeResult encodeResult = entryFormatter.encode(encodeRequest); + encodeRequest.recycle(); + requestStats.getProduceEncodeStats().registerSuccessfulEvent( + time.nanoseconds() - beforeRecordsProcess, TimeUnit.NANOSECONDS); + appendRecordsContext.getStartSendOperationForThrottling() + .accept(encodeResult.getEncodedByteBuf().readableBytes()); + if (log.isDebugEnabled()) { + log.debug("Produce messages for topic {} partition {}", + topicPartition.topic(), topicPartition.partition()); + } - publishMessages(persistentTopicOpt, - appendFuture, - encodeResult, - topicPartition, - appendRecordsContext); - }; + publishMessages(persistentTopicOpt, + appendFuture, + encodeResult, + topicPartition, + appendRecordsContext); + }; - if (topicFuture.isDone()) { - persistentTopicConsumer.accept(topicFuture.getNow(Optional.empty())); - } else { - // topic is not available now - appendRecordsContext.getPendingTopicFuturesMap() - .computeIfAbsent(topicPartition, ignored -> - new PendingTopicFutures(requestStats)) - .addListener(topicFuture, persistentTopicConsumer, appendFuture); - } + if (topicFuture.isDone()) { + persistentTopicConsumer.accept(topicFuture.getNow(Optional.empty())); + } else { + // topic is not available now + appendRecordsContext.getPendingTopicFuturesMap() + .computeIfAbsent(topicPartition, ignored -> + new PendingTopicFutures(requestStats)) + .addListener(topicFuture, persistentTopicConsumer, appendFuture); } } catch (Exception exception) { log.error("Failed to handle produce request for {}", topicPartition, exception); @@ -235,98 +221,52 @@ private void publishMessages(final Optional persistentTopicOpt, }); } - private LogAppendInfo analyzeAndValidateRecords(MemoryRecords records, - short version, - TopicPartition topicPartition, - boolean ignoreRecordSize) { - int shallowMessageCount = 0; - long lastOffset = -1L; - Optional firstOffset = Optional.empty(); - long lastOffsetOfFirstBatch = -1L; - boolean readFirstMessage = false; - boolean monotonic = true; + private static MemoryRecords validateRecords(short version, String fullPartitionName, MemoryRecords records) { + if (version >= 3) { + Iterator iterator = records.batches().iterator(); + if (!iterator.hasNext()) { + throw new InvalidRecordException("Produce requests with version " + version + " must have at least " + + "one record batch"); + } + + MutableRecordBatch entry = iterator.next(); + if (entry.magic() != RecordBatch.MAGIC_VALUE_V2) { + throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " + + "contain record batches with magic version 2"); + } + + if (iterator.hasNext()) { + throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " + + "contain exactly one record batch"); + } + } - validateRecords(version, records); int validBytesCount = 0; for (RecordBatch batch : records.batches()) { if (batch.magic() >= RecordBatch.MAGIC_VALUE_V2 && batch.baseOffset() != 0) { throw new InvalidRecordException("The baseOffset of the record batch in the append to " - + topicPartition + " should be 0, but it is " + batch.baseOffset()); - } - if (!readFirstMessage) { - if (batch.magic() >= RecordBatch.MAGIC_VALUE_V2) { - firstOffset = Optional.of(batch.baseOffset()); - } - lastOffsetOfFirstBatch = batch.lastOffset(); - readFirstMessage = true; - } - // check that offsets are monotonically increasing - if (lastOffset >= batch.lastOffset()){ - monotonic = false; - } - - // update the last offset seen - lastOffset = batch.lastOffset(); - - int batchSize = batch.sizeInBytes(); - if (!ignoreRecordSize && batchSize > kafkaConfig.getMaxMessageSize()) { - throw new RecordTooLargeException(String.format("Message batch size is %s " - + "in append to partition %s which exceeds the maximum configured size of %s .", - batchSize, topicPartition, kafkaConfig.getMaxMessageSize())); + + fullPartitionName + " should be 0, but it is " + batch.baseOffset()); } batch.ensureValid(); - shallowMessageCount += 1; - validBytesCount += batchSize; + validBytesCount += batch.sizeInBytes(); } if (validBytesCount < 0) { throw new CorruptRecordException("Cannot append record batch with illegal length " - + validBytesCount + " to log for " + topicPartition + + validBytesCount + " to log for " + fullPartitionName + ". A possible cause is corrupted produce request."); } - return new LogAppendInfo( - firstOffset, - lastOffset, - shallowMessageCount, - monotonic, - lastOffsetOfFirstBatch, - validBytesCount); - } - - private MemoryRecords trimInvalidBytes(MemoryRecords records, LogAppendInfo info) { - Integer validBytes = info.getValidBytes(); - if (validBytes < 0){ - throw new CorruptRecordException(String.format("Cannot append record batch with illegal length %s to " - + "log for %s. A possible cause is a corrupted produce request.", validBytes, topicPartition)); - } else if (validBytes == records.sizeInBytes()) { - return records; + MemoryRecords validRecords; + if (validBytesCount == records.sizeInBytes()) { + validRecords = records; } else { ByteBuffer validByteBuffer = records.buffer().duplicate(); - validByteBuffer.limit(validBytes); - return MemoryRecords.readableRecords(validByteBuffer); + validByteBuffer.limit(validBytesCount); + validRecords = MemoryRecords.readableRecords(validByteBuffer); } - } - private static void validateRecords(short version, MemoryRecords records) { - if (version >= 3) { - Iterator iterator = records.batches().iterator(); - if (!iterator.hasNext()) { - throw new InvalidRecordException("Produce requests with version " + version + " must have at least " - + "one record batch"); - } - - MutableRecordBatch entry = iterator.next(); - if (entry.magic() != RecordBatch.MAGIC_VALUE_V2) { - throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " - + "contain record batches with magic version 2"); - } - - if (iterator.hasNext()) { - throw new InvalidRecordException("Produce requests with version " + version + " are only allowed to " - + "contain exactly one record batch"); - } - } + return validRecords; } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java index 5c2fad5862..3f1b260519 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java @@ -14,7 +14,6 @@ package io.streamnative.pulsar.handlers.kop.storage; import com.google.common.collect.Maps; -import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; @@ -31,26 +30,23 @@ public class PartitionLogManager { private final Map logMap; - private final KafkaServiceConfiguration config; private final Optional transactionCoordinator; private final EntryFormatter formatter; private final Time time; - public PartitionLogManager(KafkaServiceConfiguration config, - EntryFormatter entryFormatter, + public PartitionLogManager(EntryFormatter entryFormatter, Optional transactionCoordinator, Time time) { this.logMap = Maps.newConcurrentMap(); this.transactionCoordinator = transactionCoordinator; this.formatter = entryFormatter; - this.config = config; this.time = time; } public PartitionLog getLog(TopicPartition topicPartition, String namespacePrefix) { String kopTopic = KopTopic.toString(topicPartition, namespacePrefix); return logMap.computeIfAbsent(kopTopic, key -> - new PartitionLog(config, time, topicPartition, namespacePrefix, kopTopic, formatter, + new PartitionLog(time, topicPartition, namespacePrefix, kopTopic, formatter, this.transactionCoordinator) ); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java index 52549db785..4025fc52f1 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -14,7 +14,6 @@ package io.streamnative.pulsar.handlers.kop.storage; import io.streamnative.pulsar.handlers.kop.DelayedProduceAndFetch; -import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; @@ -45,12 +44,11 @@ public class ReplicaManager { private final PartitionLogManager logManager; private final DelayedOperationPurgatory producePurgatory; - public ReplicaManager(KafkaServiceConfiguration config, - Time time, + public ReplicaManager(Time time, EntryFormatter entryFormatter, Optional transactionCoordinator, DelayedOperationPurgatory producePurgatory) { - this.logManager = new PartitionLogManager(config, entryFormatter, transactionCoordinator, time); + this.logManager = new PartitionLogManager(entryFormatter, transactionCoordinator, time); this.producePurgatory = producePurgatory; } From 13891a687b56343757d17095587ffdd83f3574d9 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Fri, 26 Nov 2021 16:21:27 +0800 Subject: [PATCH 13/15] Remove unused code --- .../handlers/kop/storage/PartitionLog.java | 25 ------------------- 1 file changed, 25 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index 9487a02c67..76b5c412fc 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -31,7 +31,6 @@ import java.util.concurrent.TimeUnit; import java.util.function.Consumer; import lombok.AllArgsConstructor; -import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.kafka.common.TopicPartition; @@ -58,30 +57,6 @@ public class PartitionLog { private EntryFormatter entryFormatter; private Optional transactionCoordinator; - - // A lock that guards all modifications to the log - private final Object lock = new Object(); - - @Data - @AllArgsConstructor - public static final class LogAppendInfo { - private Optional firstOffset; - private Long lastOffset; - private Integer shallowCount; - private Boolean offsetsMonotonic; - private Long lastOffsetOfFirstBatch; - private Integer validBytes; - - public Long numMessages() { - return firstOffset.map(firstOffsetVal -> { - if (firstOffsetVal >= 0 && lastOffset >= 0) { - return lastOffset - firstOffsetVal + 1; - } - return 0L; - }).orElse(0L); - } - } - public CompletableFuture appendRecords(final MemoryRecords records, final short version, final AppendRecordsContext appendRecordsContext) { From 636d55631ba24ea8821fe33cb405fe1271934087 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Fri, 26 Nov 2021 16:57:32 +0800 Subject: [PATCH 14/15] Fix units test --- .../pulsar/handlers/kop/KafkaProtocolHandler.java | 1 + .../pulsar/handlers/kop/storage/PartitionLog.java | 10 ++++++++++ .../handlers/kop/storage/PartitionLogManager.java | 8 ++++++-- .../pulsar/handlers/kop/storage/ReplicaManager.java | 6 ++++-- 4 files changed, 21 insertions(+), 4 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index 297f6d5715..2382ad66df 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -132,6 +132,7 @@ public ReplicaManager getReplicaManager(String tenant) { throw new IllegalStateException(e); } return new ReplicaManager( + kafkaConfig, Time.SYSTEM, entryFormatter, transactionCoordinatorOptional, diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index 76b5c412fc..5671b73666 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -14,6 +14,7 @@ package io.streamnative.pulsar.handlers.kop.storage; import io.netty.buffer.ByteBuf; +import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.KafkaTopicManager; import io.streamnative.pulsar.handlers.kop.MessagePublishContext; import io.streamnative.pulsar.handlers.kop.PendingTopicFutures; @@ -35,6 +36,7 @@ import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.CorruptRecordException; +import org.apache.kafka.common.errors.RecordTooLargeException; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.record.InvalidRecordException; import org.apache.kafka.common.record.MemoryRecords; @@ -50,6 +52,7 @@ @Slf4j @AllArgsConstructor public class PartitionLog { + private KafkaServiceConfiguration kafkaConfig; private Time time; private TopicPartition topicPartition; private String namespacePrefix; @@ -79,6 +82,13 @@ private CompletableFuture append(final MemoryRecords records, final long beforeRecordsProcess = time.nanoseconds(); try { MemoryRecords validRecords = validateRecords(version, fullPartitionName, records); + validRecords.batches().forEach(batch->{ + if (batch.sizeInBytes() > kafkaConfig.getMaxMessageSize()) { + throw new RecordTooLargeException(String.format("Message batch size is %s " + + "in append to partition %s which exceeds the maximum configured size of %s .", + batch.sizeInBytes(), topicPartition, kafkaConfig.getMaxMessageSize())); + } + }); // Append Message into pulsar final CompletableFuture> topicFuture = topicManager.getTopic(fullPartitionName); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java index 3f1b260519..1c2b04464d 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLogManager.java @@ -14,6 +14,7 @@ package io.streamnative.pulsar.handlers.kop.storage; import com.google.common.collect.Maps; +import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; @@ -29,14 +30,17 @@ @AllArgsConstructor public class PartitionLogManager { + private KafkaServiceConfiguration kafkaConfig; private final Map logMap; private final Optional transactionCoordinator; private final EntryFormatter formatter; private final Time time; - public PartitionLogManager(EntryFormatter entryFormatter, + public PartitionLogManager(KafkaServiceConfiguration kafkaConfig, + EntryFormatter entryFormatter, Optional transactionCoordinator, Time time) { + this.kafkaConfig = kafkaConfig; this.logMap = Maps.newConcurrentMap(); this.transactionCoordinator = transactionCoordinator; this.formatter = entryFormatter; @@ -46,7 +50,7 @@ public PartitionLogManager(EntryFormatter entryFormatter, public PartitionLog getLog(TopicPartition topicPartition, String namespacePrefix) { String kopTopic = KopTopic.toString(topicPartition, namespacePrefix); return logMap.computeIfAbsent(kopTopic, key -> - new PartitionLog(time, topicPartition, namespacePrefix, kopTopic, formatter, + new PartitionLog(kafkaConfig, time, topicPartition, namespacePrefix, kopTopic, formatter, this.transactionCoordinator) ); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java index 4025fc52f1..3936dfcbd3 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/ReplicaManager.java @@ -14,6 +14,7 @@ package io.streamnative.pulsar.handlers.kop.storage; import io.streamnative.pulsar.handlers.kop.DelayedProduceAndFetch; +import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; @@ -44,11 +45,12 @@ public class ReplicaManager { private final PartitionLogManager logManager; private final DelayedOperationPurgatory producePurgatory; - public ReplicaManager(Time time, + public ReplicaManager(KafkaServiceConfiguration kafkaConfig, + Time time, EntryFormatter entryFormatter, Optional transactionCoordinator, DelayedOperationPurgatory producePurgatory) { - this.logManager = new PartitionLogManager(entryFormatter, transactionCoordinator, time); + this.logManager = new PartitionLogManager(kafkaConfig, entryFormatter, transactionCoordinator, time); this.producePurgatory = producePurgatory; } From b4aea77fd55c8eca73a9ab0299a5bbea1dc57989 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Fri, 26 Nov 2021 18:39:06 +0800 Subject: [PATCH 15/15] Addressed reviewer comments --- .../handlers/kop/storage/PartitionLog.java | 26 ++++++------------- 1 file changed, 8 insertions(+), 18 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index 5671b73666..ec3c695c50 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -52,19 +52,13 @@ @Slf4j @AllArgsConstructor public class PartitionLog { - private KafkaServiceConfiguration kafkaConfig; - private Time time; - private TopicPartition topicPartition; - private String namespacePrefix; - private String fullPartitionName; - private EntryFormatter entryFormatter; - private Optional transactionCoordinator; - - public CompletableFuture appendRecords(final MemoryRecords records, - final short version, - final AppendRecordsContext appendRecordsContext) { - return append(records, version, appendRecordsContext); - } + private final KafkaServiceConfiguration kafkaConfig; + private final Time time; + private final TopicPartition topicPartition; + private final String namespacePrefix; + private final String fullPartitionName; + private final EntryFormatter entryFormatter; + private final Optional transactionCoordinator; /** * Append this message to pulsar. @@ -73,7 +67,7 @@ public CompletableFuture appendRecords(final MemoryRecords records, * @param version Inter-broker message protocol version * @param appendRecordsContext See {@link AppendRecordsContext} */ - private CompletableFuture append(final MemoryRecords records, + public CompletableFuture appendRecords(final MemoryRecords records, final short version, final AppendRecordsContext appendRecordsContext) { CompletableFuture appendFuture = new CompletableFuture<>(); @@ -123,10 +117,6 @@ private CompletableFuture append(final MemoryRecords records, time.nanoseconds() - beforeRecordsProcess, TimeUnit.NANOSECONDS); appendRecordsContext.getStartSendOperationForThrottling() .accept(encodeResult.getEncodedByteBuf().readableBytes()); - if (log.isDebugEnabled()) { - log.debug("Produce messages for topic {} partition {}", - topicPartition.topic(), topicPartition.partition()); - } publishMessages(persistentTopicOpt, appendFuture,