From 087247d96fdc09c9f8b2baf49cb2cf9582037c2e Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 2 Apr 2019 10:43:42 -0400 Subject: [PATCH] GH-259: Add pause/resume to KafkaMessageSource Resolves https://github.com/spring-projects/spring-integration-kafka/issues/259 Enable pause/resume on the polled consumer. **cherry-pick to 3.1.x** --- .../kafka/inbound/KafkaMessageSource.java | 54 ++++++++++++++----- .../kafka/inbound/MessageSourceTests.java | 11 +++- 2 files changed, 52 insertions(+), 13 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java index 1c0a92b352..9a3da386bb 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java @@ -47,6 +47,7 @@ import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.acks.AcknowledgmentCallback; import org.springframework.integration.acks.AcknowledgmentCallbackFactory; import org.springframework.integration.endpoint.AbstractMessageSource; +import org.springframework.integration.endpoint.Pausable; import org.springframework.integration.support.AbstractIntegrationMessageBuilder; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; @@ -80,7 +81,7 @@ import org.springframework.util.Assert; * @since 3.0.1 * */ -public class KafkaMessageSource extends AbstractMessageSource implements Lifecycle { +public class KafkaMessageSource extends AbstractMessageSource implements Lifecycle, Pausable { private static final long DEFAULT_POLL_TIMEOUT = 50L; @@ -117,14 +118,20 @@ public class KafkaMessageSource extends AbstractMessageSource impl private Duration commitTimeout; - private volatile Consumer consumer; - private boolean running; private boolean assigned; private Duration assignTimeout = this.minTimeoutProvider.get(); + private volatile Consumer consumer; + + private volatile Collection assignedPartitions = new ArrayList<>(); + + private volatile boolean pausing; + + private volatile boolean paused; + public KafkaMessageSource(ConsumerFactory consumerFactory, String... topics) { this(consumerFactory, new KafkaAckCallbackFactory<>(), topics); } @@ -293,12 +300,33 @@ public class KafkaMessageSource extends AbstractMessageSource impl this.running = false; } + @Override + public synchronized void pause() { + this.pausing = true; + } + + @Override + public synchronized void resume() { + this.pausing = false; + } + @Override protected synchronized Object doReceive() { if (this.consumer == null) { createConsumer(); this.running = true; } + if (this.pausing && !this.paused && this.assignedPartitions.size() > 0) { + this.consumer.pause(this.assignedPartitions); + this.paused = true; + } + else if (this.paused && !this.pausing) { + this.consumer.resume(this.assignedPartitions); + this.paused = false; + } + if (this.paused && this.logger.isDebugEnabled()) { + this.logger.debug("Consumer is paused; no records will be returned"); + } ConsumerRecord record; TopicPartition topicPartition; synchronized (this.consumerMonitor) { @@ -341,6 +369,7 @@ public class KafkaMessageSource extends AbstractMessageSource impl @Override public void onPartitionsRevoked(Collection partitions) { + KafkaMessageSource.this.assignedPartitions.clear(); if (KafkaMessageSource.this.logger.isInfoEnabled()) { KafkaMessageSource.this.logger.info("Partitions revoked: " + partitions); } @@ -351,6 +380,7 @@ public class KafkaMessageSource extends AbstractMessageSource impl @Override public void onPartitionsAssigned(Collection partitions) { + KafkaMessageSource.this.assignedPartitions = new ArrayList<>(partitions); KafkaMessageSource.this.assigned = true; if (KafkaMessageSource.this.logger.isInfoEnabled()) { KafkaMessageSource.this.logger.info("Partitions assigned: " + partitions); @@ -491,7 +521,7 @@ public class KafkaMessageSource extends AbstractMessageSource impl } else { Set> candidates = this.ackInfo.getOffsets().get(this.ackInfo.getTopicPartition()); - KafkaAckInfo ackInfo = null; + KafkaAckInfo ackInformation = null; synchronized (candidates) { if (candidates.iterator().next().equals(this.ackInfo)) { // see if there are any pending acks for higher offsets @@ -507,29 +537,29 @@ public class KafkaMessageSource extends AbstractMessageSource impl } } if (toCommit.size() > 0) { - ackInfo = toCommit.get(toCommit.size() - 1); + ackInformation = toCommit.get(toCommit.size() - 1); if (this.logger.isDebugEnabled()) { this.logger.debug("Committing pending offsets for " + record + " and all deferred to " - + ackInfo.getRecord()); + + ackInformation.getRecord()); } candidates.removeAll(toCommit); } else { - ackInfo = this.ackInfo; + ackInformation = this.ackInfo; } } else { // earlier offsets present this.ackInfo.setAckDeferred(true); } - if (ackInfo != null) { + if (ackInformation != null) { Map offset = - Collections.singletonMap(ackInfo.getTopicPartition(), - new OffsetAndMetadata(ackInfo.getRecord().offset() + 1)); + Collections.singletonMap(ackInformation.getTopicPartition(), + new OffsetAndMetadata(ackInformation.getRecord().offset() + 1)); if (this.commitTimeout == null) { - ackInfo.getConsumer().commitSync(offset); + ackInformation.getConsumer().commitSync(offset); } else { - ackInfo.getConsumer().commitSync(offset, this.commitTimeout); + ackInformation.getConsumer().commitSync(offset, this.commitTimeout); } } else { diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageSourceTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageSourceTests.java index e86350a36f..8b7324a572 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageSourceTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageSourceTests.java @@ -76,9 +76,10 @@ public class MessageSourceTests { public void testAck() { Consumer consumer = mock(Consumer.class); TopicPartition topicPartition = new TopicPartition("foo", 0); + List assigned = Collections.singletonList(topicPartition); willAnswer(i -> { ((ConsumerRebalanceListener) i.getArgument(1)) - .onPartitionsAssigned(Collections.singletonList(topicPartition)); + .onPartitionsAssigned(assigned); return null; }).given(consumer).subscribe(anyCollection(), any(ConsumerRebalanceListener.class)); AtomicReference> paused = new AtomicReference<>(new HashSet<>()); @@ -131,6 +132,10 @@ public class MessageSourceTests { .acknowledge(AcknowledgmentCallback.Status.ACCEPT); received = source.receive(); assertThat(received).isNull(); + source.pause(); + source.receive(); + source.resume(); + source.receive(); source.destroy(); InOrder inOrder = inOrder(consumer); inOrder.verify(consumer).subscribe(anyCollection(), any(ConsumerRebalanceListener.class)); @@ -143,6 +148,10 @@ public class MessageSourceTests { inOrder.verify(consumer).poll(any(Duration.class)); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(4L))); inOrder.verify(consumer).poll(any(Duration.class)); + inOrder.verify(consumer).pause(assigned); + inOrder.verify(consumer).poll(any(Duration.class)); + inOrder.verify(consumer).resume(assigned); + inOrder.verify(consumer).poll(any(Duration.class)); inOrder.verify(consumer).close(); inOrder.verifyNoMoreInteractions(); }