From 9e47f6bb5a367f4c54e6c50cf35d10c0911a44af Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Thu, 16 Dec 2021 12:39:02 +0100 Subject: [PATCH 1/7] Fetch: trigger pending Fetches when records are produced --- .../handlers/kop/DelayedCreatePartitions.java | 2 +- .../handlers/kop/DelayedCreateTopics.java | 2 +- .../pulsar/handlers/kop/DelayedFetch.java | 27 +++++++- .../handlers/kop/DelayedProduceAndFetch.java | 2 +- .../handlers/kop/KafkaRequestHandler.java | 20 ++++++ .../pulsar/handlers/kop/KopServerStats.java | 1 + .../handlers/kop/MessageFetchContext.java | 63 +++++++++++++++--- .../pulsar/handlers/kop/RequestStats.java | 8 +++ .../coordinator/group/DelayedHeartbeat.java | 2 +- .../kop/coordinator/group/DelayedJoin.java | 2 +- .../coordinator/group/InitialDelayedJoin.java | 2 +- .../kop/utils/delayed/DelayedOperation.java | 6 +- .../delayed/DelayedOperationPurgatory.java | 6 +- .../utils/delayed/DelayedOperationTest.java | 10 +-- .../pulsar/handlers/kop/KafkaApisTest.java | 65 ++++++++++++++++++- 15 files changed, 187 insertions(+), 31 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreatePartitions.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreatePartitions.java index 847920695e..39e93a5b92 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreatePartitions.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreatePartitions.java @@ -42,7 +42,7 @@ public void onComplete() { } @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { if (numTopics.get() <= 0) { forceComplete(); return true; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreateTopics.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreateTopics.java index 2e41d8ab3a..382dc6698e 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreateTopics.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreateTopics.java @@ -42,7 +42,7 @@ public void onComplete() { } @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { if (numTopics.get() <= 0) { forceComplete(); return true; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java index 105c2644a0..610725d6c2 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java @@ -15,32 +15,53 @@ import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; import java.util.Optional; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; +import lombok.extern.slf4j.Slf4j; +@Slf4j public class DelayedFetch extends DelayedOperation { private final Runnable callback; private final AtomicLong bytesReadable; private final int minBytes; + private final MessageFetchContext messageFetchContext; + private AtomicBoolean restarted = new AtomicBoolean(); - protected DelayedFetch(long delayMs, AtomicLong bytesReadable, int minBytes, Runnable callback) { + protected DelayedFetch(long delayMs, AtomicLong bytesReadable, int minBytes, + MessageFetchContext messageFetchContext) { super(delayMs, Optional.empty()); - this.callback = callback; this.bytesReadable = bytesReadable; this.minBytes = minBytes; + this.messageFetchContext = messageFetchContext; + this.callback = messageFetchContext::complete; } @Override public void onExpiration() { + if (restarted.get()) { + return; + } callback.run(); } @Override public void onComplete() { + if (restarted.get()) { + return; + } callback.run(); } @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { + if (notify) { + // if we are here then we were waiting for the condition + // someone wrote some messages to one of the topics + // trigger the Fetch from scratch + restarted.set(true); + messageFetchContext.onDataWrittenToSomePartition(); + return true; + } if (bytesReadable.get() < minBytes){ return false; } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java index be68b7b4b2..1eabb3c020 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java @@ -42,7 +42,7 @@ public void onComplete() { } @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { if (topicPartitionNum.get() <= 0) { forceComplete(); return true; 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 ed51012167..1705be45c5 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 @@ -56,6 +56,7 @@ import io.streamnative.pulsar.handlers.kop.utils.OffsetFinder; import io.streamnative.pulsar.handlers.kop.utils.TopicNameUtils; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; +import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationKey; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory; import java.net.InetSocketAddress; import java.nio.ByteBuffer; @@ -941,6 +942,11 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, mergedResponse.putAll(unauthorizedTopicResponsesMap); mergedResponse.putAll(invalidRequestResponses); resultFuture.complete(new ProduceResponse(mergedResponse)); + mergedResponse.forEach((_topicPartition, _response) -> { + if (_response.error == Errors.NONE) { + notifyPendingFetches(_topicPartition); + } + }); }); } }; @@ -979,6 +985,20 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, } + private void notifyPendingFetches(TopicPartition topicPartition) { + ctx.executor().execute( () -> { + DelayedOperationKey.TopicPartitionOperationKey key = + new DelayedOperationKey.TopicPartitionOperationKey(topicPartition); + int matches = fetchPurgatory.checkAndComplete(key); + if (matches > 0) { + requestStats.getWaitingFetchesTriggered().add(matches); + if (log.isDebugEnabled()) { + log.debug("{} DelayedFetch woke up for {}", matches, topicPartition); + } + } + }); + }; + private void validateRecords(short version, MemoryRecords records) { if (version >= 3) { Iterator iterator = records.batches().iterator(); 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 65c3cea99f..8e8324d1d6 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 @@ -75,6 +75,7 @@ public interface KopServerStats { String PREPARE_METADATA = "PREPARE_METADATA"; String MESSAGE_READ = "MESSAGE_READ"; String FETCH_DECODE = "FETCH_DECODE"; + String WAITING_FETCHES_TRIGGERED = "WAITING_FETCHES_TRIGGERED"; /** * Consumer 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 a79b836d42..52f73719f7 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 @@ -36,12 +36,14 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.util.MathUtils; @@ -68,6 +70,8 @@ import org.apache.kafka.common.requests.IsolationLevel; import org.apache.kafka.common.requests.RequestHeader; import org.apache.kafka.common.requests.ResponseCallbackWrapper; +import org.apache.kafka.common.utils.SystemTime; +import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.metadata.api.GetResult; /** @@ -83,6 +87,7 @@ protected MessageFetchContext newObject(Handle handle) { }; private final Handle recyclerHandle; + private long startTime; private Map> responseData; private ConcurrentLinkedQueue decodeResults; private KafkaRequestHandler requestHandler; @@ -95,7 +100,7 @@ protected MessageFetchContext newObject(Handle handle) { private RequestHeader header; private volatile CompletableFuture resultFuture; private AtomicBoolean hasComplete; - private AtomicLong bytesReadable; + private AtomicLong bytesRead; private DelayedOperationPurgatory fetchPurgatory; private String namespacePrefix; @@ -119,8 +124,9 @@ public static MessageFetchContext get(KafkaRequestHandler requestHandler, context.header = kafkaHeaderAndRequest.getHeader(); context.resultFuture = resultFuture; context.hasComplete = new AtomicBoolean(false); - context.bytesReadable = new AtomicLong(0); + context.bytesRead = new AtomicLong(0); context.fetchPurgatory = fetchPurgatory; + context.startTime = SystemTime.SYSTEM.hiResClockMs();; return context; } @@ -142,6 +148,7 @@ public static MessageFetchContext getForTest(FetchRequest fetchRequest, context.header = null; context.resultFuture = resultFuture; context.hasComplete = new AtomicBoolean(false); + context.startTime = SystemTime.SYSTEM.hiResClockMs();; return context; } @@ -163,7 +170,7 @@ private void recycle() { header = null; resultFuture = null; hasComplete = null; - bytesReadable = null; + bytesRead = null; fetchPurgatory = null; namespacePrefix = null; recyclerHandle.recycle(this); @@ -195,15 +202,48 @@ private void addErrorPartitionResponse(TopicPartition topicPartition, Errors err private void tryComplete() { if (resultFuture != null && responseData.size() >= fetchRequest.fetchData().size() && hasComplete.compareAndSet(false, true)) { - DelayedFetch delayedFetch = new DelayedFetch(fetchRequest.maxWait(), bytesReadable, - fetchRequest.minBytes(), this::complete); - List delayedFetchKeys = - fetchRequest.fetchData().keySet().stream() - .map(DelayedOperationKey.TopicPartitionOperationKey::new).collect(Collectors.toList()); - fetchPurgatory.tryCompleteElseWatch(delayedFetch, delayedFetchKeys); + boolean errorsOccurred = false; + if (responseData + .values() + .stream() + .anyMatch(p->p.error != Errors.NONE)) { + // if there is an error no need to wait, the fetch must fail + // as soon as possible + errorsOccurred = true; + } + long now = SystemTime.SYSTEM.hiResClockMs(); + long currentWait = now - this.startTime; + long remainingMaxWait = fetchRequest.maxWait() - currentWait; + long maxWait = Math.min(remainingMaxWait, fetchRequest.maxWait()); + if (bytesRead.get() < fetchRequest.minBytes() && !errorsOccurred && maxWait > 0) { + // we haven't read enough data, need to wait + DelayedFetch delayedFetch = new DelayedFetch(maxWait, bytesRead, + fetchRequest.minBytes(), this); + List delayedFetchKeys = + fetchRequest.fetchData().keySet().stream() + .map(DelayedOperationKey.TopicPartitionOperationKey::new).collect(Collectors.toList()); + fetchPurgatory.tryCompleteElseWatch(delayedFetch, delayedFetchKeys); + } else { + this.complete(); + } } } + /** + * Restart this Fetch, we were waiting for some data (minBytes) + * and someone wrote something on any of the watched partitions. + */ + public void onDataWrittenToSomePartition() { + decodeResults.forEach(DecodeResult::recycle); + decodeResults.clear(); + bytesRead.set(0); + hasComplete.set(false); + log.info("onDataWrittenToSomePartition current resposnses {}", responseData); + responseData.clear(); + handleFetch(); + } + + public void complete() { if (resultFuture == null) { // the context has been recycled @@ -329,6 +369,9 @@ private void handlePartitionData(final TopicPartition topicPartition, // the future that is returned by getTopicConsumerManager is always completed normally topicManager.getTopicConsumerManager(fullTopicName).thenAccept(tcm -> { if (tcm == null) { + if (log.isDebugEnabled()) { + log.debug("Fetch for {}: failed, topic not owned .", topicPartition); + } registerPrepareMetadataFailedEvent(startPrepareMetadataNanos); // remove null future cache KafkaTopicConsumerManagerCache.getInstance().removeAndCloseByTopic(fullTopicName); @@ -491,7 +534,7 @@ private void handleEntries(final List entries, highWatermark, // TODO: should it be changed to the logStartOffset? abortedTransactions, kafkaRecords)); - bytesReadable.getAndAdd(kafkaRecords.sizeInBytes()); + bytesRead.getAndAdd(kafkaRecords.sizeInBytes()); tryComplete(); }); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/RequestStats.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/RequestStats.java index cff8d88a1c..99e79faa48 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/RequestStats.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/RequestStats.java @@ -28,6 +28,7 @@ import static io.streamnative.pulsar.handlers.kop.KopServerStats.RESPONSE_BLOCKED_LATENCY; import static io.streamnative.pulsar.handlers.kop.KopServerStats.RESPONSE_BLOCKED_TIMES; import static io.streamnative.pulsar.handlers.kop.KopServerStats.SERVER_SCOPE; +import static io.streamnative.pulsar.handlers.kop.KopServerStats.WAITING_FETCHES_TRIGGERED; import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; import java.util.concurrent.atomic.AtomicInteger; @@ -112,6 +113,12 @@ public class RequestStats { ) private final OpStatsLogger fetchDecodeStats; + @StatsDoc( + name = WAITING_FETCHES_TRIGGERED, + help = "number of pending fetches that woke up due to some data produced" + ) + private final Counter waitingFetchesTriggered; + public RequestStats(StatsLogger statsLogger) { this.statsLogger = statsLogger; @@ -127,6 +134,7 @@ public RequestStats(StatsLogger statsLogger) { this.prepareMetadataStats = statsLogger.getOpStatsLogger(PREPARE_METADATA); this.messageReadStats = statsLogger.getOpStatsLogger(MESSAGE_READ); this.fetchDecodeStats = statsLogger.getOpStatsLogger(FETCH_DECODE); + this.waitingFetchesTriggered = statsLogger.getCounter(WAITING_FETCHES_TRIGGERED); statsLogger.registerGauge(REQUEST_QUEUE_SIZE, new Gauge() { @Override diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedHeartbeat.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedHeartbeat.java index 59ea23207c..442c6cd688 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedHeartbeat.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedHeartbeat.java @@ -51,7 +51,7 @@ public void onComplete() { } @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { return coordinator.tryCompleteHeartbeat(group, member, heartbeatDeadline, () -> forceComplete()); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedJoin.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedJoin.java index 94c08113bf..0ef343ffb9 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedJoin.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedJoin.java @@ -50,7 +50,7 @@ public void onComplete() { } @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { return coordinator.tryCompleteJoin(group, () -> forceComplete()); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/InitialDelayedJoin.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/InitialDelayedJoin.java index 8b6070562c..5e224dc1ee 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/InitialDelayedJoin.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/InitialDelayedJoin.java @@ -74,7 +74,7 @@ public void onComplete() { } @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { return false; } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java index b03234cc26..6e9dad2a36 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java @@ -96,7 +96,7 @@ public boolean isCompleted() { * *

This function needs to be defined in subclasses. */ - public abstract boolean tryComplete(); + public abstract boolean tryComplete(boolean notify); /** * Thread-safe variant of tryComplete() that attempts completion only if the lock can be acquired @@ -110,14 +110,14 @@ public boolean isCompleted() { * every invocation of `maybeTryComplete` is followed by at least one invocation of `tryComplete` until * the operation is actually completed. */ - boolean maybeTryComplete() { + boolean maybeTryComplete(boolean notify) { boolean retry = false; boolean done = false; do { if (lock.tryLock()) { try { tryCompletePending.set(false); - done = tryComplete(); + done = tryComplete(notify); } finally { lock.unlock(); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java index fafa361754..71c5d2df06 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java @@ -175,7 +175,7 @@ public boolean tryCompleteElseWatch(T operation, List watchKeys) { // At this point the only thread that can attempt this operation is this current thread // Hence it is safe to tryComplete() without a lock - boolean isCompletedByMe = operation.tryComplete(); + boolean isCompletedByMe = operation.tryComplete(false); if (isCompletedByMe) { return true; } @@ -194,7 +194,7 @@ public boolean tryCompleteElseWatch(T operation, List watchKeys) { } } - isCompletedByMe = operation.maybeTryComplete(); + isCompletedByMe = operation.maybeTryComplete(false); if (isCompletedByMe) { return true; } @@ -349,7 +349,7 @@ public int tryCompleteWatched() { if (curr.isCompleted()) { // another thread has completed this operation, just remove it iter.remove(); - } else if (curr.maybeTryComplete()) { + } else if (curr.maybeTryComplete(true)) { iter.remove(); completed += 1; } diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationTest.java index 88f3d5de61..5c07a71e33 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationTest.java @@ -92,7 +92,7 @@ public void onComplete() { } @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { if (completable) { return forceComplete(); } else { @@ -120,7 +120,7 @@ static class TestDelayOperation extends MockDelayedOperation { @SneakyThrows @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { boolean shouldComplete = completable; Thread.sleep(ThreadLocalRandom.current().nextInt(maxDelayMs)); if (shouldComplete) { @@ -225,13 +225,13 @@ public void testRequestPurge() { // complete the operations, it should immediately be purged from the delayed operation r2.completable = true; - r2.tryComplete(); + r2.tryComplete(false); assertEquals( "Purgatory should have 2 total delayed operations instead of " + purgatory.delayed(), 2, purgatory.delayed()); r3.completable = true; - r3.tryComplete(); + r3.tryComplete(false); assertEquals( "Purgatory should have 1 total delayed operations instead of " + purgatory.delayed(), 1, purgatory.delayed()); @@ -285,7 +285,7 @@ public void testTryCompleteLockContention() throws Exception { MockDelayedOperation op = new MockDelayedOperation(100000L) { @SneakyThrows @Override - public boolean tryComplete() { + public boolean tryComplete(boolean notify) { boolean shouldComplete = completionAttemptsRemaining.decrementAndGet() <= 0; tryCompleteSemaphore.acquire(); try { diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java index 530904a429..b80181eb0d 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java @@ -30,9 +30,12 @@ import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; import io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.KafkaHeaderAndRequest; + +import java.io.InputStream; import java.net.InetSocketAddress; import java.net.SocketAddress; import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.ArrayList; import java.util.Collection; @@ -46,7 +49,12 @@ import java.util.stream.Collectors; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.io.IOUtils; import org.apache.commons.lang3.tuple.Pair; +import org.apache.http.HttpResponse; +import org.apache.http.client.HttpClient; +import org.apache.http.client.methods.HttpGet; +import org.apache.http.impl.client.HttpClientBuilder; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -380,17 +388,21 @@ public void testFetchMinBytes() throws Exception { int maxWaitMs = 3000; int minBytes = 1; // case1: consuming an empty topic. + @Cleanup KafkaConsumer consumer1 = createKafkaConsumer(maxWaitMs, minBytes); consumer1.assign(topicPartitions); Long startTime1 = System.currentTimeMillis(); - consumer1.poll(Duration.ofMillis(maxWaitMs)); + ConsumerRecords emptyResult = consumer1.poll(Duration.ofMillis(maxWaitMs)); Long endTime1 = System.currentTimeMillis(); log.info("cost time1:" + (endTime1 - startTime1)); + assertEquals(0, emptyResult.count()); // case2: consuming an topic after producing data. + @Cleanup KafkaProducer kProducer = createKafkaProducer(); produceData(kProducer, topicPartitions, 10); + @Cleanup KafkaConsumer consumer2 = createKafkaConsumer(maxWaitMs, minBytes); consumer2.assign(topicPartitions); consumer2.seekToBeginning(topicPartitions); @@ -407,6 +419,57 @@ public void testFetchMinBytes() throws Exception { assertTrue(endTime2 - startTime2 < maxWaitMs); } + /** + * Test the sending speed of fetch request when the readable data is less than fetch.minBytes. + */ + @Test(timeOut = 60000) + public void testFetchMinBytesSingleConsumer() throws Exception { + String topicName = "testMinBytesTopic"; + TopicPartition tp = new TopicPartition(topicName, 0); + + // create partitioned topic. + admin.topics().createPartitionedTopic(topicName, 1); + List topicPartitions = new ArrayList<>(); + topicPartitions.add(tp); + + int maxWaitMs = 3000; // very long time + int minBytes = 1; + // case1: consuming an empty topic. + @Cleanup + KafkaConsumer consumer1 = createKafkaConsumer(maxWaitMs, minBytes); + consumer1.assign(topicPartitions); + ConsumerRecords emptyResult = consumer1.poll(Duration.ofMillis(maxWaitMs)); + assertEquals(0, emptyResult.count()); + + // case2: consuming an topic after producing data. + @Cleanup + KafkaProducer kProducer = createKafkaProducer(); + produceData(kProducer, topicPartitions, 10); + + int totalRead = 0; + do { + // Consumer1 is able to eventually read the data + // please note that we are passing 100 as max pool time + ConsumerRecords goodResultFrom1 = consumer1.poll(Duration.ofMillis(100)); + totalRead += goodResultFrom1.count(); + log.info("read {} records totalRead {}", goodResultFrom1.count(), totalRead); + // we require that every pool returns at least one record + // that is that we NEVER hit the maxWait timeout and also the pool timeout + assertTrue(goodResultFrom1.count() > 0); + } while (totalRead < 10); + assertEquals(10, totalRead); + + + HttpClient httpClient = HttpClientBuilder.create().build(); + final String metricsEndPoint = pulsar.getWebServiceAddress() + "/metrics"; + HttpResponse response = httpClient.execute(new HttpGet(metricsEndPoint)); + InputStream inputStream = response.getEntity().getContent(); + String metrics = IOUtils.toString(inputStream, StandardCharsets.UTF_8); + log.info("metrics {}", metrics); + + assertTrue(metrics.contains("kop_server_WAITING_FETCHES_TRIGGERED 1")); + } + @Test(timeOut = 80000) public void testConsumerListOffset() throws Exception { String topicName = "listOffset"; From 6f361392c634e2717bf18b559ef9854dbbb70cb3 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Thu, 16 Dec 2021 12:57:12 +0100 Subject: [PATCH 2/7] fix checkstyle --- .../pulsar/handlers/kop/KafkaRequestHandler.java | 2 +- .../pulsar/handlers/kop/MessageFetchContext.java | 7 ++----- .../io/streamnative/pulsar/handlers/kop/KafkaApisTest.java | 1 - 3 files changed, 3 insertions(+), 7 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 1705be45c5..a87daefd11 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 @@ -986,7 +986,7 @@ protected void handleProduceRequest(KafkaHeaderAndRequest produceHar, } private void notifyPendingFetches(TopicPartition topicPartition) { - ctx.executor().execute( () -> { + ctx.executor().execute(() -> { DelayedOperationKey.TopicPartitionOperationKey key = new DelayedOperationKey.TopicPartitionOperationKey(topicPartition); int matches = fetchPurgatory.checkAndComplete(key); 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 52f73719f7..3dabcb64bb 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 @@ -36,14 +36,12 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; -import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; -import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.util.MathUtils; @@ -71,7 +69,6 @@ import org.apache.kafka.common.requests.RequestHeader; import org.apache.kafka.common.requests.ResponseCallbackWrapper; import org.apache.kafka.common.utils.SystemTime; -import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.metadata.api.GetResult; /** @@ -126,7 +123,7 @@ public static MessageFetchContext get(KafkaRequestHandler requestHandler, context.hasComplete = new AtomicBoolean(false); context.bytesRead = new AtomicLong(0); context.fetchPurgatory = fetchPurgatory; - context.startTime = SystemTime.SYSTEM.hiResClockMs();; + context.startTime = SystemTime.SYSTEM.hiResClockMs(); return context; } @@ -148,7 +145,7 @@ public static MessageFetchContext getForTest(FetchRequest fetchRequest, context.header = null; context.resultFuture = resultFuture; context.hasComplete = new AtomicBoolean(false); - context.startTime = SystemTime.SYSTEM.hiResClockMs();; + context.startTime = SystemTime.SYSTEM.hiResClockMs(); return context; } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java index b80181eb0d..2887bb16f5 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java @@ -30,7 +30,6 @@ import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; import io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.KafkaHeaderAndRequest; - import java.io.InputStream; import java.net.InetSocketAddress; import java.net.SocketAddress; From 203f5ca051d0f7f3e9fe4cc43f4e6fb1838c9b84 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 20 Dec 2021 09:23:25 +0100 Subject: [PATCH 3/7] Add docs --- docs/reference-metrics.md | 1 + 1 file changed, 1 insertion(+) diff --git a/docs/reference-metrics.md b/docs/reference-metrics.md index 8a83a0b93e..b19cfd6188 100644 --- a/docs/reference-metrics.md +++ b/docs/reference-metrics.md @@ -72,6 +72,7 @@ The KoP metrics are exposed under "/metrics" at port `8000` along with Pulsar me | 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_server_WAITING_FETCHES_TRIGGERED | Counter | Number of fetches that have been delayed due to not enough data, and that have been unblocked because some message has been produced| ### Kop event metrics From 605a342c8b942efdddea44341edcad186fe73127 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 20 Dec 2021 09:23:44 +0100 Subject: [PATCH 4/7] Move tryComplete(boolean) to wakeup() --- .../handlers/kop/DelayedCreatePartitions.java | 2 +- .../pulsar/handlers/kop/DelayedCreateTopics.java | 2 +- .../pulsar/handlers/kop/DelayedFetch.java | 14 ++++++++++++-- .../handlers/kop/DelayedProduceAndFetch.java | 2 +- .../pulsar/handlers/kop/KafkaRequestHandler.java | 2 +- .../kop/coordinator/group/DelayedHeartbeat.java | 2 +- .../kop/coordinator/group/DelayedJoin.java | 2 +- .../kop/coordinator/group/InitialDelayedJoin.java | 2 +- .../kop/utils/delayed/DelayedOperation.java | 13 ++++++++++--- .../utils/delayed/DelayedOperationPurgatory.java | 10 +++++----- .../kop/utils/delayed/DelayedOperationTest.java | 10 +++++----- 11 files changed, 39 insertions(+), 22 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreatePartitions.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreatePartitions.java index 39e93a5b92..847920695e 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreatePartitions.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreatePartitions.java @@ -42,7 +42,7 @@ public void onComplete() { } @Override - public boolean tryComplete(boolean notify) { + public boolean tryComplete() { if (numTopics.get() <= 0) { forceComplete(); return true; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreateTopics.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreateTopics.java index 382dc6698e..2e41d8ab3a 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreateTopics.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedCreateTopics.java @@ -42,7 +42,7 @@ public void onComplete() { } @Override - public boolean tryComplete(boolean notify) { + public boolean tryComplete() { if (numTopics.get() <= 0) { forceComplete(); return true; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java index 610725d6c2..4971ea7787 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java @@ -26,6 +26,7 @@ public class DelayedFetch extends DelayedOperation { private final int minBytes; private final MessageFetchContext messageFetchContext; private AtomicBoolean restarted = new AtomicBoolean(); + private final AtomicBoolean someMessageProduced = new AtomicBoolean(); protected DelayedFetch(long delayMs, AtomicLong bytesReadable, int minBytes, MessageFetchContext messageFetchContext) { @@ -53,8 +54,8 @@ public void onComplete() { } @Override - public boolean tryComplete(boolean notify) { - if (notify) { + public boolean tryComplete() { + if (someMessageProduced.get()) { // if we are here then we were waiting for the condition // someone wrote some messages to one of the topics // trigger the Fetch from scratch @@ -68,4 +69,13 @@ public boolean tryComplete(boolean notify) { callback.run(); return true; } + + @Override + public boolean wakeup(Object topicPartition) { + // In the future we could notify the MessageFetchContext that the + // new data is only on this partition and not + // on other partitions + someMessageProduced.set(true); + return true; + } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java index 1eabb3c020..be68b7b4b2 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedProduceAndFetch.java @@ -42,7 +42,7 @@ public void onComplete() { } @Override - public boolean tryComplete(boolean notify) { + public boolean tryComplete() { if (topicPartitionNum.get() <= 0) { forceComplete(); return true; 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 a87daefd11..4b46613d50 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 @@ -997,7 +997,7 @@ private void notifyPendingFetches(TopicPartition topicPartition) { } } }); - }; + } private void validateRecords(short version, MemoryRecords records) { if (version >= 3) { diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedHeartbeat.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedHeartbeat.java index 442c6cd688..59ea23207c 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedHeartbeat.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedHeartbeat.java @@ -51,7 +51,7 @@ public void onComplete() { } @Override - public boolean tryComplete(boolean notify) { + public boolean tryComplete() { return coordinator.tryCompleteHeartbeat(group, member, heartbeatDeadline, () -> forceComplete()); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedJoin.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedJoin.java index 0ef343ffb9..94c08113bf 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedJoin.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/DelayedJoin.java @@ -50,7 +50,7 @@ public void onComplete() { } @Override - public boolean tryComplete(boolean notify) { + public boolean tryComplete() { return coordinator.tryCompleteJoin(group, () -> forceComplete()); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/InitialDelayedJoin.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/InitialDelayedJoin.java index 5e224dc1ee..8b6070562c 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/InitialDelayedJoin.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/group/InitialDelayedJoin.java @@ -74,7 +74,7 @@ public void onComplete() { } @Override - public boolean tryComplete(boolean notify) { + public boolean tryComplete() { return false; } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java index 6e9dad2a36..43b5423ccc 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java @@ -96,7 +96,14 @@ public boolean isCompleted() { * *

This function needs to be defined in subclasses. */ - public abstract boolean tryComplete(boolean notify); + public abstract boolean tryComplete(); + + /** + * Try to wake up the operation. + */ + public boolean wakeup(Object eventKey) { + return true; + } /** * Thread-safe variant of tryComplete() that attempts completion only if the lock can be acquired @@ -110,14 +117,14 @@ public boolean isCompleted() { * every invocation of `maybeTryComplete` is followed by at least one invocation of `tryComplete` until * the operation is actually completed. */ - boolean maybeTryComplete(boolean notify) { + boolean maybeTryComplete() { boolean retry = false; boolean done = false; do { if (lock.tryLock()) { try { tryCompletePending.set(false); - done = tryComplete(notify); + done = tryComplete(); } finally { lock.unlock(); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java index 71c5d2df06..5276991917 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java @@ -175,7 +175,7 @@ public boolean tryCompleteElseWatch(T operation, List watchKeys) { // At this point the only thread that can attempt this operation is this current thread // Hence it is safe to tryComplete() without a lock - boolean isCompletedByMe = operation.tryComplete(false); + boolean isCompletedByMe = operation.tryComplete(); if (isCompletedByMe) { return true; } @@ -194,7 +194,7 @@ public boolean tryCompleteElseWatch(T operation, List watchKeys) { } } - isCompletedByMe = operation.maybeTryComplete(false); + isCompletedByMe = operation.maybeTryComplete(); if (isCompletedByMe) { return true; } @@ -226,7 +226,7 @@ public int checkAndComplete(Object key) { if (null == watchers) { return 0; } else { - return watchers.tryCompleteWatched(); + return watchers.tryCompleteWatched(key); } } @@ -340,7 +340,7 @@ public void watch(T t) { } // traverse the list and try to complete some watched elements - public int tryCompleteWatched() { + public int tryCompleteWatched(Object key) { int completed = 0; Iterator iter = operations.iterator(); @@ -349,7 +349,7 @@ public int tryCompleteWatched() { if (curr.isCompleted()) { // another thread has completed this operation, just remove it iter.remove(); - } else if (curr.maybeTryComplete(true)) { + } else if (curr.wakeup(key) && curr.maybeTryComplete()) { iter.remove(); completed += 1; } diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationTest.java index 5c07a71e33..88f3d5de61 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationTest.java @@ -92,7 +92,7 @@ public void onComplete() { } @Override - public boolean tryComplete(boolean notify) { + public boolean tryComplete() { if (completable) { return forceComplete(); } else { @@ -120,7 +120,7 @@ static class TestDelayOperation extends MockDelayedOperation { @SneakyThrows @Override - public boolean tryComplete(boolean notify) { + public boolean tryComplete() { boolean shouldComplete = completable; Thread.sleep(ThreadLocalRandom.current().nextInt(maxDelayMs)); if (shouldComplete) { @@ -225,13 +225,13 @@ public void testRequestPurge() { // complete the operations, it should immediately be purged from the delayed operation r2.completable = true; - r2.tryComplete(false); + r2.tryComplete(); assertEquals( "Purgatory should have 2 total delayed operations instead of " + purgatory.delayed(), 2, purgatory.delayed()); r3.completable = true; - r3.tryComplete(false); + r3.tryComplete(); assertEquals( "Purgatory should have 1 total delayed operations instead of " + purgatory.delayed(), 1, purgatory.delayed()); @@ -285,7 +285,7 @@ public void testTryCompleteLockContention() throws Exception { MockDelayedOperation op = new MockDelayedOperation(100000L) { @SneakyThrows @Override - public boolean tryComplete(boolean notify) { + public boolean tryComplete() { boolean shouldComplete = completionAttemptsRemaining.decrementAndGet() <= 0; tryCompleteSemaphore.acquire(); try { From 464c2602be5893c674f5d934d91fe8afbe29a7fb Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 21 Dec 2021 10:34:28 +0100 Subject: [PATCH 5/7] Update kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java Co-authored-by: Yunze Xu --- .../io/streamnative/pulsar/handlers/kop/MessageFetchContext.java | 1 - 1 file changed, 1 deletion(-) 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 3dabcb64bb..d1f974c6a3 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 @@ -235,7 +235,6 @@ public void onDataWrittenToSomePartition() { decodeResults.clear(); bytesRead.set(0); hasComplete.set(false); - log.info("onDataWrittenToSomePartition current resposnses {}", responseData); responseData.clear(); handleFetch(); } From ec48967ebc46f311355545718f0bbe16b2aa6872 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 21 Dec 2021 10:36:38 +0100 Subject: [PATCH 6/7] Update kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java Co-authored-by: Yunze Xu --- .../java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java index 4971ea7787..e1c1a2632f 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java @@ -25,7 +25,7 @@ public class DelayedFetch extends DelayedOperation { private final AtomicLong bytesReadable; private final int minBytes; private final MessageFetchContext messageFetchContext; - private AtomicBoolean restarted = new AtomicBoolean(); + private final AtomicBoolean restarted = new AtomicBoolean(); private final AtomicBoolean someMessageProduced = new AtomicBoolean(); protected DelayedFetch(long delayMs, AtomicLong bytesReadable, int minBytes, From 4aae40a51e340f15ff3d9c0bf8ee4ec0d51e2c22 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 21 Dec 2021 12:35:32 +0100 Subject: [PATCH 7/7] Address comments --- .../java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java | 2 +- .../pulsar/handlers/kop/utils/delayed/DelayedOperation.java | 2 +- .../handlers/kop/utils/delayed/DelayedOperationPurgatory.java | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java index e1c1a2632f..811b36de2b 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/DelayedFetch.java @@ -71,7 +71,7 @@ public boolean tryComplete() { } @Override - public boolean wakeup(Object topicPartition) { + public boolean wakeup() { // In the future we could notify the MessageFetchContext that the // new data is only on this partition and not // on other partitions diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java index 43b5423ccc..6ff448a26d 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperation.java @@ -101,7 +101,7 @@ public boolean isCompleted() { /** * Try to wake up the operation. */ - public boolean wakeup(Object eventKey) { + public boolean wakeup() { return true; } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java index 5276991917..e1d6ef5a0c 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/delayed/DelayedOperationPurgatory.java @@ -349,7 +349,7 @@ public int tryCompleteWatched(Object key) { if (curr.isCompleted()) { // another thread has completed this operation, just remove it iter.remove(); - } else if (curr.wakeup(key) && curr.maybeTryComplete()) { + } else if (curr.wakeup() && curr.maybeTryComplete()) { iter.remove(); completed += 1; }