Code ví dụ Spring Boot Kafka (Producer, Consumer Kafka Spring)

Code ví dụ Spring Boot Kafka (Producer, Consumer Kafka Spring)

(Xem lại: Cài đặt, chạy Apache Kafka, Apache Zookeeper trên windows)

(Xem lại: Cài đặt, cấu hình Apache Kafka, Apache Zookeeper trên Ubuntu)

(Xem lại: Code ví dụ Spring Boot Intellij)

1. Cấu trúc Project

Code ví dụ Spring Boot Kafka (Producer, Consumer Kafka Spring)

Thư viện sử dụng:

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>2.3.0.RELEASE</version>
        <relativePath/> <!-- lookup parent from repository -->
    </parent>
    <groupId>stackjava.com</groupId>
    <artifactId>spring-kafka</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <name>spring-kafka</name>
    <description>Demo project for Spring Boot</description>

    <properties>
        <java.version>1.8</java.version>
    </properties>

    <dependencies>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-test</artifactId>
            <scope>test</scope>
            <exclusions>
                <exclusion>
                    <groupId>org.junit.vintage</groupId>
                    <artifactId>junit-vintage-engine</artifactId>
                </exclusion>
            </exclusions>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.springframework.kafka/spring-kafka -->
        <dependency>
            <groupId>org.springframework.kafka</groupId>
            <artifactId>spring-kafka</artifactId>
            <version>2.5.1.RELEASE</version>
        </dependency>

    </dependencies>

    <build>
        <plugins>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
            </plugin>
        </plugins>
    </build>

</project>

Trong ví dụ này mình dùng thêm thư viện spring-kafka

Cấu hình Consumer

package stackjava.com.springbootkafka.configuration;


import java.util.HashMap;
import java.util.Map;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;

@EnableKafka
@Configuration
public class KafkaConsumerConfig {
    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "group-id");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);

        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "100");
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        return new DefaultKafkaConsumerFactory<>(props);
    }
    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String>
                factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}
  • annotation @EnableKafka được dùng để detect các annotation @KafkaListener
  • method consumerFactory cấu hình consumer

Sau khi có class KafkaConsumerConfig.java, muốn consume 1 topic nào đó ta chỉ cần sử dụng annotation @KafkaListener

cho method tương ứng, ví dụ:

@KafkaListener(topics = "demo", groupId = "group-id")
public void listen(String message) {
    System.out.println("Received Message in group - group-id: " + message);
}

Cấu hình producer

package stackjava.com.springbootkafka.configuration;
import java.util.HashMap;
import java.util.Map;

import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;

@Configuration
public class KafkaProducerConfig {
    @Bean
    public ProducerFactory<String, String> producerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return new DefaultKafkaProducerFactory<>(configProps);
    }
    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

Khi muốn send message tới topic nào ta chỉ cần Autowire KafkaTemplate và gọi hàm send. Ví dụ: hàm gửi message tới topic test.

@Autowired
private KafkaTemplate<String, String> kafkaTemplate;

public void sendMessage(String msg) {
    kafkaTemplate.send("test", msg);
}

2. Demo:

package stackjava.com.springbootkafka.listener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;

import java.util.Date;

@Component
public class StartupListener implements ApplicationRunner {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    public void sendMessage(String msg) {
        kafkaTemplate.send("test", msg);
    }

    @KafkaListener(topics = "demo", groupId = "group-id")
    public void listen(String message) {
        System.out.println("Received Message in group - group-id: " + message);
    }
    @Override
    public void run(ApplicationArguments args) throws Exception {
        for (int i = 0; i < 1000; i++) {
            sendMessage("Now: " + new Date());
            Thread.sleep(2000);
        }
    }
}
  • Class StartupListener thực hiện implements ApplicationRunner nên nó sẽ được chạy hàm run ngay khi start project.
  • Trong ví dụ này StartupListener thực hiện:
    • consume topic demo và in ra message
    • cứ 2 giây thì gửi 1 message tới topic test.

Trong ví dụ này mình thực hiện connect tới kafka ở địa chỉ localhost:9092

Demo: start zookeeper và kafka.

Start project (chạy file SpringBootKafkaApplication.java) và mở command line consume topic test:

Code ví dụ Spring Boot Kafka (Produce, Consume Kafka Spring)

 

Mở command line và tạo producer để gửi message vào topic test:

Code ví dụ Spring Boot Kafka (Produce, Consume Kafka Spring)

 

Okay, Done!

Download code ví dụ trên tại đây hoặc trên github: https://github.com/stackjava/spring-boot-kafka

References: https://docs.spring.io/spring-kafka/reference/html/

stackjava.com