s-k 1.0.0.BUILD-SNAPSHOT Compatibility

This commit is contained in:
Gary Russell
2016-06-02 13:17:53 -04:00
committed by Artem Bilan
parent 66e0772da0
commit 490296c6b8
3 changed files with 11 additions and 4 deletions

View File

@@ -46,10 +46,11 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
public KafkaMessageDrivenChannelAdapter(AbstractMessageListenerContainer<K, V> messageListenerContainer) {
Assert.notNull(messageListenerContainer, "messageListenerContainer is required");
Assert.isNull(messageListenerContainer.getMessageListener(), "Container must not already have a listener");
Assert.isNull(messageListenerContainer.getContainerProperties().getMessageListener(),
"Container must not already have a listener");
this.messageListenerContainer = messageListenerContainer;
this.messageListenerContainer.setAutoStartup(false);
this.messageListenerContainer.setMessageListener(this.listener);
this.messageListenerContainer.getContainerProperties().setMessageListener(this.listener);
}
public void setMessageConverter(MessageConverter messageConverter) {

View File

@@ -35,7 +35,11 @@
</constructor-arg>
</bean>
</constructor-arg>
<constructor-arg name="topics" value="foo" />
<constructor-arg>
<bean class="org.springframework.kafka.listener.config.ContainerProperties">
<constructor-arg name="topics" value="foo" />
</bean>
</constructor-arg>
</bean>
</beans>

View File

@@ -33,6 +33,7 @@ import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.listener.KafkaMessageListenerContainer;
import org.springframework.kafka.listener.config.ContainerProperties;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.kafka.support.converter.MessagingMessageConverter;
@@ -59,8 +60,9 @@ public class MessageDrivenAdapterTests {
Map<String, Object> props = KafkaTestUtils.consumerProps("test1", "true", embeddedKafka);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
DefaultKafkaConsumerFactory<Integer, String> cf = new DefaultKafkaConsumerFactory<Integer, String>(props);
ContainerProperties containerProps = new ContainerProperties(topic1);
KafkaMessageListenerContainer<Integer, String> container =
new KafkaMessageListenerContainer<>(cf, topic1);
new KafkaMessageListenerContainer<>(cf, containerProps);
KafkaMessageDrivenChannelAdapter<Integer, String> adapter = new KafkaMessageDrivenChannelAdapter<>(container);
QueueChannel out = new QueueChannel();
adapter.setOutputChannel(out);