diff --git a/samples/sample-01/pom.xml b/samples/sample-01/pom.xml index 90953b89..d3abe008 100644 --- a/samples/sample-01/pom.xml +++ b/samples/sample-01/pom.xml @@ -5,7 +5,7 @@ com.example kafka-sample-01 - 0.0.1-SNAPSHOT + 2.3.1.BUILD-SNAPSHOT jar kafka-sample-01 @@ -14,7 +14,7 @@ org.springframework.boot spring-boot-starter-parent - 2.1.0.RELEASE + 2.2.0.BUILD-SNAPSHOT diff --git a/samples/sample-01/src/main/java/com/example/Application.java b/samples/sample-01/src/main/java/com/example/Application.java index 4c81f6e5..5219c152 100644 --- a/samples/sample-01/src/main/java/com/example/Application.java +++ b/samples/sample-01/src/main/java/com/example/Application.java @@ -32,6 +32,7 @@ import org.springframework.kafka.listener.DeadLetterPublishingRecoverer; import org.springframework.kafka.listener.SeekToCurrentErrorHandler; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.kafka.support.converter.StringJsonMessageConverter; +import org.springframework.util.backoff.FixedBackOff; import com.common.Foo2; @@ -58,7 +59,7 @@ public class Application { ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); configurer.configure(factory, kafkaConsumerFactory); factory.setErrorHandler(new SeekToCurrentErrorHandler( - new DeadLetterPublishingRecoverer(template), 3)); // dead-letter after 3 tries + new DeadLetterPublishingRecoverer(template), new FixedBackOff(0L, 2))); // dead-letter after 3 tries return factory; } diff --git a/samples/sample-02/pom.xml b/samples/sample-02/pom.xml index 4ae2b989..c4994274 100644 --- a/samples/sample-02/pom.xml +++ b/samples/sample-02/pom.xml @@ -5,7 +5,7 @@ com.example kafka-sample-02 - 0.0.1-SNAPSHOT + 2.3.1.BUILD-SNAPSHOT jar kafka-sample-02 @@ -14,7 +14,7 @@ org.springframework.boot spring-boot-starter-parent - 2.1.0.RELEASE + 2.2.0.BUILD-SNAPSHOT diff --git a/samples/sample-02/src/main/java/com/example/Application.java b/samples/sample-02/src/main/java/com/example/Application.java index c65050e6..33ae67cf 100644 --- a/samples/sample-02/src/main/java/com/example/Application.java +++ b/samples/sample-02/src/main/java/com/example/Application.java @@ -34,6 +34,7 @@ import org.springframework.kafka.support.converter.DefaultJackson2JavaTypeMapper import org.springframework.kafka.support.converter.Jackson2JavaTypeMapper.TypePrecedence; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.kafka.support.converter.StringJsonMessageConverter; +import org.springframework.util.backoff.FixedBackOff; import com.common.Bar2; import com.common.Foo2; @@ -53,7 +54,7 @@ public class Application { ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); configurer.configure(factory, kafkaConsumerFactory); factory.setErrorHandler(new SeekToCurrentErrorHandler( - new DeadLetterPublishingRecoverer(template), 3)); + new DeadLetterPublishingRecoverer(template), new FixedBackOff(0L, 2))); return factory; } diff --git a/samples/sample-03/pom.xml b/samples/sample-03/pom.xml index d8b949b9..e50ef42d 100644 --- a/samples/sample-03/pom.xml +++ b/samples/sample-03/pom.xml @@ -5,7 +5,7 @@ com.example kafka-sample-03 - 0.0.1-SNAPSHOT + 2.3.1.BUILD-SNAPSHOT jar kafka-sample-03 @@ -14,7 +14,7 @@ org.springframework.boot spring-boot-starter-parent - 2.1.0.RELEASE + 2.2.0.BUILD-SNAPSHOT diff --git a/samples/sample-03/src/main/java/com/example/Application.java b/samples/sample-03/src/main/java/com/example/Application.java index 34219843..b0554a54 100644 --- a/samples/sample-03/src/main/java/com/example/Application.java +++ b/samples/sample-03/src/main/java/com/example/Application.java @@ -30,6 +30,7 @@ import org.springframework.boot.autoconfigure.kafka.ConcurrentKafkaListenerConta import org.springframework.context.annotation.Bean; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; +import org.springframework.kafka.config.TopicBuilder; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.converter.BatchMessagingMessageConverter; @@ -56,9 +57,10 @@ public class Application { @Bean public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory( ConcurrentKafkaListenerContainerFactoryConfigurer configurer, - ConsumerFactory kafkaConsumerFactory, - KafkaTemplate template) { - ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); + ConsumerFactory kafkaConsumerFactory) { + + ConcurrentKafkaListenerContainerFactory factory = + new ConcurrentKafkaListenerContainerFactory<>(); configurer.configure(factory, kafkaConsumerFactory); factory.setBatchListener(true); factory.setMessageConverter(batchConverter()); @@ -92,8 +94,13 @@ public class Application { } @Bean - public NewTopic topic() { - return new NewTopic("topic2", 1, (short) 1); + public NewTopic topic2() { + return TopicBuilder.name("topic2").partitions(1).replicas(1).build(); + } + + @Bean + public NewTopic topic3() { + return TopicBuilder.name("topic3").partitions(1).replicas(1).build(); } }