diff --git a/docs/reference-metrics.md b/docs/reference-metrics.md
index 22582dfe33..8a83a0b93e 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*.
- *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*, *partition*. - *topic*: the topic name to produce.
- *partition*: the partition id for the topic to produce
|
### 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*. - *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_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
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 5731e34a96..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,10 +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.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;
@@ -167,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;
@@ -917,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());
+ encodeResult.updateProducerStats(topicPartition, requestStats);
// publish
final CompletableFuture offsetFuture = new CompletableFuture<>();
@@ -2584,22 +2576,6 @@ private static MemoryRecords validateRecords(short version, TopicPartition topic
return validRecords;
}
- private void updateProducerStats(final TopicPartition topicPartition, final int numMessages, final int numBytes) {
- 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);
-
- 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/KopServerStats.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopServerStats.java
index 5372863470..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
@@ -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 CONSUME_MESSAGE_CONVERSIONS = "CONSUME_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 70c909c296..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,12 +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.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;
@@ -28,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;
@@ -447,7 +440,7 @@ private void handleEntries(final List entries,
MathUtils.elapsedNanos(startDecodingEntriesNanos), TimeUnit.NANOSECONDS);
decodeResults.add(decodeResult);
- MemoryRecords kafkaRecords = decodeResult.getRecords();
+ final MemoryRecords kafkaRecords = decodeResult.getRecords();
CompletableFuture groupNameFuture = requestHandler
.getCurrentConnectedGroup()
@@ -476,10 +469,10 @@ private void handleEntries(final List entries,
groupName = "";
}
// collect consumer metrics
- updateConsumerStats(topicPartition,
- kafkaRecords,
+ decodeResult.updateConsumerStats(topicPartition,
entries.size(),
- groupName);
+ groupName,
+ statsLogger);
final List abortedTransactions =
(readCommitted ? tc.getAbortedIndexList(partitionData.fetchOffset) : null);
responseData.put(topicPartition, new PartitionData<>(
@@ -598,30 +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 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);
- }
}
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..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
@@ -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.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;
+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;
/**
@@ -28,18 +39,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 +75,7 @@ public void recycle() {
releasedByteBuf.release();
releasedByteBuf = null;
}
+ conversionCount = -1;
recyclerHandle.recycle(this);
}
@@ -71,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(CONSUME_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 9bc747fce7..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
@@ -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.
@@ -27,16 +39,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,7 +73,27 @@ public void recycle() {
encodedByteBuf = null;
}
numMessages = -1;
+ conversionCount = -1;
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, numBytes);
+ 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);
+ }
+
}
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