From 20db3689e6b63e6f3353c815b45c4dc9d3647a5d Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Wed, 20 Jan 2021 21:17:36 +0800 Subject: [PATCH 1/8] Upgrade pulsar with LightProto support --- .../pulsar/handlers/kop/InternalProducer.java | 4 +- .../handlers/kop/KafkaRequestHandler.java | 17 +++--- .../handlers/kop/MessageFetchContext.java | 2 +- .../kop/format/KafkaEntryFormatter.java | 21 ++++--- .../kop/format/PulsarEntryFormatter.java | 55 ++++++++----------- .../handlers/kop/utils/ByteBufUtils.java | 12 ++-- .../kop/utils/OffsetSearchPredicate.java | 4 +- .../handlers/kop/utils/OffsetFinderTest.java | 15 ++--- pom.xml | 2 +- 9 files changed, 60 insertions(+), 72 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/InternalProducer.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/InternalProducer.java index 6f35741d9c..2b916d8d84 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/InternalProducer.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/InternalProducer.java @@ -19,7 +19,7 @@ import org.apache.pulsar.broker.service.Producer; import org.apache.pulsar.broker.service.ServerCnx; import org.apache.pulsar.broker.service.Topic; -import org.apache.pulsar.common.api.proto.PulsarApi; +import org.apache.pulsar.common.api.proto.ProducerAccessMode; /** * InternalServerCnx, this only used to construct internalProducer / internalConsumer. @@ -34,7 +34,7 @@ public InternalProducer(Topic topic, ServerCnx cnx, long producerId, String producerName) { super(topic, cnx, producerId, producerName, null, false, null, null, 0, false, - PulsarApi.ProducerAccessMode.Shared, Optional.empty()); + ProducerAccessMode.Shared, Optional.empty()); this.serverCnx = cnx; } 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 d53078dc4c..50da0e868b 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 @@ -136,8 +136,8 @@ import org.apache.pulsar.broker.loadbalance.LoadManager; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.admin.PulsarAdmin; -import org.apache.pulsar.common.api.proto.PulsarApi; -import org.apache.pulsar.common.api.proto.PulsarMarkers; +import org.apache.pulsar.common.api.proto.MarkerType; +import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.partition.PartitionedTopicMetadata; @@ -1438,26 +1438,25 @@ private CompletableFuture writeTxnMarker(TopicPartition topicPartition, private ByteBuf generateTxnMarker(TransactionResult transactionResult, long producerId, short producerEpoch) { ControlRecordType controlRecordType; - PulsarMarkers.MarkerType markerType; + MarkerType markerType; if (transactionResult.equals(TransactionResult.COMMIT)) { - markerType = PulsarMarkers.MarkerType.TXN_COMMIT; + markerType = MarkerType.TXN_COMMIT; controlRecordType = ControlRecordType.COMMIT; } else { - markerType = PulsarMarkers.MarkerType.TXN_ABORT; + markerType = MarkerType.TXN_ABORT; controlRecordType = ControlRecordType.ABORT; } EndTransactionMarker marker = new EndTransactionMarker(controlRecordType, 0); MemoryRecords memoryRecords = MemoryRecords.withEndTransactionMarker(producerId, producerEpoch, marker); ByteBuf byteBuf = Unpooled.wrappedBuffer(memoryRecords.buffer()); - PulsarApi.MessageMetadata messageMetadata = PulsarApi.MessageMetadata.newBuilder() + MessageMetadata messageMetadata = new MessageMetadata() .setTxnidMostBits(producerId) .setTxnidLeastBits(producerEpoch) - .setMarkerType(markerType.getNumber()) + .setMarkerType(markerType.getValue()) .setPublishTime(SystemTime.SYSTEM.milliseconds()) .setProducerName("") - .setSequenceId(0L) - .build(); + .setSequenceId(0L); return Commands.serializeMetadataAndPayload(Commands.ChecksumType.None, messageMetadata, byteBuf); } 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 6c8dd9b90b..90542dd017 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 @@ -455,7 +455,7 @@ public void readEntriesFailed(ManagedLedgerException e, Object o) { readFuture.completeExceptionally(e); } - }, null); + }, null, PositionImpl.latest); readFutures.putIfAbsent(cursorOffsetPair.getKey(), readFuture); }); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaEntryFormatter.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaEntryFormatter.java index d5b648a0c9..46f61c82f7 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaEntryFormatter.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/format/KafkaEntryFormatter.java @@ -26,7 +26,7 @@ import org.apache.kafka.common.record.MemoryRecordsBuilder; import org.apache.kafka.common.record.Record; import org.apache.kafka.common.record.TimestampType; -import org.apache.pulsar.common.api.proto.PulsarApi; +import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.protocol.Commands; /** @@ -73,17 +73,16 @@ public MemoryRecords decode(List entries, byte magic) { return builder.build(); } - private static PulsarApi.MessageMetadata getMessageMetadataWithNumberMessages(int numMessages) { - final PulsarApi.MessageMetadata.Builder builder = PulsarApi.MessageMetadata.newBuilder(); - builder.addProperties(PulsarApi.KeyValue.newBuilder() + private static MessageMetadata getMessageMetadataWithNumberMessages(int numMessages) { + final MessageMetadata metadata = new MessageMetadata(); + metadata.addProperty() .setKey("entry.format") - .setValue(EntryFormatterFactory.EntryFormat.KAFKA.name().toLowerCase()) - .build()); - builder.setProducerName(""); - builder.setSequenceId(0L); - builder.setPublishTime(System.currentTimeMillis()); - builder.setNumMessagesInBatch(numMessages); - return builder.build(); + .setValue(EntryFormatterFactory.EntryFormat.KAFKA.name().toLowerCase()); + metadata.setProducerName(""); + metadata.setSequenceId(0L); + metadata.setPublishTime(System.currentTimeMillis()); + metadata.setNumMessagesInBatch(numMessages); + return metadata; } } 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 917d8e6c29..c346617916 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 @@ -41,11 +41,11 @@ import org.apache.pulsar.client.impl.MessageImpl; import org.apache.pulsar.client.impl.TypedMessageBuilderImpl; import org.apache.pulsar.common.allocator.PulsarByteBufAllocator; -import org.apache.pulsar.common.api.proto.PulsarApi; -import org.apache.pulsar.common.api.proto.PulsarApi.KeyValue; -import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata; -import org.apache.pulsar.common.api.proto.PulsarApi.SingleMessageMetadata; -import org.apache.pulsar.common.api.proto.PulsarMarkers; +import org.apache.pulsar.common.api.proto.CompressionType; +import org.apache.pulsar.common.api.proto.KeyValue; +import org.apache.pulsar.common.api.proto.MarkerType; +import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.apache.pulsar.common.api.proto.SingleMessageMetadata; import org.apache.pulsar.common.compression.CompressionCodec; import org.apache.pulsar.common.compression.CompressionCodecProvider; import org.apache.pulsar.common.protocol.Commands; @@ -70,30 +70,30 @@ public ByteBuf encode(final MemoryRecords records, final int numMessages) { long sequenceId = -1; // TODO: handle different compression type - PulsarApi.CompressionType compressionType = PulsarApi.CompressionType.NONE; + CompressionType compressionType = CompressionType.NONE; ByteBuf batchedMessageMetadataAndPayload = PulsarByteBufAllocator.DEFAULT .buffer(Math.min(INITIAL_BATCH_BUFFER_SIZE, MAX_MESSAGE_BATCH_SIZE_BYTES)); List> messages = Lists.newArrayListWithExpectedSize(numMessages); - MessageMetadata.Builder messageMetaBuilder = MessageMetadata.newBuilder(); + final MessageMetadata msgMetadata = new MessageMetadata(); records.batches().forEach(recordBatch -> { StreamSupport.stream(recordBatch.spliterator(), true).forEachOrdered(record -> { MessageImpl message = recordToEntry(record); messages.add(message); - if (messageMetaBuilder.getPublishTime() <= 0) { - messageMetaBuilder.setPublishTime(message.getPublishTime()); + if (msgMetadata.getPublishTime() <= 0) { + msgMetadata.setPublishTime(message.getPublishTime()); } if (recordBatch.isTransactional()) { - messageMetaBuilder.setTxnidMostBits(recordBatch.producerId()); - messageMetaBuilder.setTxnidLeastBits(recordBatch.producerEpoch()); + msgMetadata.setTxnidMostBits(recordBatch.producerId()); + msgMetadata.setTxnidLeastBits(recordBatch.producerEpoch()); } }); }); for (MessageImpl message : messages) { if (++numMessagesInBatch == 1) { - sequenceId = Commands.initBatchMessageMetadata(messageMetaBuilder, message.getMessageBuilder()); + sequenceId = Commands.initBatchMessageMetadata(msgMetadata, message.getMessageBuilder()); } currentBatchSizeBytes += message.getDataBuffer().readableBytes(); if (log.isDebugEnabled()) { @@ -101,29 +101,24 @@ public ByteBuf encode(final MemoryRecords records, final int numMessages) { sequenceId, numMessagesInBatch, currentBatchSizeBytes); } - PulsarApi.MessageMetadata.Builder msgBuilder = message.getMessageBuilder(); + final MessageMetadata msgBuilder = message.getMessageBuilder(); batchedMessageMetadataAndPayload = Commands.serializeSingleMessageInBatchWithPayload(msgBuilder, message.getDataBuffer(), batchedMessageMetadataAndPayload); - msgBuilder.recycle(); } int uncompressedSize = batchedMessageMetadataAndPayload.readableBytes(); - if (PulsarApi.CompressionType.NONE != compressionType) { - messageMetaBuilder.setCompression(compressionType); - messageMetaBuilder.setUncompressedSize(uncompressedSize); + if (CompressionType.NONE != compressionType) { + msgMetadata.setCompression(compressionType); + msgMetadata.setUncompressedSize(uncompressedSize); } - messageMetaBuilder.setNumMessagesInBatch(numMessagesInBatch); - - MessageMetadata msgMetadata = messageMetaBuilder.build(); + msgMetadata.setNumMessagesInBatch(numMessagesInBatch); ByteBuf buf = Commands.serializeMetadataAndPayload(ChecksumType.Crc32c, msgMetadata, batchedMessageMetadataAndPayload); - messageMetaBuilder.recycle(); - msgMetadata.recycle(); batchedMessageMetadataAndPayload.release(); return buf; @@ -141,8 +136,8 @@ public MemoryRecords decode(final List entries, final byte magic) { // Uncompress the payload if necessary MessageMetadata msgMetadata = Commands.parseMessageMetadata(metadataAndPayload); - if (msgMetadata.getMarkerType() == PulsarMarkers.MarkerType.TXN_COMMIT_VALUE - || msgMetadata.getMarkerType() == PulsarMarkers.MarkerType.TXN_ABORT_VALUE) { + if (msgMetadata.getMarkerType() == MarkerType.TXN_COMMIT_VALUE + || msgMetadata.getMarkerType() == MarkerType.TXN_ABORT_VALUE) { MemoryRecords memoryRecords = MemoryRecords.withEndTransactionMarker( baseOffset, msgMetadata.getPublishTime(), @@ -150,7 +145,7 @@ public MemoryRecords decode(final List entries, final byte magic) { msgMetadata.getTxnidMostBits(), (short) msgMetadata.getTxnidLeastBits(), new EndTransactionMarker( - msgMetadata.getMarkerType() == PulsarMarkers.MarkerType.TXN_COMMIT_VALUE + msgMetadata.getMarkerType() == MarkerType.TXN_COMMIT_VALUE ? ControlRecordType.COMMIT : ControlRecordType.ABORT, 0)); byteBuffer.put(memoryRecords.buffer()); return; @@ -204,15 +199,13 @@ public MemoryRecords decode(final List entries, final byte magic) { log.debug(" processing message num - {} in batch", i); } try { - SingleMessageMetadata.Builder singleMessageMetadataBuilder = SingleMessageMetadata - .newBuilder(); + final SingleMessageMetadata singleMessageMetadata = new SingleMessageMetadata(); ByteBuf singleMessagePayload = Commands.deSerializeSingleMessageInBatch(payload, - singleMessageMetadataBuilder, i, numMessages); + singleMessageMetadata, i, numMessages); - SingleMessageMetadata singleMessageMetadata = singleMessageMetadataBuilder.build(); Header[] headers = getHeadersFromMetadata(singleMessageMetadata.getPropertiesList()); - final ByteBuffer value = (singleMessageMetadata.getNullValue()) + final ByteBuffer value = (singleMessageMetadata.isNullValue()) ? null : ByteBufUtils.getNioBuffer(singleMessagePayload); builder.appendWithOffset( @@ -223,7 +216,6 @@ public MemoryRecords decode(final List entries, final byte magic) { value, headers); singleMessagePayload.release(); - singleMessageMetadataBuilder.recycle(); } catch (IOException e) { log.error("Meet IOException: {}", e); throw new UncheckedIOException(e); @@ -240,7 +232,6 @@ public MemoryRecords decode(final List entries, final byte magic) { headers); } - msgMetadata.recycle(); payload.release(); entry.release(); 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 0450ba8963..5ea7576ed9 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 @@ -18,7 +18,9 @@ import io.netty.buffer.ByteBuf; import java.nio.ByteBuffer; import java.util.Base64; -import org.apache.pulsar.common.api.proto.PulsarApi; + +import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.apache.pulsar.common.api.proto.SingleMessageMetadata; /** @@ -26,9 +28,9 @@ */ public class ByteBufUtils { - public static ByteBuffer getKeyByteBuffer(PulsarApi.SingleMessageMetadata messageMetadata) { + public static ByteBuffer getKeyByteBuffer(SingleMessageMetadata messageMetadata) { if (messageMetadata.hasOrderingKey()) { - return messageMetadata.getOrderingKey().asReadOnlyByteBuffer(); + return ByteBuffer.wrap(messageMetadata.getOrderingKey()).asReadOnlyBuffer(); } String key = messageMetadata.getPartitionKey(); @@ -40,9 +42,9 @@ public static ByteBuffer getKeyByteBuffer(PulsarApi.SingleMessageMetadata messag } } - public static ByteBuffer getKeyByteBuffer(PulsarApi.MessageMetadata messageMetadata) { + public static ByteBuffer getKeyByteBuffer(MessageMetadata messageMetadata) { if (messageMetadata.hasOrderingKey()) { - return messageMetadata.getOrderingKey().asReadOnlyByteBuffer(); + return ByteBuffer.wrap(messageMetadata.getOrderingKey()).asReadOnlyBuffer(); } String key = messageMetadata.getPartitionKey(); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/OffsetSearchPredicate.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/OffsetSearchPredicate.java index a367eccebf..03cf301821 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/OffsetSearchPredicate.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/OffsetSearchPredicate.java @@ -14,7 +14,7 @@ package io.streamnative.pulsar.handlers.kop.utils; import org.apache.bookkeeper.mledger.Entry; -import org.apache.pulsar.common.api.proto.PulsarApi; +import org.apache.pulsar.common.api.proto.BrokerEntryMetadata; import org.apache.pulsar.common.protocol.Commands; import org.checkerframework.checker.nullness.qual.Nullable; import org.slf4j.Logger; @@ -34,7 +34,7 @@ public OffsetSearchPredicate(long indexToSearch) { @Override public boolean apply(@Nullable Entry entry) { try { - PulsarApi.BrokerEntryMetadata brokerEntryMetadata = + BrokerEntryMetadata brokerEntryMetadata = Commands.parseBrokerEntryMetadataIfExist(entry.getDataBuffer()); return brokerEntryMetadata.getIndex() < indexToSearch; } catch (Exception e) { diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/OffsetFinderTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/OffsetFinderTest.java index 18d150d8c6..db42f414e5 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/OffsetFinderTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/OffsetFinderTest.java @@ -36,9 +36,8 @@ import org.apache.bookkeeper.test.MockedBookKeeperTestCase; import org.apache.pulsar.broker.service.persistent.PersistentMessageExpiryMonitor; import org.apache.pulsar.common.allocator.PulsarByteBufAllocator; -import org.apache.pulsar.common.api.proto.PulsarApi; +import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.protocol.ByteBufPair; -import org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStream; import org.testng.annotations.Test; /** @@ -47,11 +46,10 @@ public class OffsetFinderTest extends MockedBookKeeperTestCase { public static byte[] createMessageWrittenToLedger(String msg) throws Exception { - PulsarApi.MessageMetadata.Builder messageMetadataBuilder = PulsarApi.MessageMetadata.newBuilder(); - messageMetadataBuilder.setPublishTime(System.currentTimeMillis()); - messageMetadataBuilder.setProducerName("createMessageWrittenToLedger"); - messageMetadataBuilder.setSequenceId(1); - PulsarApi.MessageMetadata messageMetadata = messageMetadataBuilder.build(); + final MessageMetadata messageMetadata = new MessageMetadata(); + messageMetadata.setPublishTime(System.currentTimeMillis()); + messageMetadata.setProducerName("createMessageWrittenToLedger"); + messageMetadata.setSequenceId(1); ByteBuf data = UnpooledByteBufAllocator.DEFAULT.heapBuffer().writeBytes(msg.getBytes()); int msgMetadataSize = messageMetadata.getSerializedSize(); @@ -59,9 +57,8 @@ public static byte[] createMessageWrittenToLedger(String msg) throws Exception { int totalSize = 4 + msgMetadataSize + payloadSize; ByteBuf headers = PulsarByteBufAllocator.DEFAULT.heapBuffer(totalSize, totalSize); - ByteBufCodedOutputStream outStream = ByteBufCodedOutputStream.get(headers); headers.writeInt(msgMetadataSize); - messageMetadata.writeTo(outStream); + messageMetadata.writeTo(headers); ByteBuf headersAndPayload = ByteBufPair.coalesce(ByteBufPair.get(headers, data)); byte[] byteMessage = headersAndPayload.nioBuffer().array(); headersAndPayload.release(); diff --git a/pom.xml b/pom.xml index 3c937d1d58..4e6c6b213f 100644 --- a/pom.xml +++ b/pom.xml @@ -47,7 +47,7 @@ 2.13.3 1.18.4 2.22.0 - 2.8.0-rc-202101052233 + 2.8.0-rc-202101192246 1.7.25 3.1.8 1.12.5 From 1e9601b2417308aca81c5c91b80e1a2e8dee982b Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 21 Jan 2021 20:12:52 +0800 Subject: [PATCH 2/8] Fix PulsarEntryFormatter --- .../kop/format/PulsarEntryFormatter.java | 20 ++++++------------- 1 file changed, 6 insertions(+), 14 deletions(-) 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 c346617916..93952457bb 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 @@ -41,7 +41,6 @@ import org.apache.pulsar.client.impl.MessageImpl; import org.apache.pulsar.client.impl.TypedMessageBuilderImpl; import org.apache.pulsar.common.allocator.PulsarByteBufAllocator; -import org.apache.pulsar.common.api.proto.CompressionType; import org.apache.pulsar.common.api.proto.KeyValue; import org.apache.pulsar.common.api.proto.MarkerType; import org.apache.pulsar.common.api.proto.MessageMetadata; @@ -69,8 +68,6 @@ public ByteBuf encode(final MemoryRecords records, final int numMessages) { int numMessagesInBatch = 0; long sequenceId = -1; - // TODO: handle different compression type - CompressionType compressionType = CompressionType.NONE; ByteBuf batchedMessageMetadataAndPayload = PulsarByteBufAllocator.DEFAULT .buffer(Math.min(INITIAL_BATCH_BUFFER_SIZE, MAX_MESSAGE_BATCH_SIZE_BYTES)); @@ -81,9 +78,6 @@ public ByteBuf encode(final MemoryRecords records, final int numMessages) { StreamSupport.stream(recordBatch.spliterator(), true).forEachOrdered(record -> { MessageImpl message = recordToEntry(record); messages.add(message); - if (msgMetadata.getPublishTime() <= 0) { - msgMetadata.setPublishTime(message.getPublishTime()); - } if (recordBatch.isTransactional()) { msgMetadata.setTxnidMostBits(recordBatch.producerId()); msgMetadata.setTxnidLeastBits(recordBatch.producerEpoch()); @@ -93,6 +87,7 @@ public ByteBuf encode(final MemoryRecords records, final int numMessages) { for (MessageImpl message : messages) { if (++numMessagesInBatch == 1) { + // msgMetadata will set publish time here sequenceId = Commands.initBatchMessageMetadata(msgMetadata, message.getMessageBuilder()); } currentBatchSizeBytes += message.getDataBuffer().readableBytes(); @@ -106,13 +101,6 @@ public ByteBuf encode(final MemoryRecords records, final int numMessages) { message.getDataBuffer(), batchedMessageMetadataAndPayload); } - int uncompressedSize = batchedMessageMetadataAndPayload.readableBytes(); - - if (CompressionType.NONE != compressionType) { - msgMetadata.setCompression(compressionType); - msgMetadata.setUncompressedSize(uncompressedSize); - } - msgMetadata.setNumMessagesInBatch(numMessagesInBatch); ByteBuf buf = Commands.serializeMetadataAndPayload(ChecksumType.Crc32c, @@ -267,9 +255,13 @@ private static MessageImpl recordToEntry(Record record) { builder.value(null); } - // sequence + // Following fields are required in `Commands.initBatchMessageMetadata`, but since we write to + // bookie directly, broker won't make use of them. So here we just set trivial values. + builder.getMetadataBuilder().setProducerName(""); if (record.sequence() >= 0) { builder.sequenceId(record.sequence()); + } else { + builder.sequenceId(0L); } // timestamp From c8248c1f6f5adeda53d663533a4f89c4ee31fa63 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 21 Jan 2021 21:28:27 +0800 Subject: [PATCH 3/8] Fix PulsarEntryFormatter#decode --- .../pulsar/handlers/kop/format/PulsarEntryFormatter.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) 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 93952457bb..544a2cae18 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 @@ -122,10 +122,12 @@ public MemoryRecords decode(final List entries, final byte magic) { ByteBuf metadataAndPayload = entry.getDataBuffer(); // Uncompress the payload if necessary + Commands.skipBrokerEntryMetadataIfExist(metadataAndPayload); MessageMetadata msgMetadata = Commands.parseMessageMetadata(metadataAndPayload); - if (msgMetadata.getMarkerType() == MarkerType.TXN_COMMIT_VALUE - || msgMetadata.getMarkerType() == MarkerType.TXN_ABORT_VALUE) { + if (msgMetadata.hasMarkerType() + && (msgMetadata.getMarkerType() == MarkerType.TXN_COMMIT_VALUE + || msgMetadata.getMarkerType() == MarkerType.TXN_ABORT_VALUE)) { MemoryRecords memoryRecords = MemoryRecords.withEndTransactionMarker( baseOffset, msgMetadata.getPublishTime(), From 06f843892cd4b23483e3498c3342d9f3c9b0e472 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 21 Jan 2021 22:49:36 +0800 Subject: [PATCH 4/8] Fix exception for message without key --- .../handlers/kop/format/PulsarEntryFormatter.java | 1 - .../pulsar/handlers/kop/utils/ByteBufUtils.java | 14 +++++++++----- 2 files changed, 9 insertions(+), 6 deletions(-) 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 544a2cae18..a9dcb76f18 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 @@ -121,7 +121,6 @@ public MemoryRecords decode(final List entries, final byte magic) { // each entry is a batched message ByteBuf metadataAndPayload = entry.getDataBuffer(); - // Uncompress the payload if necessary Commands.skipBrokerEntryMetadataIfExist(metadataAndPayload); MessageMetadata msgMetadata = Commands.parseMessageMetadata(metadataAndPayload); 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 5ea7576ed9..fbc249b64f 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 @@ -33,12 +33,16 @@ public static ByteBuffer getKeyByteBuffer(SingleMessageMetadata messageMetadata) return ByteBuffer.wrap(messageMetadata.getOrderingKey()).asReadOnlyBuffer(); } - String key = messageMetadata.getPartitionKey(); - if (messageMetadata.hasPartitionKeyB64Encoded()) { - return ByteBuffer.wrap(Base64.getDecoder().decode(key)); + if (messageMetadata.hasPartitionKey()) { + final String key = messageMetadata.getPartitionKey(); + if (messageMetadata.hasPartitionKeyB64Encoded()) { + return ByteBuffer.wrap(Base64.getDecoder().decode(key)).asReadOnlyBuffer(); + } else { + // for Base64 not encoded string, convert to UTF_8 chars + return ByteBuffer.wrap(key.getBytes(UTF_8)); + } } else { - // for Base64 not encoded string, convert to UTF_8 chars - return ByteBuffer.wrap(key.getBytes(UTF_8)); + return ByteBuffer.allocate(0); } } From 141c4c5eb8caac4bf545c04bb11837723732b6ea Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 21 Jan 2021 23:22:24 +0800 Subject: [PATCH 5/8] Bump pulsar version to 2.8.0-SNAPSHOT --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index 4e6c6b213f..356bb7b16a 100644 --- a/pom.xml +++ b/pom.xml @@ -47,7 +47,7 @@ 2.13.3 1.18.4 2.22.0 - 2.8.0-rc-202101192246 + 2.8.0-SNAPSHOT 1.7.25 3.1.8 1.12.5 From 9d893b4c94c84f16df445933a47c3357af686479 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Mon, 25 Jan 2021 18:51:47 +0800 Subject: [PATCH 6/8] Fix tests setup failure --- .../pulsar/handlers/kop/KopProtocolHandlerTestBase.java | 2 ++ 1 file changed, 2 insertions(+) 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 3c9065ffae..77008a7ea9 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 @@ -71,6 +71,7 @@ import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.compaction.Compactor; +import org.apache.pulsar.metadata.impl.ZKMetadataStore; import org.apache.pulsar.zookeeper.ZooKeeperClientFactory; import org.apache.pulsar.zookeeper.ZookeeperClientFactoryImpl; import org.apache.zookeeper.CreateMode; @@ -326,6 +327,7 @@ protected void setupBrokerMocks(PulsarService pulsar) throws Exception { // Override default providers with mocked ones doReturn(mockZooKeeperClientFactory).when(pulsar).getZooKeeperClientFactory(); doReturn(mockBookKeeperClientFactory).when(pulsar).newBookKeeperClientFactory(); + doReturn(new ZKMetadataStore(mockZooKeeper)).when(pulsar).createLocalMetadataStore(); Supplier namespaceServiceSupplier = () -> spy(new NamespaceService(pulsar)); doReturn(namespaceServiceSupplier).when(pulsar).getNamespaceServiceProvider(); From 67569892da862957650d3ba091c43b17a24f068d Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 26 Jan 2021 10:26:42 +0800 Subject: [PATCH 7/8] Bump pulsar to rc version --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index 356bb7b16a..7a7b5a30c4 100644 --- a/pom.xml +++ b/pom.xml @@ -47,7 +47,7 @@ 2.13.3 1.18.4 2.22.0 - 2.8.0-SNAPSHOT + 2.8.0-rc-202101252233 1.7.25 3.1.8 1.12.5 From 3ef7fbccdb77260bcb2c319fcc643ee98cdc4592 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 26 Jan 2021 17:37:43 +0800 Subject: [PATCH 8/8] Fix timestamp related tests --- .../handlers/kop/utils/MessageIdUtils.java | 17 ++++++++++++++--- .../pulsar/handlers/kop/utils/OffsetFinder.java | 12 +++++------- .../kop/EntryPublishTimePulsarFormatTest.java | 2 +- 3 files changed, 20 insertions(+), 11 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MessageIdUtils.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MessageIdUtils.java index 379c6ea7c5..a6ada71b05 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MessageIdUtils.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MessageIdUtils.java @@ -13,8 +13,11 @@ */ package io.streamnative.pulsar.handlers.kop.utils; +import io.netty.buffer.ByteBuf; + import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; + import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedLedger; @@ -22,7 +25,7 @@ import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.pulsar.broker.intercept.ManagedLedgerInterceptorImpl; -import org.apache.pulsar.client.impl.MessageImpl; +import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.protocol.Commands; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -45,6 +48,14 @@ public static long getLogEndOffset(ManagedLedger managedLedger) { return getCurrentOffset(managedLedger) + 1; } + public static long getPublishTime(final ByteBuf byteBuf) { + final int readerIndex = byteBuf.readerIndex(); + Commands.skipBrokerEntryMetadataIfExist(byteBuf); + final MessageMetadata metadata = Commands.parseMessageMetadata(byteBuf); + byteBuf.readerIndex(readerIndex); + return metadata.getPublishTime(); + } + public static CompletableFuture getOffsetOfPosition(ManagedLedgerImpl managedLedger, PositionImpl position, boolean needCheckMore, @@ -61,8 +72,8 @@ public void readEntryComplete(Entry entry, Object ctx) { try { if (needCheckMore) { long offset = peekOffsetFromEntry(entry); - MessageImpl msg = MessageImpl.deserialize(entry.getDataBuffer()); - if (msg.getPublishTime() >= timestamp) { + final long publishTime = getPublishTime(entry.getDataBuffer()); + if (publishTime >= timestamp) { future.complete(offset); } else { future.complete(offset + 1); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/OffsetFinder.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/OffsetFinder.java index 0c1ddb17a4..5bc13ebd44 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/OffsetFinder.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/OffsetFinder.java @@ -19,6 +19,7 @@ import com.google.common.base.Predicate; import java.util.Optional; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; + import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.AsyncCallbacks.FindEntryCallback; @@ -29,7 +30,6 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; -import org.apache.pulsar.client.impl.MessageImpl; /** * given a timestamp find the first message (position) (published) at or before the timestamp. @@ -59,17 +59,15 @@ public void findMessages(final long timestamp, AsyncCallbacks.FindEntryCallback } asyncFindNewestMatching(ManagedCursor.FindPositionConstraint.SearchAllAvailableEntries, entry -> { - MessageImpl msg = null; + if (entry == null) { + return false; + } try { - msg = MessageImpl.deserialize(entry.getDataBuffer()); - return msg.getPublishTime() <= timestamp; + return MessageIdUtils.getPublishTime(entry.getDataBuffer()) <= timestamp; } catch (Exception e) { log.error("[{}] Error deserialize message for message position find", managedLedger.getName(), e); } finally { entry.release(); - if (msg != null) { - msg.recycle(); - } } return false; }, this, callback); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimePulsarFormatTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimePulsarFormatTest.java index 6e7de4cc98..dcd251abf3 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimePulsarFormatTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimePulsarFormatTest.java @@ -34,7 +34,7 @@ public class EntryPublishTimePulsarFormatTest extends EntryPublishTimeTest { private static final Logger log = LoggerFactory.getLogger(EntryPublishTimePulsarFormatTest.class); - public EntryPublishTimePulsarFormatTest(String format) { + public EntryPublishTimePulsarFormatTest() { super("pulsar"); }