Restore PartitionedConsumerTest.
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2015-2019 the original author or authors.
|
||||
* Copyright 2015-2022 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -16,20 +16,21 @@
|
||||
|
||||
package org.springframework.cloud.stream.partitioning;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.mockito.ArgumentCaptor;
|
||||
import org.mockito.ArgumentMatcher;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.cloud.stream.annotation.EnableBinding;
|
||||
import org.springframework.cloud.stream.binder.Binder;
|
||||
import org.springframework.cloud.stream.binder.BinderFactory;
|
||||
import org.springframework.cloud.stream.binder.ConsumerProperties;
|
||||
import org.springframework.cloud.stream.config.BinderFactoryAutoConfiguration;
|
||||
import org.springframework.cloud.stream.messaging.Sink;
|
||||
import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Import;
|
||||
import org.springframework.context.annotation.PropertySource;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
@@ -45,47 +46,39 @@ import static org.mockito.Mockito.verifyNoMoreInteractions;
|
||||
* @author Marius Bogoevici
|
||||
* @author Ilayaperumal Gopinathan
|
||||
* @author Janne Valkealahti
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
// @checkstyle:off
|
||||
@SpringBootTest(classes = PartitionedConsumerTest.TestSink.class, properties = "spring.cloud.stream.default-binder=mock")
|
||||
public class PartitionedConsumerTest {
|
||||
|
||||
// @checkstyle:on
|
||||
@Autowired
|
||||
private BinderFactory binderFactory;
|
||||
|
||||
@Autowired
|
||||
private Sink testSink;
|
||||
|
||||
@Test
|
||||
@SuppressWarnings({ "rawtypes", "unchecked" })
|
||||
public void testBindingPartitionedConsumer() {
|
||||
Binder binder = this.binderFactory.getBinder(null, MessageChannel.class);
|
||||
ArgumentCaptor<ConsumerProperties> argumentCaptor = ArgumentCaptor
|
||||
.forClass(ConsumerProperties.class);
|
||||
verify(binder).bindConsumer(eq("partIn"), isNull(), eq(this.testSink.input()),
|
||||
ArgumentCaptor<DirectWithAttributesChannel> directWithAttributesChannelArgumentCaptor =
|
||||
ArgumentCaptor.forClass(DirectWithAttributesChannel.class);
|
||||
verify(binder).bindConsumer(eq("partIn"), isNull(), directWithAttributesChannelArgumentCaptor.capture(),
|
||||
argumentCaptor.capture());
|
||||
assertThat(directWithAttributesChannelArgumentCaptor.getValue().getBeanName()).isEqualTo("sink-in-0");
|
||||
assertThat(argumentCaptor.getValue().getInstanceIndex()).isEqualTo(0);
|
||||
assertThat(argumentCaptor.getValue().getInstanceCount()).isEqualTo(2);
|
||||
verifyNoMoreInteractions(binder);
|
||||
}
|
||||
|
||||
@EnableBinding(Sink.class)
|
||||
@EnableAutoConfiguration
|
||||
@Import({ BinderFactoryAutoConfiguration.class })
|
||||
@PropertySource("classpath:/org/springframework/cloud/stream/binder/partitioned-consumer-test.properties")
|
||||
public static class TestSink {
|
||||
|
||||
}
|
||||
|
||||
class PropertiesArgumentMatcher implements ArgumentMatcher<ConsumerProperties> {
|
||||
|
||||
@Override
|
||||
public boolean matches(ConsumerProperties argument) {
|
||||
return argument instanceof ConsumerProperties;
|
||||
@Bean
|
||||
public Consumer<String> sink() {
|
||||
return System.out::println;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
spring.cloud.stream.bindings.input.destination=partIn
|
||||
spring.cloud.stream.bindings.input.consumer.partitioned=true
|
||||
spring.cloud.stream.bindings.sink-in-0.destination=partIn
|
||||
spring.cloud.stream.bindings.sink-in-0.consumer.partitioned=true
|
||||
spring.cloud.stream.instanceCount=2
|
||||
spring.cloud.stream.instanceIndex=0
|
||||
|
||||
|
||||
Reference in New Issue
Block a user