package de.juplo.kafka;
-import org.apache.kafka.clients.consumer.ConsumerRecord;
+import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import java.util.Properties;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
-import java.util.function.Consumer;
@Configuration
public class ApplicationConfiguration
{
@Bean
- public EndlessConsumer endlessConsumer(
+ public WordcountRecordHandler wordcountRecordHandler()
+ {
+ return new WordcountRecordHandler();
+ }
+
+ @Bean
+ public WordcountRebalanceListener wordcountRebalanceListener(
+ WordcountRecordHandler wordcountRecordHandler,
+ PartitionStatisticsRepository repository,
+ Consumer<String, String> consumer,
+ ApplicationProperties properties)
+ {
+ return new WordcountRebalanceListener(
+ wordcountRecordHandler,
+ repository,
+ properties.getClientId(),
+ properties.getTopic(),
+ Clock.systemDefaultZone(),
+ properties.getCommitInterval(),
+ consumer);
+ }
+
+ @Bean
+ public EndlessConsumer<String, String> endlessConsumer(
KafkaConsumer<String, String> kafkaConsumer,
ExecutorService executor,
- PartitionStatisticsRepository repository,
+ WordcountRebalanceListener wordcountRebalanceListener,
+ WordcountRecordHandler wordcountRecordHandler,
ApplicationProperties properties)
{
return
- new EndlessConsumer(
+ new EndlessConsumer<>(
executor,
- repository,
properties.getClientId(),
properties.getTopic(),
- Clock.systemDefaultZone(),
- properties.getCommitInterval(),
- kafkaConsumer);
+ kafkaConsumer,
+ wordcountRebalanceListener,
+ wordcountRecordHandler);
}
@Bean