From 36361d19bf9494b4e0f8b637819e1e67b2c4768f Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 6 Sep 2018 12:59:57 -0400 Subject: [PATCH] GH-1459: Pollable Consumer and Requeue Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/1459 Requires https://github.com/spring-cloud/spring-cloud-stream/pull/1467 --- .../stream/binder/kafka/KafkaBinderTests.java | 34 +++++++++++++++++++ 1 file changed, 34 insertions(+) diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 83ca0b920..29d0bc149 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -35,6 +35,7 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import com.fasterxml.jackson.databind.ObjectMapper; + import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.CreateTopicsResult; @@ -75,6 +76,7 @@ import org.springframework.cloud.stream.binder.HeaderMode; import org.springframework.cloud.stream.binder.PartitionCapableBinderTests; import org.springframework.cloud.stream.binder.PartitionTestSupport; import org.springframework.cloud.stream.binder.PollableSource; +import org.springframework.cloud.stream.binder.RequeueCurrentMessageException; import org.springframework.cloud.stream.binder.Spy; import org.springframework.cloud.stream.binder.TestUtils; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; @@ -2533,6 +2535,38 @@ public class KafkaBinderTests extends binding.unbind(); } + @SuppressWarnings({ "rawtypes", "unchecked" }) + @Test + public void testPolledConsumerRequeue() throws Exception { + KafkaTestBinder binder = getBinder(); + PollableSource inboundBindTarget = new DefaultPollableMessageSource(this.messageConverter); + ExtendedConsumerProperties properties = createConsumerProperties(); + Binding> binding = binder.bindPollableConsumer("pollableRequeue", "group", + inboundBindTarget, properties); + Map producerProps = KafkaTestUtils.producerProps(embeddedKafka.getEmbeddedKafka()); + KafkaTemplate template = new KafkaTemplate(new DefaultKafkaProducerFactory<>(producerProps)); + template.send("pollableRequeue", "testPollable"); + try { + boolean polled = false; + int n = 0; + while (n++ < 100 && !polled) { + polled = inboundBindTarget.poll(m -> { + assertThat(m.getPayload()).isEqualTo("testPollable".getBytes()); + throw new RequeueCurrentMessageException(); + }); + } + fail("Expected exception"); + } + catch (MessageHandlingException e) { + assertThat(e.getCause()).isInstanceOf(RequeueCurrentMessageException.class); + } + boolean polled = inboundBindTarget.poll(m -> { + assertThat(m.getPayload()).isEqualTo("testPollable".getBytes()); + }); + assertThat(polled).isTrue(); + binding.unbind(); + } + @SuppressWarnings({ "rawtypes", "unchecked" }) @Test public void testPolledConsumerWithDlq() throws Exception {