From 5a39c69da4034ca135203fbdebece258eab5fe0f Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Sun, 24 Oct 2021 16:54:41 +0800 Subject: [PATCH 1/5] Add message Conversions prometheus metrics for produce & consume request --- docs/reference-metrics.md | 2 + .../handlers/kop/KafkaRequestHandler.java | 16 ++- .../pulsar/handlers/kop/KopServerStats.java | 2 + .../handlers/kop/MessageFetchContext.java | 16 ++- .../kop/format/AbstractEntryFormatter.java | 7 +- .../handlers/kop/format/DecodeResult.java | 9 +- .../handlers/kop/format/EncodeResult.java | 6 +- .../kop/format/KafkaMixedEntryFormatter.java | 25 ++-- .../kop/format/KafkaV1EntryFormatter.java | 2 +- .../kop/format/PulsarEntryFormatter.java | 2 +- .../ValidationAndOffsetAssignResult.java | 58 ++++++++ .../handlers/kop/utils/ByteBufUtils.java | 5 +- .../handlers/kop/utils/KopLogValidator.java | 125 ++++++++++-------- .../kop/format/EntryFormatterTest.java | 18 ++- .../format/NoHeaderKafkaEntryFormatter.java | 2 +- 15 files changed, 213 insertions(+), 82 deletions(-) create mode 100644 kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/ValidationAndOffsetAssignResult.java diff --git a/docs/reference-metrics.md b/docs/reference-metrics.md index d4defa6662..777d78bc20 100644 --- a/docs/reference-metrics.md +++ b/docs/reference-metrics.md @@ -58,6 +58,7 @@ The KoP metrics are exposed under "/metrics" at port `8000` along with Pulsar me | kop_server_BYTES_IN | Counter | The producer bytes in stats.
Available labels: *topic*, *partition*.
| | kop_server_MESSAGE_IN | Counter | The producer message in stats.
Available labels: *topic*, *partition*.
| | kop_server_BATCH_COUNT_PER_MEMORYRECORDS | Gauge | The number of batches in each memory records| +| kop_server_PRODUCE_MESSAGE_CONVERSIONS | Counter | The producer message conversions in stats.
Available labels: *topic*.
| ### Consumer metrics @@ -70,6 +71,7 @@ The KoP metrics are exposed under "/metrics" at port `8000` along with Pulsar me | kop_server_BYTES_OUT | Counter | The consumer bytes out stats.
Available labels: *topic*, *partition*, *group*.
| | kop_server_MESSAGE_OUT | Counter | The consumer message out stats.
Available labels: *topic*, *partition*, *group*.
| | kop_server_ENTRIES_OUT | Counter | The consumer entries out stats.
Available labels: *topic*, *partition*, *group*.
| +| kop_server_CONSUMER_MESSAGE_CONVERSIONS | Counter | The consumer message conversions in stats.
Available labels: *topic*.
| ### Kop event metrics 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 f49d287ae7..3b1d9d220a 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 @@ -20,6 +20,7 @@ import static io.streamnative.pulsar.handlers.kop.KopServerStats.BYTES_IN; import static io.streamnative.pulsar.handlers.kop.KopServerStats.MESSAGE_IN; import static io.streamnative.pulsar.handlers.kop.KopServerStats.PARTITION_SCOPE; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.PRODUCE_MESSAGE_CONVERSIONS; import static io.streamnative.pulsar.handlers.kop.KopServerStats.TOPIC_SCOPE; import static java.nio.charset.StandardCharsets.UTF_8; import static org.apache.kafka.common.internals.Topic.GROUP_METADATA_TOPIC_NAME; @@ -924,6 +925,7 @@ private void publishMessages(final Optional persistentTopicOpt, final MemoryRecords records = encodeResult.getRecords(); final int numMessages = encodeResult.getNumMessages(); final ByteBuf byteBuf = encodeResult.getEncodedByteBuf(); + final int conversionCount = encodeResult.getConversionCount(); if (!persistentTopicOpt.isPresent()) { encodeResult.recycle(); // It will trigger a retry send of Kafka client @@ -944,7 +946,7 @@ private void publishMessages(final Optional persistentTopicOpt, final Producer producer = KafkaTopicManager.getReferenceProducer(partitionName); producer.updateRates(numMessages, byteBuf.readableBytes()); producer.getTopic().incrementPublishCount(numMessages, byteBuf.readableBytes()); - updateProducerStats(topicPartition, numMessages, byteBuf.readableBytes()); + updateProducerStats(topicPartition, numMessages, byteBuf.readableBytes(), conversionCount); // publish final CompletableFuture offsetFuture = new CompletableFuture<>(); @@ -2725,7 +2727,10 @@ private static MemoryRecords validateRecords(short version, TopicPartition topic return validRecords; } - private void updateProducerStats(final TopicPartition topicPartition, final int numMessages, final int numBytes) { + private void updateProducerStats(final TopicPartition topicPartition, + final int numMessages, + final int numBytes, + final int conversionCount) { requestStats.getStatsLogger() .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) .scopeLabel(PARTITION_SCOPE, String.valueOf((topicPartition.partition()))) @@ -2738,6 +2743,13 @@ private void updateProducerStats(final TopicPartition topicPartition, final int .getCounter(MESSAGE_IN) .add(numMessages); + if (conversionCount > 0) { + requestStats.getStatsLogger() + .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) + .getCounter(PRODUCE_MESSAGE_CONVERSIONS) + .add(conversionCount); + } + RequestStats.BATCH_COUNT_PER_MEMORY_RECORDS_INSTANCE.set(numMessages); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopServerStats.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopServerStats.java index 5372863470..52513a3a62 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopServerStats.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopServerStats.java @@ -61,6 +61,7 @@ public interface KopServerStats { String BYTES_IN = "BYTES_IN"; String MESSAGE_IN = "MESSAGE_IN"; String BATCH_COUNT_PER_MEMORYRECORDS = "BATCH_COUNT_PER_MEMORYRECORDS"; + String PRODUCE_MESSAGE_CONVERSIONS = "PRODUCE_MESSAGE_CONVERSIONS"; /** * FETCH stats. @@ -81,6 +82,7 @@ public interface KopServerStats { String BYTES_OUT = "BYTES_OUT"; String MESSAGE_OUT = "MESSAGE_OUT"; String ENTRIES_OUT = "ENTRIES_OUT"; + String CONSUMER_MESSAGE_CONVERSIONS = "CONSUMER_MESSAGE_CONVERSIONS"; /** * Kop event queue stats. diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java index 126e5cd52e..a113474cec 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java @@ -14,6 +14,7 @@ package io.streamnative.pulsar.handlers.kop; import static io.streamnative.pulsar.handlers.kop.KopServerStats.BYTES_OUT; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.CONSUMER_MESSAGE_CONVERSIONS; import static io.streamnative.pulsar.handlers.kop.KopServerStats.ENTRIES_OUT; import static io.streamnative.pulsar.handlers.kop.KopServerStats.GROUP_SCOPE; import static io.streamnative.pulsar.handlers.kop.KopServerStats.MESSAGE_OUT; @@ -435,7 +436,8 @@ private void handleEntries(final List entries, MathUtils.elapsedNanos(startDecodingEntriesNanos), TimeUnit.NANOSECONDS); decodeResults.add(decodeResult); - MemoryRecords kafkaRecords = decodeResult.getRecords(); + final MemoryRecords kafkaRecords = decodeResult.getRecords(); + final int conversionCount = decodeResult.getConversionCount(); CompletableFuture groupNameFuture = requestHandler .getCurrentConnectedGroup() @@ -467,7 +469,8 @@ private void handleEntries(final List entries, updateConsumerStats(topicPartition, kafkaRecords, entries.size(), - groupName); + groupName, + conversionCount); final List abortedTransactions = (readCommitted ? tc.getAbortedIndexList(partitionData.fetchOffset) : null); responseData.put(topicPartition, new PartitionData<>( @@ -588,7 +591,7 @@ public void markDeleteFailed(ManagedLedgerException e, Object ctx) { } private void updateConsumerStats(final TopicPartition topicPartition, final MemoryRecords records, - int entrySize, final String groupId) { + int entrySize, final String groupId, int conversionCount) { int numMessages = EntryFormatter.parseNumMessages(records); statsLogger.getStatsLogger() @@ -611,5 +614,12 @@ private void updateConsumerStats(final TopicPartition topicPartition, final Memo .scopeLabel(GROUP_SCOPE, groupId) .getCounter(ENTRIES_OUT) .add(entrySize); + + if (conversionCount > 0) { + statsLogger.getStatsLogger() + .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) + .getCounter(CONSUMER_MESSAGE_CONVERSIONS) + .add(conversionCount); + } } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/AbstractEntryFormatter.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/AbstractEntryFormatter.java index d6dbfe8ed9..0163549e71 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/AbstractEntryFormatter.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/AbstractEntryFormatter.java @@ -45,6 +45,7 @@ public abstract class AbstractEntryFormatter implements EntryFormatter { @Override public DecodeResult decode(List entries, byte magic) { int totalSize = 0; + int conversionCount = 0; // batched ByteBuf should be released after sending to client ByteBuf batchedByteBuf = PulsarByteBufAllocator.DEFAULT.directBuffer(totalSize); for (Entry entry : entries) { @@ -63,6 +64,7 @@ public DecodeResult decode(List entries, byte magic) { // down converted, batch magic will be set to client magic ConvertedRecords convertedRecords = memoryRecords.downConvert(magic, startOffset, time); + conversionCount += convertedRecords.recordConversionStats().numRecordsConverted(); final ByteBuf kafkaBuffer = Unpooled.wrappedBuffer(convertedRecords.records().buffer()); totalSize += kafkaBuffer.readableBytes(); @@ -83,6 +85,7 @@ public DecodeResult decode(List entries, byte magic) { } else { final DecodeResult decodeResult = ByteBufUtils.decodePulsarEntryToKafkaRecords(metadata, byteBuf, startOffset, magic); + conversionCount += decodeResult.getConversionCount(); final ByteBuf kafkaBuffer = decodeResult.getOrCreateByteBuf(); totalSize += kafkaBuffer.readableBytes(); batchedByteBuf.writeBytes(kafkaBuffer); @@ -101,7 +104,9 @@ public DecodeResult decode(List entries, byte magic) { } return DecodeResult.get( - MemoryRecords.readableRecords(ByteBufUtils.getNioBuffer(batchedByteBuf)), batchedByteBuf); + MemoryRecords.readableRecords(ByteBufUtils.getNioBuffer(batchedByteBuf)), + batchedByteBuf, + conversionCount); } protected static boolean isKafkaEntryFormat(final MessageMetadata messageMetadata) { diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java index 6c6596847a..b6c074b217 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java @@ -28,18 +28,22 @@ public class DecodeResult { @Getter private MemoryRecords records; private ByteBuf releasedByteBuf; + @Getter + private int conversionCount; private final Recycler.Handle recyclerHandle; public static DecodeResult get(MemoryRecords records) { - return get(records, null); + return get(records, null, 0); } public static DecodeResult get(MemoryRecords records, - ByteBuf releasedByteBuf) { + ByteBuf releasedByteBuf, + int conversionCount) { DecodeResult decodeResult = RECYCLER.get(); decodeResult.records = records; decodeResult.releasedByteBuf = releasedByteBuf; + decodeResult.conversionCount = conversionCount; return decodeResult; } @@ -60,6 +64,7 @@ public void recycle() { releasedByteBuf.release(); releasedByteBuf = null; } + conversionCount = -1; recyclerHandle.recycle(this); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java index 9bc747fce7..5512aa2a86 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java @@ -27,16 +27,19 @@ public class EncodeResult { private MemoryRecords records; private ByteBuf encodedByteBuf; private int numMessages; + private int conversionCount; private final Recycler.Handle recyclerHandle; public static EncodeResult get(MemoryRecords records, ByteBuf encodedByteBuf, - int numMessages) { + int numMessages, + int conversionCount) { EncodeResult encodeResult = RECYCLER.get(); encodeResult.records = records; encodeResult.encodedByteBuf = encodedByteBuf; encodeResult.numMessages = numMessages; + encodeResult.conversionCount = conversionCount; return encodeResult; } @@ -58,6 +61,7 @@ public void recycle() { encodedByteBuf = null; } numMessages = -1; + conversionCount = -1; recyclerHandle.recycle(this); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaMixedEntryFormatter.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaMixedEntryFormatter.java index 455659deb2..e5191104e4 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaMixedEntryFormatter.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaMixedEntryFormatter.java @@ -52,15 +52,19 @@ public EncodeResult encode(final EncodeRequest encodeRequest) { final KopLogValidator.CompressionCodec sourceCodec = getSourceCodec(records); final KopLogValidator.CompressionCodec targetCodec = getTargetCodec(sourceCodec); - final MemoryRecords validRecords = KopLogValidator.validateMessagesAndAssignOffsets(records, - offset, - System.currentTimeMillis(), - sourceCodec, - targetCodec, - false, - RecordBatch.MAGIC_VALUE_V2, - TimestampType.CREATE_TIME, - Long.MAX_VALUE); + final ValidationAndOffsetAssignResult validationAndOffsetAssignResult = + KopLogValidator.validateMessagesAndAssignOffsets(records, + offset, + System.currentTimeMillis(), + sourceCodec, + targetCodec, + false, + RecordBatch.MAGIC_VALUE_V2, + TimestampType.CREATE_TIME, + Long.MAX_VALUE); + + MemoryRecords validRecords = validationAndOffsetAssignResult.getRecords(); + int conversionCount = validationAndOffsetAssignResult.getConversionCount(); final int numMessages = EntryFormatter.parseNumMessages(validRecords); final ByteBuf recordsWrapper = Unpooled.wrappedBuffer(validRecords.buffer()); @@ -69,8 +73,9 @@ public EncodeResult encode(final EncodeRequest encodeRequest) { getMessageMetadataWithNumberMessages(numMessages), recordsWrapper); recordsWrapper.release(); + validationAndOffsetAssignResult.recycle(); - return EncodeResult.get(validRecords, buf, numMessages); + return EncodeResult.get(validRecords, buf, numMessages, conversionCount); } @Override diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaV1EntryFormatter.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaV1EntryFormatter.java index 81efaee8a8..7349980455 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaV1EntryFormatter.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaV1EntryFormatter.java @@ -41,7 +41,7 @@ public EncodeResult encode(final EncodeRequest encodeRequest) { recordsWrapper); recordsWrapper.release(); - return EncodeResult.get(records, buf, numMessages); + return EncodeResult.get(records, buf, numMessages, 0); } @Override diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/PulsarEntryFormatter.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/PulsarEntryFormatter.java index 62bf0874c7..f8ad6fd99f 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/PulsarEntryFormatter.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/PulsarEntryFormatter.java @@ -90,7 +90,7 @@ public EncodeResult encode(final EncodeRequest encodeRequest) { batchedMessageMetadataAndPayload.release(); - return EncodeResult.get(records, buf, numMessages); + return EncodeResult.get(records, buf, numMessages, numMessagesInBatch); } @Override diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/ValidationAndOffsetAssignResult.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/ValidationAndOffsetAssignResult.java new file mode 100644 index 0000000000..95fb8ae388 --- /dev/null +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/ValidationAndOffsetAssignResult.java @@ -0,0 +1,58 @@ +/** + * 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.format; + +import io.netty.util.Recycler; +import lombok.Getter; +import org.apache.kafka.common.record.MemoryRecords; + +/** + * Result of KopLogValidator validateMessagesAndAssignOffsets in KafkaMixedEntryFormatter. + */ +@Getter +public class ValidationAndOffsetAssignResult { + + private MemoryRecords records; + private int conversionCount; + + private final Recycler.Handle recyclerHandle; + + public static ValidationAndOffsetAssignResult get(MemoryRecords records, + int conversionCount) { + ValidationAndOffsetAssignResult validationAndOffsetAssignResult = RECYCLER.get(); + validationAndOffsetAssignResult.records = records; + validationAndOffsetAssignResult.conversionCount = conversionCount; + return validationAndOffsetAssignResult; + } + + private ValidationAndOffsetAssignResult(Recycler.Handle recyclerHandle) { + this.recyclerHandle = recyclerHandle; + } + + private static final Recycler RECYCLER = + new Recycler() { + @Override + protected ValidationAndOffsetAssignResult newObject( + Recycler.Handle handle) { + return new ValidationAndOffsetAssignResult(handle); + } + }; + + public void recycle() { + records = null; + conversionCount = -1; + recyclerHandle.recycle(this); + } + +} diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/ByteBufUtils.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/ByteBufUtils.java index f5df33cebf..8a6ba2fd02 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/ByteBufUtils.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/ByteBufUtils.java @@ -133,8 +133,10 @@ public static DecodeResult decodePulsarEntryToKafkaRecords(final MessageMetadata builder.setProducerState(metadata.getTxnidMostBits(), (short) metadata.getTxnidLeastBits(), 0, true); } + int conversionCount = 0; if (metadata.hasNumMessagesInBatch()) { final int numMessages = metadata.getNumMessagesInBatch(); + conversionCount += numMessages; for (int i = 0; i < numMessages; i++) { final SingleMessageMetadata singleMessageMetadata = new SingleMessageMetadata(); final ByteBuf singleMessagePayload = Commands.deSerializeSingleMessageInBatch( @@ -163,6 +165,7 @@ public static DecodeResult decodePulsarEntryToKafkaRecords(final MessageMetadata singleMessagePayload.release(); } } else { + conversionCount += 1; final long timestamp = (metadata.getEventTime() > 0) ? metadata.getEventTime() : metadata.getPublishTime(); @@ -183,7 +186,7 @@ public static DecodeResult decodePulsarEntryToKafkaRecords(final MessageMetadata final MemoryRecords records = builder.build(); uncompressedPayload.release(); - return DecodeResult.get(records, directBufferOutputStream.getByteBuf()); + return DecodeResult.get(records, directBufferOutputStream.getByteBuf(), conversionCount); } @NonNull diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/KopLogValidator.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/KopLogValidator.java index 9d72e3f714..bd8169adee 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/KopLogValidator.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/KopLogValidator.java @@ -13,6 +13,7 @@ */ package io.streamnative.pulsar.handlers.kop.utils; +import io.streamnative.pulsar.handlers.kop.format.ValidationAndOffsetAssignResult; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Iterator; @@ -46,7 +47,7 @@ public class KopLogValidator { * the offset of the shallow message with the max timestamp and a boolean indicating * whether the message sizes may have changed. */ - public static MemoryRecords validateMessagesAndAssignOffsets(MemoryRecords records, + public static ValidationAndOffsetAssignResult validateMessagesAndAssignOffsets(MemoryRecords records, LongRef offsetCounter, long now, CompressionCodec sourceCodec, @@ -89,13 +90,13 @@ public static MemoryRecords validateMessagesAndAssignOffsets(MemoryRecords recor } } - private static MemoryRecords convertAndAssignOffsetsNonCompressed(MemoryRecords records, - LongRef offsetCounter, - boolean compactedTopic, - long now, - TimestampType timestampType, - long timestampDiffMaxMs, - byte toMagicValue) { + private static ValidationAndOffsetAssignResult convertAndAssignOffsetsNonCompressed(MemoryRecords records, + LongRef offsetCounter, + boolean compactedTopic, + long now, + TimestampType timestampType, + long timestampDiffMaxMs, + byte toMagicValue) { int sizeInBytesAfterConversion = AbstractRecords.estimateSizeInBytes(toMagicValue, offsetCounter.value(), CompressionType.NONE, records.records()); @@ -130,16 +131,19 @@ private static MemoryRecords convertAndAssignOffsetsNonCompressed(MemoryRecords } }); - return builder.build(); + MemoryRecords memoryRecords = builder.build(); + int conversionCount = builder.numRecords(); + + return ValidationAndOffsetAssignResult.get(memoryRecords, conversionCount); } - private static MemoryRecords assignOffsetsNonCompressed(MemoryRecords records, - LongRef offsetCounter, - long now, - boolean compactedTopic, - TimestampType timestampType, - long timestampDiffMaxMs, - byte magic) { + private static ValidationAndOffsetAssignResult assignOffsetsNonCompressed(MemoryRecords records, + LongRef offsetCounter, + long now, + boolean compactedTopic, + TimestampType timestampType, + long timestampDiffMaxMs, + byte magic) { long maxTimestamp = RecordBatch.NO_TIMESTAMP; for (MutableRecordBatch batch : records.batches()) { validateBatch(batch, magic); @@ -172,7 +176,7 @@ private static MemoryRecords assignOffsetsNonCompressed(MemoryRecords records, } } - return records; + return ValidationAndOffsetAssignResult.get(records, 0); } /** @@ -182,15 +186,17 @@ private static MemoryRecords assignOffsetsNonCompressed(MemoryRecords records, * 3. When magic value to use is above 0, but some fields of inner messages need to be overwritten. * 4. Message format conversion is needed. */ - private static MemoryRecords validateMessagesAndAssignOffsetsCompressed(MemoryRecords records, - LongRef offsetCounter, - long now, - CompressionCodec sourceCodec, - CompressionCodec targetCodec, - boolean compactedTopic, - byte toMagic, - TimestampType timestampType, - long timestampDiffMaxMs) { + private static ValidationAndOffsetAssignResult validateMessagesAndAssignOffsetsCompressed( + MemoryRecords records, + LongRef offsetCounter, + long now, + CompressionCodec sourceCodec, + CompressionCodec targetCodec, + boolean compactedTopic, + byte toMagic, + TimestampType timestampType, + long timestampDiffMaxMs) { + // No in place assignment situation 1 and 2 boolean inPlaceAssignment = sourceCodec == targetCodec && toMagic > RecordBatch.MAGIC_VALUE_V0; @@ -246,15 +252,15 @@ private static MemoryRecords validateMessagesAndAssignOffsetsCompressed(MemoryRe } - private static MemoryRecords buildIfPlaceAssignment(boolean inPlaceAssignment, - MemoryRecords records, - ArrayList validatedRecords, - LongRef offsetCounter, - long now, - byte toMagic, - TimestampType timestampType, - long maxTimestamp, - CompressionCodec targetCodec) { + private static ValidationAndOffsetAssignResult buildIfPlaceAssignment(boolean inPlaceAssignment, + MemoryRecords records, + ArrayList validatedRecords, + LongRef offsetCounter, + long now, + byte toMagic, + TimestampType timestampType, + long maxTimestamp, + CompressionCodec targetCodec) { if (inPlaceAssignment) { return buildInPlaceAssignment(records, validatedRecords, @@ -274,13 +280,13 @@ private static MemoryRecords buildIfPlaceAssignment(boolean inPlaceAssignment, } } - private static MemoryRecords buildNoInPlaceAssignment(MemoryRecords records, - ArrayList validatedRecords, - LongRef offsetCounter, - long now, - CompressionCodec targetCodec, - byte toMagic, - TimestampType timestampType) { + private static ValidationAndOffsetAssignResult buildNoInPlaceAssignment(MemoryRecords records, + ArrayList validatedRecords, + LongRef offsetCounter, + long now, + CompressionCodec targetCodec, + byte toMagic, + TimestampType timestampType) { // note that we only reassign offsets for requests coming straight from a producer. // For records with magic V2, there should be exactly one RecordBatch per request, // so the following is all we need to do. For Records @@ -296,13 +302,13 @@ private static MemoryRecords buildNoInPlaceAssignment(MemoryRecords records, first); } - private static MemoryRecords buildInPlaceAssignment(MemoryRecords records, - ArrayList validatedRecords, - LongRef offsetCounter, - long now, - byte toMagic, - TimestampType timestampType, - long maxTimestamp) { + private static ValidationAndOffsetAssignResult buildInPlaceAssignment(MemoryRecords records, + ArrayList validatedRecords, + LongRef offsetCounter, + long now, + byte toMagic, + TimestampType timestampType, + long maxTimestamp) { long currentMaxTimestamp = maxTimestamp; // we can update the batch only and write the compressed payload as is MutableRecordBatch batch = records.batches().iterator().next(); @@ -322,16 +328,16 @@ private static MemoryRecords buildInPlaceAssignment(MemoryRecords records, batch.setPartitionLeaderEpoch(RecordBatch.NO_PARTITION_LEADER_EPOCH); } - return records; + return ValidationAndOffsetAssignResult.get(records, 0); } - private static MemoryRecords buildRecordsAndAssignOffsets(byte magic, - LongRef offsetCounter, - TimestampType timestampType, - CompressionType compressionType, - long logAppendTime, - ArrayList validatedRecords, - MutableRecordBatch first) { + private static ValidationAndOffsetAssignResult buildRecordsAndAssignOffsets(byte magic, + LongRef offsetCounter, + TimestampType timestampType, + CompressionType compressionType, + long logAppendTime, + ArrayList validatedRecords, + MutableRecordBatch first) { long producerId = first.producerId(); short producerEpoch = first.producerEpoch(); int baseSequence = first.baseSequence(); @@ -356,7 +362,10 @@ private static MemoryRecords buildRecordsAndAssignOffsets(byte magic, builder.appendWithOffset(offsetCounter.getAndIncrement(), record); }); - return builder.build(); + MemoryRecords memoryRecords = builder.build(); + int conversionCount = builder.numRecords(); + + return ValidationAndOffsetAssignResult.get(memoryRecords, conversionCount); } private static void validateBatch(RecordBatch batch, byte toMagic) { diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/format/EntryFormatterTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/format/EntryFormatterTest.java index 7531343355..5ff8f45438 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/format/EntryFormatterTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/format/EntryFormatterTest.java @@ -54,7 +54,13 @@ public class EntryFormatterTest { private static EntryFormatter pulsarFormatter; private static EntryFormatter kafkaV1Formatter; private static EntryFormatter kafkaMixedFormatter; - private static long baseOffset = 100; + // Kafka's absolute message offset is allocated by the Kafka server, + // so each batch of messages on the client always starts from 0. + // If it starts from non-zero, KopLogValidator will think that + // the internal offset of these messages is damaged and will reconstruct the message. + // I found this problem after adding the message conversion metrics, + // so here I changed it to 0. + private static final long baseOffset = 0; private void init() { pulsarServiceConfiguration.setEntryFormat("pulsar"); @@ -104,14 +110,24 @@ public void testEntryFormatterEncode(CompressionType compressionType, byte magic EncodeResult encodeResult; // Verify that KafkaV1EntryFormatter cannot fix the wrong relative offset. encodeResult = kafkaV1Formatter.encode(EncodeRequest.get(records)); + Assert.assertEquals(0, encodeResult.getConversionCount()); checkWrongOffset(encodeResult.getRecords(), compressionType, magic); // Verify that PulsarEntryFormatter cannot fix the wrong relative offset. encodeResult = pulsarFormatter.encode(EncodeRequest.get(records)); + Assert.assertEquals(NUM_MESSAGES, encodeResult.getConversionCount()); checkWrongOffset(encodeResult.getRecords(), compressionType, magic); // Verify that KafkaMixedEntryFormatter can fix incorrect relative offset. encodeResult = kafkaMixedFormatter.encode(EncodeRequest.get(records)); + if (magic == RecordBatch.MAGIC_VALUE_V2) { + // After changing baseOffset to 0, + // KafkaMixedFormatter will not reconstruct the message with magic=2, + // so the conversion count here is always 0. + Assert.assertEquals(0, encodeResult.getConversionCount()); + } else { + Assert.assertEquals(NUM_MESSAGES, encodeResult.getConversionCount()); + } checkCorrectOffset(encodeResult.getRecords()); } diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/format/NoHeaderKafkaEntryFormatter.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/format/NoHeaderKafkaEntryFormatter.java index db7d43b4a9..67268eb248 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/format/NoHeaderKafkaEntryFormatter.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/format/NoHeaderKafkaEntryFormatter.java @@ -28,7 +28,7 @@ public EncodeResult encode(final EncodeRequest encodeRequest) { final MemoryRecords records = encodeRequest.getRecords(); final int numMessages = EntryFormatter.parseNumMessages(records); // The difference from KafkaEntryFormatter is here we don't add the header - return EncodeResult.get(records, Unpooled.wrappedBuffer(records.buffer()), numMessages); + return EncodeResult.get(records, Unpooled.wrappedBuffer(records.buffer()), numMessages, 0); } @Override From a18b8645b3cc6ca5b8e67a029cab3270f22747f9 Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Mon, 1 Nov 2021 21:25:41 +0800 Subject: [PATCH 2/5] Addressed comments --- .../handlers/kop/KafkaRequestHandler.java | 38 +-------------- .../handlers/kop/MessageFetchContext.java | 47 +------------------ .../handlers/kop/format/DecodeResult.java | 31 ++++++++++++ .../handlers/kop/format/EncodeResult.java | 31 ++++++++++++ 4 files changed, 65 insertions(+), 82 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 fed8ee7ccc..d38f8a26e1 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 @@ -17,11 +17,6 @@ import static com.google.common.base.Preconditions.checkState; import static io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration.TENANT_ALLNAMESPACES_PLACEHOLDER; import static io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration.TENANT_PLACEHOLDER; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.BYTES_IN; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.MESSAGE_IN; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.PARTITION_SCOPE; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.PRODUCE_MESSAGE_CONVERSIONS; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.TOPIC_SCOPE; 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; @@ -168,7 +163,6 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.Utils; import org.apache.pulsar.broker.PulsarService; -import org.apache.pulsar.broker.service.Producer; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; @@ -901,7 +895,6 @@ private void publishMessages(final Optional persistentTopicOpt, final MemoryRecords records = encodeResult.getRecords(); final int numMessages = encodeResult.getNumMessages(); final ByteBuf byteBuf = encodeResult.getEncodedByteBuf(); - final int conversionCount = encodeResult.getConversionCount(); if (!persistentTopicOpt.isPresent()) { encodeResult.recycle(); // It will trigger a retry send of Kafka client @@ -919,10 +912,7 @@ private void publishMessages(final Optional persistentTopicOpt, topicManager.registerProducerInPersistentTopic(partitionName, persistentTopic); // collect metrics - final Producer producer = KafkaTopicManager.getReferenceProducer(partitionName); - producer.updateRates(numMessages, byteBuf.readableBytes()); - producer.getTopic().incrementPublishCount(numMessages, byteBuf.readableBytes()); - updateProducerStats(topicPartition, numMessages, byteBuf.readableBytes(), conversionCount); + encodeResult.updateProducerStats(topicPartition, requestStats); // publish final CompletableFuture offsetFuture = new CompletableFuture<>(); @@ -2586,32 +2576,6 @@ private static MemoryRecords validateRecords(short version, TopicPartition topic return validRecords; } - private void updateProducerStats(final TopicPartition topicPartition, - final int numMessages, - final int numBytes, - final int conversionCount) { - requestStats.getStatsLogger() - .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) - .scopeLabel(PARTITION_SCOPE, String.valueOf((topicPartition.partition()))) - .getCounter(BYTES_IN) - .add(numBytes); - - requestStats.getStatsLogger() - .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) - .scopeLabel(PARTITION_SCOPE, String.valueOf((topicPartition.partition()))) - .getCounter(MESSAGE_IN) - .add(numMessages); - - if (conversionCount > 0) { - requestStats.getStatsLogger() - .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) - .getCounter(PRODUCE_MESSAGE_CONVERSIONS) - .add(conversionCount); - } - - RequestStats.BATCH_COUNT_PER_MEMORY_RECORDS_INSTANCE.set(numMessages); - } - @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/MessageFetchContext.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java index d08d5d8beb..e3662e754a 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java @@ -13,13 +13,6 @@ */ package io.streamnative.pulsar.handlers.kop; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.BYTES_OUT; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.CONSUMER_MESSAGE_CONVERSIONS; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.ENTRIES_OUT; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.GROUP_SCOPE; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.MESSAGE_OUT; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.PARTITION_SCOPE; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.TOPIC_SCOPE; import static org.apache.kafka.common.protocol.CommonFields.THROTTLE_TIME_MS; import com.google.common.collect.Lists; @@ -29,7 +22,6 @@ import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.exceptions.KoPMessageMetadataNotFoundException; import io.streamnative.pulsar.handlers.kop.format.DecodeResult; -import io.streamnative.pulsar.handlers.kop.format.EntryFormatter; import io.streamnative.pulsar.handlers.kop.security.auth.Resource; import io.streamnative.pulsar.handlers.kop.security.auth.ResourceType; import io.streamnative.pulsar.handlers.kop.utils.GroupIdUtils; @@ -449,7 +441,6 @@ private void handleEntries(final List entries, decodeResults.add(decodeResult); final MemoryRecords kafkaRecords = decodeResult.getRecords(); - final int conversionCount = decodeResult.getConversionCount(); CompletableFuture groupNameFuture = requestHandler .getCurrentConnectedGroup() @@ -478,11 +469,10 @@ private void handleEntries(final List entries, groupName = ""; } // collect consumer metrics - updateConsumerStats(topicPartition, - kafkaRecords, + decodeResult.updateConsumerStats(topicPartition, entries.size(), groupName, - conversionCount); + statsLogger); final List abortedTransactions = (readCommitted ? tc.getAbortedIndexList(partitionData.fetchOffset) : null); responseData.put(topicPartition, new PartitionData<>( @@ -601,37 +591,4 @@ public void markDeleteFailed(ManagedLedgerException e, Object ctx) { } }, null); } - - private void updateConsumerStats(final TopicPartition topicPartition, final MemoryRecords records, - int entrySize, final String groupId, int conversionCount) { - int numMessages = EntryFormatter.parseNumMessages(records); - - statsLogger.getStatsLogger() - .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) - .scopeLabel(PARTITION_SCOPE, String.valueOf(topicPartition.partition())) - .scopeLabel(GROUP_SCOPE, groupId) - .getCounter(BYTES_OUT) - .add(records.sizeInBytes()); - - statsLogger.getStatsLogger() - .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) - .scopeLabel(PARTITION_SCOPE, String.valueOf(topicPartition.partition())) - .scopeLabel(GROUP_SCOPE, groupId) - .getCounter(MESSAGE_OUT) - .add(numMessages); - - statsLogger.getStatsLogger() - .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) - .scopeLabel(PARTITION_SCOPE, String.valueOf(topicPartition.partition())) - .scopeLabel(GROUP_SCOPE, groupId) - .getCounter(ENTRIES_OUT) - .add(entrySize); - - if (conversionCount > 0) { - statsLogger.getStatsLogger() - .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) - .getCounter(CONSUMER_MESSAGE_CONVERSIONS) - .add(conversionCount); - } - } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java index b6c074b217..198408d68b 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java @@ -13,11 +13,22 @@ */ package io.streamnative.pulsar.handlers.kop.format; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.BYTES_OUT; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.CONSUMER_MESSAGE_CONVERSIONS; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.ENTRIES_OUT; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.GROUP_SCOPE; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.MESSAGE_OUT; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.PARTITION_SCOPE; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.TOPIC_SCOPE; + import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; import io.netty.util.Recycler; +import io.streamnative.pulsar.handlers.kop.RequestStats; +import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; import lombok.Getter; import lombok.NonNull; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.record.MemoryRecords; /** @@ -76,4 +87,24 @@ public void recycle() { } } + public void updateConsumerStats(final TopicPartition topicPartition, + int entrySize, + final String groupId, + RequestStats statsLogger) { + final int numMessages = EntryFormatter.parseNumMessages(records); + + final StatsLogger statsLoggerForThisPartition = statsLogger.getStatsLogger() + .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) + .scopeLabel(PARTITION_SCOPE, String.valueOf(topicPartition.partition())); + + statsLoggerForThisPartition.getCounter(CONSUMER_MESSAGE_CONVERSIONS).add(conversionCount); + + final StatsLogger statsLoggerForThisGroup = statsLoggerForThisPartition.scopeLabel(GROUP_SCOPE, groupId); + + statsLoggerForThisGroup.getCounter(BYTES_OUT).add(records.sizeInBytes()); + statsLoggerForThisGroup.getCounter(MESSAGE_OUT).add(numMessages); + statsLoggerForThisGroup.getCounter(ENTRIES_OUT).add(entrySize); + + } + } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java index 5512aa2a86..2dd894b3ac 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java @@ -13,10 +13,22 @@ */ package io.streamnative.pulsar.handlers.kop.format; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.BYTES_IN; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.MESSAGE_IN; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.PARTITION_SCOPE; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.PRODUCE_MESSAGE_CONVERSIONS; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.TOPIC_SCOPE; + import io.netty.buffer.ByteBuf; import io.netty.util.Recycler; +import io.streamnative.pulsar.handlers.kop.KafkaTopicManager; +import io.streamnative.pulsar.handlers.kop.RequestStats; +import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; +import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import lombok.Getter; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.record.MemoryRecords; +import org.apache.pulsar.broker.service.Producer; /** * Result of encode in entry formatter. @@ -65,4 +77,23 @@ public void recycle() { recyclerHandle.recycle(this); } + public void updateProducerStats(final TopicPartition topicPartition, + final RequestStats requestStats) { + final int numBytes = encodedByteBuf.readableBytes(); + + final Producer producer = KafkaTopicManager.getReferenceProducer(KopTopic.toString(topicPartition)); + producer.updateRates(numMessages, numMessages); + producer.getTopic().incrementPublishCount(numMessages, numBytes); + + final StatsLogger statsLoggerForThisPartition = requestStats.getStatsLogger() + .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) + .scopeLabel(PARTITION_SCOPE, String.valueOf(topicPartition.partition())); + + statsLoggerForThisPartition.getCounter(BYTES_IN).add(numBytes); + statsLoggerForThisPartition.getCounter(MESSAGE_IN).add(numMessages); + statsLoggerForThisPartition.getCounter(PRODUCE_MESSAGE_CONVERSIONS).add(conversionCount); + + RequestStats.BATCH_COUNT_PER_MEMORY_RECORDS_INSTANCE.set(numMessages); + } + } From 83b8050a3df840b652e3a2ce2cc722f1e1cadf16 Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Tue, 2 Nov 2021 13:03:56 +0800 Subject: [PATCH 3/5] rename metric name --- .../io/streamnative/pulsar/handlers/kop/KopServerStats.java | 2 +- .../streamnative/pulsar/handlers/kop/format/DecodeResult.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopServerStats.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopServerStats.java index 52513a3a62..65c3cea99f 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopServerStats.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopServerStats.java @@ -82,7 +82,7 @@ public interface KopServerStats { String BYTES_OUT = "BYTES_OUT"; String MESSAGE_OUT = "MESSAGE_OUT"; String ENTRIES_OUT = "ENTRIES_OUT"; - String CONSUMER_MESSAGE_CONVERSIONS = "CONSUMER_MESSAGE_CONVERSIONS"; + String CONSUME_MESSAGE_CONVERSIONS = "CONSUME_MESSAGE_CONVERSIONS"; /** * Kop event queue stats. diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java index 198408d68b..c7390077ad 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/DecodeResult.java @@ -14,7 +14,7 @@ package io.streamnative.pulsar.handlers.kop.format; import static io.streamnative.pulsar.handlers.kop.KopServerStats.BYTES_OUT; -import static io.streamnative.pulsar.handlers.kop.KopServerStats.CONSUMER_MESSAGE_CONVERSIONS; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.CONSUME_MESSAGE_CONVERSIONS; import static io.streamnative.pulsar.handlers.kop.KopServerStats.ENTRIES_OUT; import static io.streamnative.pulsar.handlers.kop.KopServerStats.GROUP_SCOPE; import static io.streamnative.pulsar.handlers.kop.KopServerStats.MESSAGE_OUT; @@ -97,7 +97,7 @@ public void updateConsumerStats(final TopicPartition topicPartition, .scopeLabel(TOPIC_SCOPE, topicPartition.topic()) .scopeLabel(PARTITION_SCOPE, String.valueOf(topicPartition.partition())); - statsLoggerForThisPartition.getCounter(CONSUMER_MESSAGE_CONVERSIONS).add(conversionCount); + statsLoggerForThisPartition.getCounter(CONSUME_MESSAGE_CONVERSIONS).add(conversionCount); final StatsLogger statsLoggerForThisGroup = statsLoggerForThisPartition.scopeLabel(GROUP_SCOPE, groupId); From 251b096c3b640f272c82facc9ff50e3ca825358e Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Tue, 2 Nov 2021 13:13:32 +0800 Subject: [PATCH 4/5] fix docs --- docs/reference-metrics.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/reference-metrics.md b/docs/reference-metrics.md index b81cc97d56..8a83a0b93e 100644 --- a/docs/reference-metrics.md +++ b/docs/reference-metrics.md @@ -58,7 +58,7 @@ The KoP metrics are exposed under "/metrics" at port `8000` along with Pulsar me | kop_server_BYTES_IN | Counter | The producer bytes in stats.
Available labels: *topic*, *partition*.
  • *topic*: the topic name to produce.
  • *partition*: the partition id for the topic to produce
