WIP
[demos/kafka/training] / src / main / java / de / juplo / kafka / Application.java
1 package de.juplo.kafka;
2
3 import lombok.extern.slf4j.Slf4j;
4 import org.apache.kafka.clients.consumer.Consumer;
5 import org.springframework.beans.factory.annotation.Autowired;
6 import org.springframework.boot.ApplicationArguments;
7 import org.springframework.boot.ApplicationRunner;
8 import org.springframework.boot.SpringApplication;
9 import org.springframework.boot.autoconfigure.SpringBootApplication;
10 import org.springframework.context.annotation.Bean;
11 import org.springframework.kafka.core.ConsumerFactory;
12
13 import javax.annotation.PreDestroy;
14 import java.util.concurrent.ExecutorService;
15 import java.util.concurrent.TimeUnit;
16
17
18 @SpringBootApplication
19 @Slf4j
20 public class Application implements ApplicationRunner
21 {
22   @Autowired
23   Consumer<String, Message> consumer;
24   @Autowired
25   SimpleConsumer simpleConsumer;
26
27
28   @Override
29   public void run(ApplicationArguments args) throws Exception
30   {
31     log.info("Starting EndlessConsumer");
32     simpleConsumer.start();
33   }
34
35   @PreDestroy
36   public void shutdown()
37   {
38     log.info("Signaling the consumer to quit its work");
39     consumer.wakeup();
40   }
41
42   @Bean(destroyMethod = "close")
43   public Consumer<String, Message> kafkaConsumer(ConsumerFactory<String, Message> factory)
44   {
45     return factory.createConsumer();
46   }
47
48
49   public static void main(String[] args)
50   {
51     SpringApplication.run(Application.class, args);
52   }
53 }