From 17ea7dace7ffe42a4a0b738f8f4058255cecb1f4 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 22 Feb 2018 08:48:58 -0500 Subject: [PATCH] GH-198: Support pause/resume on inbound adapter Resolves https://github.com/spring-projects/spring-integration-kafka/issues/198 You can now pause/resume the consumer. Any records previously fetched will be processed before the pause takes effect. The listener container requires an `idleEventInterval`, a resume will take effect on the next `ListenerContainerIdleEvent`. * Delegate Pausable methods to listener container Requires https://github.com/spring-projects/spring-kafka/pull/584 * Polishing --- .../KafkaMessageDrivenChannelAdapter.java | 21 +++++- .../inbound/MessageDrivenAdapterTests.java | 75 +++++++++++++++++++ 2 files changed, 92 insertions(+), 4 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java index 40fd7c4601..e0ad7493c8 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java @@ -26,6 +26,7 @@ import org.springframework.core.AttributeAccessor; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.context.OrderlyShutdownCapable; import org.springframework.integration.endpoint.MessageProducerSupport; +import org.springframework.integration.endpoint.Pausable; import org.springframework.integration.kafka.support.RawRecordHeaderErrorMessageStrategy; import org.springframework.integration.support.ErrorMessageStrategy; import org.springframework.integration.support.ErrorMessageUtils; @@ -66,7 +67,8 @@ import org.springframework.util.Assert; * @author Artem Bilan * */ -public class KafkaMessageDrivenChannelAdapter extends MessageProducerSupport implements OrderlyShutdownCapable { +public class KafkaMessageDrivenChannelAdapter extends MessageProducerSupport implements OrderlyShutdownCapable, + Pausable { private static final ThreadLocal attributesHolder = new ThreadLocal<>(); @@ -223,6 +225,11 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo this.batchListener.setFallbackType(payloadType); } + @Override + public String getComponentType() { + return "kafka:message-driven-channel-adapter"; + } + @Override protected void onInit() { super.onInit(); @@ -280,8 +287,13 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo } @Override - public String getComponentType() { - return "kafka:message-driven-channel-adapter"; + public void pause() { + this.messageListenerContainer.pause(); + } + + @Override + public void resume() { + this.messageListenerContainer.resume(); } @Override @@ -434,7 +446,8 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo @Override public void onMessage(List> records, Acknowledgment acknowledgment, Consumer consumer) { - Message message = null; + + Message message = null; try { message = toMessagingMessage(records, acknowledgment, consumer); setAttributesIfNecessary(records, message); diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java index 6435c149e6..6129d8a744 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java @@ -17,17 +17,32 @@ package org.springframework.integration.kafka.inbound; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.isNull; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.willAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import java.lang.reflect.Type; import java.util.Arrays; import java.util.Collections; +import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.header.Headers; import org.apache.kafka.common.header.internals.RecordHeaders; import org.junit.ClassRule; @@ -40,16 +55,19 @@ import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAd import org.springframework.integration.kafka.support.RawRecordHeaderErrorMessageStrategy; import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.support.StaticMessageHeaderAccessor; +import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; +import org.springframework.kafka.listener.AbstractMessageListenerContainer.AckMode; import org.springframework.kafka.listener.KafkaMessageListenerContainer; import org.springframework.kafka.listener.config.ContainerProperties; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.DefaultKafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.KafkaNull; +import org.springframework.kafka.support.TopicPartitionInitialOffset; import org.springframework.kafka.support.converter.BatchMessageConverter; import org.springframework.kafka.support.converter.BatchMessagingMessageConverter; import org.springframework.kafka.support.converter.ConversionException; @@ -411,6 +429,63 @@ public class MessageDrivenAdapterTests { adapter.stop(); } + @SuppressWarnings({ "unchecked", "rawtypes" }) + @Test + public void testPauseResume() throws Exception { + ConsumerFactory cf = mock(ConsumerFactory.class); + Consumer consumer = mock(Consumer.class); + given(cf.createConsumer(isNull(), eq("clientId"), isNull())).willReturn(consumer); + final Map>> records = new HashMap<>(); + records.put(new TopicPartition("foo", 0), Arrays.asList( + new ConsumerRecord<>("foo", 0, 0L, 1, "foo"), + new ConsumerRecord<>("foo", 0, 1L, 1, "bar"))); + ConsumerRecords consumerRecords = new ConsumerRecords<>(records); + ConsumerRecords emptyRecords = new ConsumerRecords<>(Collections.emptyMap()); + AtomicBoolean first = new AtomicBoolean(true); + given(consumer.poll(anyLong())).willAnswer(i -> { + Thread.sleep(50); + return first.getAndSet(false) ? consumerRecords : emptyRecords; + }); + final CountDownLatch commitLatch = new CountDownLatch(2); + willAnswer(i -> { + commitLatch.countDown(); + return null; + }).given(consumer).commitSync(any(Map.class)); + given(consumer.assignment()).willReturn(records.keySet()); + final CountDownLatch pauseLatch = new CountDownLatch(1); + willAnswer(i -> { + pauseLatch.countDown(); + return null; + }).given(consumer).pause(records.keySet()); + given(consumer.paused()).willReturn(records.keySet()); + final CountDownLatch resumeLatch = new CountDownLatch(1); + willAnswer(i -> { + resumeLatch.countDown(); + return null; + }).given(consumer).resume(records.keySet()); + TopicPartitionInitialOffset[] topicPartition = new TopicPartitionInitialOffset[] { + new TopicPartitionInitialOffset("foo", 0) }; + ContainerProperties containerProps = new ContainerProperties(topicPartition); + containerProps.setAckMode(AckMode.RECORD); + containerProps.setClientId("clientId"); + containerProps.setIdleEventInterval(100L); + KafkaMessageListenerContainer container = + new KafkaMessageListenerContainer<>(cf, containerProps); + KafkaMessageDrivenChannelAdapter adapter = new KafkaMessageDrivenChannelAdapter(container); + QueueChannel outputChannel = new QueueChannel(); + adapter.setOutputChannel(outputChannel); + adapter.afterPropertiesSet(); + adapter.start(); + assertThat(commitLatch.await(10, TimeUnit.SECONDS)).isTrue(); + verify(consumer, times(2)).commitSync(any(Map.class)); + assertThat(outputChannel.getQueueSize()).isEqualTo(2); + adapter.pause(); + assertThat(pauseLatch.await(10, TimeUnit.SECONDS)).isTrue(); + adapter.resume(); + assertThat(resumeLatch.await(10, TimeUnit.SECONDS)).isTrue(); + adapter.stop(); + } + public static class Foo { private String bar;