| | kop_server_MESSAGE_IN | Counter | The producer message in stats.
Available labels: *topic*, *partition*.
  • *topic*: the topic name to produce.
  • *partition*: the partition id for the topic to produce
| | kop_server_BATCH_COUNT_PER_MEMORYRECORDS | Gauge | The number of batches in each memory records| -| kop_server_PRODUCE_MESSAGE_CONVERSIONS | Counter | The producer message conversions in stats.
Available labels: *topic*.
  • *topic*: the topic name to produce.
| +| kop_server_PRODUCE_MESSAGE_CONVERSIONS | Counter | The producer message conversions in stats.
Available labels: *topic*, *partition*.
  • *topic*: the topic name to produce.
  • *partition*: the partition id for the topic to produce
| ### Consumer metrics @@ -71,7 +71,7 @@ The KoP metrics are exposed under "/metrics" at port `8000` along with Pulsar me | kop_server_BYTES_OUT | Counter | The consumer bytes out stats.
Available labels: *topic*, *partition*, *group*.
  • *topic*: the topic name to consume.
  • *partition*: the partition id for the topic to consume
  • *group*: the group id for consumer to consumer message from topic-partition
| | kop_server_MESSAGE_OUT | Counter | The consumer message out stats.
Available labels: *topic*, *partition*, *group*.
  • *topic*: the topic name to consume.
  • *partition*: the partition id for the topic to consume
  • *group*: the group id for consumer to consumer message from topic-partition
