From fe87a4146e3b07d9eaf5836ef4f3e563ff8aba4b Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Thu, 10 Feb 2022 20:29:11 +0800 Subject: [PATCH] Do not override isMarkerMessage --- .../pulsar/handlers/kop/MessagePublishContext.java | 9 --------- .../pulsar/handlers/kop/storage/PartitionLog.java | 1 - 2 files changed, 10 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessagePublishContext.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessagePublishContext.java index db6a5a3f72..f9b54c34a7 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessagePublishContext.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessagePublishContext.java @@ -44,7 +44,6 @@ public final class MessagePublishContext implements PublishContext { private long sequenceId; private long highestSequenceId; private String producerName; - private boolean isControlBatch; @Override public long getSequenceId() { @@ -56,11 +55,6 @@ public long getHighestSequenceId() { return this.highestSequenceId; } - @Override - public boolean isMarkerMessage() { - return this.isControlBatch; - } - private MetadataCorruptedException peekOffsetError; @Override @@ -126,7 +120,6 @@ public static MessagePublishContext get(CompletableFuture offsetFuture, long sequenceId, long highestSequenceId, int numberOfMessages, - boolean isControlBatch, long startTimeNs) { MessagePublishContext callback = RECYCLER.get(); callback.offsetFuture = offsetFuture; @@ -138,7 +131,6 @@ public static MessagePublishContext get(CompletableFuture offsetFuture, callback.sequenceId = sequenceId; callback.highestSequenceId = highestSequenceId; callback.peekOffsetError = null; - callback.isControlBatch = isControlBatch; return callback; } @@ -174,7 +166,6 @@ public void recycle() { sequenceId = -1; highestSequenceId = -1; peekOffsetError = null; - isControlBatch = false; recyclerHandle.recycle(this); } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java index 436a344317..566086ac25 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/storage/PartitionLog.java @@ -340,7 +340,6 @@ private CompletableFuture publishNormalMessage(final PersistentTopic persi appendInfo.firstSequence(), appendInfo.lastSequence(), appendInfo.numMessages(), - appendInfo.isControlBatch(), time.nanoseconds())); return offsetFuture; }