From 9997ffdd74523b4a2e670bae6efe42f6ddafd678 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Sat, 25 Aug 2018 12:37:49 -0400 Subject: [PATCH] Refactor KafkaAwareTransactionManager Sub-interface of `PlatformTransactionManager` to facilitate Boot auto-configuration. Deprecate `ResourceTransactionManager` inheritance. --- .../ChainedKafkaTransactionManager.java | 14 +++++ .../KafkaAwareTransactionManager.java | 5 +- .../transaction/KafkaTransactionManager.java | 10 +++- .../EnableKafkaIntegrationTests.java | 54 +++++++++++++++++-- 4 files changed, 78 insertions(+), 5 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/transaction/ChainedKafkaTransactionManager.java b/spring-kafka/src/main/java/org/springframework/kafka/transaction/ChainedKafkaTransactionManager.java index 0da7c82f..0b1ad68a 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/transaction/ChainedKafkaTransactionManager.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/transaction/ChainedKafkaTransactionManager.java @@ -55,4 +55,18 @@ public class ChainedKafkaTransactionManager extends ChainedTransactionMana return this.kafkaTransactionManager.getProducerFactory(); } + /** + * Return the producer factory. + * @return the producer factory. + * @deprecated - in a future release {@link KafkaAwareTransactionManager} will not be + * a sub interface of + * {@link org.springframework.transaction.support.ResourceTransactionManager}. + * TODO: Remove in 3.0 + */ + @Deprecated + @Override + public Object getResourceFactory() { + return getProducerFactory(); + } + } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/transaction/KafkaAwareTransactionManager.java b/spring-kafka/src/main/java/org/springframework/kafka/transaction/KafkaAwareTransactionManager.java index be51bb44..1e8be508 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/transaction/KafkaAwareTransactionManager.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/transaction/KafkaAwareTransactionManager.java @@ -17,9 +17,12 @@ package org.springframework.kafka.transaction; import org.springframework.kafka.core.ProducerFactory; +import org.springframework.transaction.support.ResourceTransactionManager; /** * A transaction manager that can provide a {@link ProducerFactory}. + * Currently a sub-interface of {@link ResourceTransactionManager} + * for backwards compatibility. * * @param the key type. * @param the value type. @@ -28,7 +31,7 @@ import org.springframework.kafka.core.ProducerFactory; * @since 2.1.3 * */ -public interface KafkaAwareTransactionManager { +public interface KafkaAwareTransactionManager extends ResourceTransactionManager { /** * Get the producer factory. diff --git a/spring-kafka/src/main/java/org/springframework/kafka/transaction/KafkaTransactionManager.java b/spring-kafka/src/main/java/org/springframework/kafka/transaction/KafkaTransactionManager.java index 64e1c983..830d4ed4 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/transaction/KafkaTransactionManager.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/transaction/KafkaTransactionManager.java @@ -66,7 +66,7 @@ import org.springframework.util.Assert; */ @SuppressWarnings("serial") public class KafkaTransactionManager extends AbstractPlatformTransactionManager - implements ResourceTransactionManager, KafkaAwareTransactionManager { + implements KafkaAwareTransactionManager { private final ProducerFactory producerFactory; @@ -93,6 +93,14 @@ public class KafkaTransactionManager extends AbstractPlatformTransactionMa return this.producerFactory; } + /** + * Return the producer factory. + * @return the producer factory. + * @deprecated - in a future release {@link KafkaAwareTransactionManager} will + * not be a sub interface of {@link ResourceTransactionManager}. + * TODO: Remove in 3.0 + */ + @Deprecated @Override public Object getResourceFactory() { return getProducerFactory(); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java index 3a878327..bee493a3 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java @@ -47,9 +47,11 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mockito; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Primary; import org.springframework.context.event.EventListener; import org.springframework.context.support.PropertySourcesPlaceholderConfigurer; import org.springframework.core.convert.converter.Converter; @@ -88,6 +90,9 @@ import org.springframework.kafka.support.converter.StringJsonMessageConverter; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.test.rule.EmbeddedKafkaRule; import org.springframework.kafka.test.utils.KafkaTestUtils; +import org.springframework.kafka.transaction.ChainedKafkaTransactionManager; +import org.springframework.kafka.transaction.KafkaAwareTransactionManager; +import org.springframework.kafka.transaction.KafkaTransactionManager; import org.springframework.lang.NonNull; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; @@ -173,6 +178,9 @@ public class EnableKafkaIntegrationTests { @Autowired private FooConverter fooConverter; + @Autowired + private ConcurrentKafkaListenerContainerFactory transactionalFactory; + @Test public void testAnonymous() { MessageListenerContainer container = this.registry @@ -669,6 +677,12 @@ public class EnableKafkaIntegrationTests { assertThat(this.listener.consumerRecords.iterator().next().value()).isEqualTo("allRecords"); } + @Test + public void testAutoConfigTm() { + assertThat(this.transactionalFactory.getContainerProperties().getTransactionManager()) + .isInstanceOf(ChainedKafkaTransactionManager.class); + } + @Configuration @EnableKafka @EnableTransactionManagement(proxyTargetClass = true) @@ -686,11 +700,23 @@ public class EnableKafkaIntegrationTests { return Mockito.mock(PlatformTransactionManager.class); } + @Bean + public KafkaTransactionManager ktm() { + return new KafkaTransactionManager<>(txProducerFactory()); + } + + @Bean + @Primary + public ChainedKafkaTransactionManager cktm() { + return new ChainedKafkaTransactionManager<>(ktm(), transactionManager()); + } + private Throwable globalErrorThrowable; @Bean public KafkaListenerContainerFactory> - kafkaListenerContainerFactory() { + kafkaListenerContainerFactory() { + ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); @@ -705,13 +731,28 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - withNoReplyTemplateContainerFactory() { + withNoReplyTemplateContainerFactory() { + ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; } + @Bean + public KafkaListenerContainerFactory> + transactionalFactory(ObjectProvider> tm) { + + ConcurrentKafkaListenerContainerFactory factory = + new ConcurrentKafkaListenerContainerFactory<>(); + factory.setConsumerFactory(consumerFactory()); + KafkaAwareTransactionManager ktm = tm.getIfUnique(); + if (ktm != null) { + factory.getContainerProperties().setTransactionManager(ktm); + } + return factory; + } + @Bean public RecordPassAllFilter recordFilter() { return new RecordPassAllFilter(); @@ -926,6 +967,13 @@ public class EnableKafkaIntegrationTests { return new DefaultKafkaProducerFactory<>(producerConfigs()); } + @Bean + public ProducerFactory txProducerFactory() { + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(producerConfigs()); + pf.setTransactionIdPrefix("tx-"); + return pf; + } + @Bean public Map producerConfigs() { return KafkaTestUtils.producerProps(embeddedKafka); @@ -1481,7 +1529,7 @@ public class EnableKafkaIntegrationTests { } @KafkaListener(id = "ifctx", topics = "annotated9") - @Transactional + @Transactional(transactionManager = "transactionManager") public void listenTx(String foo) { latch2.countDown(); }