| | kop_server_ENTRIES_OUT | Counter | The consumer entries out stats.
Available labels: *topic*, *partition*, *group*.
  • *topic*: the topic name to consume.
  • *partition*: the partition id for the topic to consume
  • *group*: the group id for consumer to consumer message from topic-partition
| -| kop_server_CONSUMER_MESSAGE_CONVERSIONS | Counter | The consumer message conversions in stats.
Available labels: *topic*.
  • *topic*: the topic name to consume.
| +| kop_server_CONSUME_MESSAGE_CONVERSIONS | Counter | The consumer message conversions in stats.
Available labels: *topic*, *partition*.
  • *topic*: the topic name to consume.
  • *partition*: the partition id for the topic to consume
| ### Kop event metrics From 46297cfb476c29c6e6fecb3ffe7e3399ebaff906 Mon Sep 17 00:00:00 2001 From: wenbingshen Date: Tue, 2 Nov 2021 14:37:08 +0800 Subject: [PATCH 5/5] fix failed test --- .../streamnative/pulsar/handlers/kop/format/EncodeResult.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java index 2dd894b3ac..0225b05aee 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/EncodeResult.java @@ -82,7 +82,7 @@ public void updateProducerStats(final TopicPartition topicPartition, final int numBytes = encodedByteBuf.readableBytes(); final Producer producer = KafkaTopicManager.getReferenceProducer(KopTopic.toString(topicPartition)); - producer.updateRates(numMessages, numMessages); + producer.updateRates(numMessages, numBytes); producer.getTopic().incrementPublishCount(numMessages, numBytes); final StatsLogger statsLoggerForThisPartition = requestStats.getStatsLogger()