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 262b0bea2d..b4e1ad67db 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 @@ -41,6 +41,7 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.WakeupException; import org.springframework.beans.factory.DisposableBean; +import org.springframework.context.Lifecycle; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.endpoint.AbstractMessageSource; import org.springframework.integration.support.AbstractIntegrationMessageBuilder; @@ -75,10 +76,12 @@ import org.springframework.util.Assert; * */ public class KafkaMessageSource extends AbstractMessageSource - implements DisposableBean { + implements DisposableBean, Lifecycle { private static final long DEFAULT_POLL_TIMEOUT = 50L; + private final Log logger = LogFactory.getLog(getClass()); + private final ConsumerFactory consumerFactory; private final KafkaAckCallbackFactory ackCallbackFactory; @@ -105,7 +108,7 @@ public class KafkaMessageSource extends AbstractMessageSource private volatile Consumer consumer; - private volatile Collection partitions; + private volatile boolean running; public KafkaMessageSource(ConsumerFactory consumerFactory, String... topics) { this(consumerFactory, new KafkaAckCallbackFactory<>(), topics); @@ -244,6 +247,26 @@ public class KafkaMessageSource extends AbstractMessageSource } } + @Override + public synchronized boolean isRunning() { + return this.running; + } + + @Override + public synchronized void start() { + this.running = true; + } + + @Override + public synchronized void stop() { + synchronized (this.consumerMonitor) { + if (this.consumer != null) { + this.consumer.close(30, TimeUnit.SECONDS); + } + } + this.running = false; + } + @Override protected synchronized Object doReceive() { if (this.consumer == null) { @@ -252,12 +275,7 @@ public class KafkaMessageSource extends AbstractMessageSource ConsumerRecord record; TopicPartition topicPartition; synchronized (this.consumerMonitor) { - Set paused = this.consumer.paused(); - if (paused.size() > 0) { - this.consumer.resume(paused); - } ConsumerRecords records = this.consumer.poll(this.pollTimeout); - this.consumer.pause(this.partitions); if (records == null || records.count() == 0) { return null; } @@ -295,7 +313,9 @@ public class KafkaMessageSource extends AbstractMessageSource @Override public void onPartitionsRevoked(Collection partitions) { - KafkaMessageSource.this.partitions = Collections.emptyList(); + if (KafkaMessageSource.this.logger.isInfoEnabled()) { + KafkaMessageSource.this.logger.info("Partitions revoked: " + partitions); + } if (KafkaMessageSource.this.rebalanceListener != null) { KafkaMessageSource.this.rebalanceListener.onPartitionsRevoked(partitions); } @@ -303,7 +323,9 @@ public class KafkaMessageSource extends AbstractMessageSource @Override public void onPartitionsAssigned(Collection partitions) { - KafkaMessageSource.this.partitions = new ArrayList<>(partitions); + if (KafkaMessageSource.this.logger.isInfoEnabled()) { + KafkaMessageSource.this.logger.info("Partitions assigned: " + partitions); + } if (KafkaMessageSource.this.rebalanceListener != null) { KafkaMessageSource.this.rebalanceListener.onPartitionsAssigned(partitions); } 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 b819797163..638ec6e61e 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 @@ -136,29 +136,15 @@ public class MessageSourceTests { source.destroy(); InOrder inOrder = inOrder(consumer); inOrder.verify(consumer).subscribe(anyCollection(), any(ConsumerRebalanceListener.class)); - inOrder.verify(consumer).paused(); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(1L))); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(2L))); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(3L))); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(4L))); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).close(30, TimeUnit.SECONDS); inOrder.verifyNoMoreInteractions(); } @@ -212,10 +198,15 @@ public class MessageSourceTests { KafkaMessageSource source = new KafkaMessageSource(consumerFactory, "foo"); Message received1 = source.receive(); + consumer.paused(); // need some other interaction with mock between polls for InOrder Message received2 = source.receive(); + consumer.paused(); // need some other interaction with mock between polls for InOrder Message received3 = source.receive(); + consumer.paused(); // need some other interaction with mock between polls for InOrder Message received4 = source.receive(); + consumer.paused(); // need some other interaction with mock between polls for InOrder Message received5 = source.receive(); + consumer.paused(); // need some other interaction with mock between polls for InOrder Message received6 = source.receive(); StaticMessageHeaderAccessor.getAcknowledgmentCallback(received3) .acknowledge(Status.ACCEPT); @@ -233,27 +224,17 @@ public class MessageSourceTests { source.destroy(); InOrder inOrder = inOrder(consumer); inOrder.verify(consumer).subscribe(anyCollection(), any(ConsumerRebalanceListener.class)); + inOrder.verify(consumer).poll(anyLong()); inOrder.verify(consumer).paused(); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); - inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(3L))); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(6L))); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).close(30, TimeUnit.SECONDS); inOrder.verifyNoMoreInteractions(); } @@ -310,28 +291,15 @@ public class MessageSourceTests { source.destroy(); assertThat(received).isNull(); InOrder inOrder = inOrder(consumer); - inOrder.verify(consumer).paused(); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).seek(topicPartition, 0L); // rollback - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(1L))); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); inOrder.verify(consumer).seek(topicPartition, 1L); // rollback - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(2L))); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).close(30, TimeUnit.SECONDS); inOrder.verifyNoMoreInteractions(); } @@ -354,8 +322,7 @@ public class MessageSourceTests { willAnswer(i -> paused.get()).given(consumer).paused(); Map> records1 = new LinkedHashMap<>(); records1.put(topicPartition, Arrays.asList( - new ConsumerRecord("foo", 0, 0L, 0L, TimestampType.NO_TIMESTAMP_TYPE, 0, 0, 0, null, "foo"), - new ConsumerRecord("foo", 0, 1L, 0L, TimestampType.NO_TIMESTAMP_TYPE, 0, 0, 0, null, "bar"))); + new ConsumerRecord("foo", 0, 0L, 0L, TimestampType.NO_TIMESTAMP_TYPE, 0, 0, 0, null, "foo"))); ConsumerRecords cr1 = new ConsumerRecords(records1); Map> records2 = new LinkedHashMap<>(); records2.put(topicPartition, Collections.singletonList( @@ -370,6 +337,7 @@ public class MessageSourceTests { KafkaMessageSource source = new KafkaMessageSource(consumerFactory, "foo"); Message received1 = source.receive(); + consumer.paused(); // need some other interaction with mock between polls for InOrder Message received2 = source.receive(); // inflight assertThat(received1.getHeaders().get(KafkaHeaders.OFFSET)).isEqualTo(0L); AcknowledgmentCallback ack1 = StaticMessageHeaderAccessor.getAcknowledgmentCallback(received1); @@ -397,12 +365,9 @@ public class MessageSourceTests { source.destroy(); assertThat(received1).isNull(); InOrder inOrder = inOrder(consumer, log1, log2); + inOrder.verify(consumer).poll(anyLong()); inOrder.verify(consumer).paused(); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); // in flight inOrder.verify(consumer).seek(topicPartition, 0L); // rollback inOrder.verify(log1).isWarnEnabled(); ArgumentCaptor captor = ArgumentCaptor.forClass(String.class); @@ -416,20 +381,11 @@ public class MessageSourceTests { assertThat(captor.getValue()) .contains("Cannot commit offset for ConsumerRecord") .contains("; an earlier offset was rolled back"); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(1L))); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(2L))); - inOrder.verify(consumer).paused(); - inOrder.verify(consumer).resume(anyCollection()); inOrder.verify(consumer).poll(anyLong()); - inOrder.verify(consumer).pause(anyCollection()); inOrder.verify(consumer).close(30, TimeUnit.SECONDS); inOrder.verifyNoMoreInteractions(); }