最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

java集成kafka實例代碼

 更新時間:2024年12月30日 14:15:59   作者:沉墨的夜  
文章介紹了如何在Java項目中集成Apache Kafka以實現(xiàn)消息的生產(chǎn)和消費,通過添加Maven依賴、配置生產(chǎn)者和消費者、使用SpringBoot簡化集成以及控制消費者的啟動和停止,可以實現(xiàn)高效的消息處理

java集成kafka

要在 Java 項目中集成 Apache Kafka 以實現(xiàn)消息的生產(chǎn)和消費,步驟如下:

1. 引入 Maven 依賴

在您的 pom.xml 文件中添加以下依賴,以包含 Kafka 客戶端庫:

<dependencies>
    <!-- Kafka Clients -->
    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
        <version>2.8.0</version>
    </dependency>
    <!-- 如果使用 Spring Boot,可添加以下依賴 -->
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
        <version>2.7.0</version>
    </dependency>
</dependencies>

2. 配置 Kafka 生產(chǎn)者

首先,設(shè)置生產(chǎn)者的配置屬性:

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;

public class KafkaProducerExample {
    public static void main(String[] args) {
        // 配置屬性
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        // 創(chuàng)建生產(chǎn)者
        Producer<String, String> producer = new KafkaProducer<>(props);

        // 發(fā)送消息
        for (int i = 0; i < 10; i++) {
            ProducerRecord<String, String> record = new ProducerRecord<>("your_topic", "key" + i, "value" + i);
            producer.send(record);
        }

        // 關(guān)閉生產(chǎn)者
        producer.close();
    }
}

3. 配置 Kafka 消費者

接下來,設(shè)置消費者的配置屬性,并訂閱主題以消費消息:

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class KafkaConsumerExample {
    public static void main(String[] args) {
        // 配置屬性
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "your_group_id");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());

        // 創(chuàng)建消費者
        Consumer<String, String> consumer = new KafkaConsumer<>(props);

        // 訂閱主題
        consumer.subscribe(Collections.singletonList("your_topic"));

        // 持續(xù)消費消息
        try {
            while (true) {
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
                records.forEach(record -> {
                    System.out.printf("Consumed message: key = %s, value = %s, offset = %d%n",
                            record.key(), record.value(), record.offset());
                });
            }
        } finally {
            // 關(guān)閉消費者
            consumer.close();
        }
    }
}

4. 使用 Spring Boot 集成 Kafka

如果您使用 Spring Boot,可以通過配置 KafkaTemplate(用于生產(chǎn)消息)和使用 @KafkaListener 注解(用于消費消息)來簡化 Kafka 的集成。

生產(chǎn)者配置:

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;

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

@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.getName());
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

使用 KafkaTemplate 發(fā)送消息:

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;

@Service
public class KafkaProducerService {

    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    public void sendMessage(String topic, String key, String value) {
        kafkaTemplate.send(topic, key, value);
    }
}

消費者配置:

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;

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

@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, "your_group_id");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}

使用 @KafkaListener 消費消息:

在 Spring Boot 中,@KafkaListener 注解用于監(jiān)聽指定的 Kafka 主題,并在收到消息時觸發(fā)相應(yīng)的方法。

以下是一個基本示例:

import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Service;

@Service
public class KafkaConsumerService {

    @KafkaListener(topics = "your_topic", groupId = "your_group_id")
    public void listen(String message) {
        System.out.println("Received message: " + message);
        // 在此處添加處理邏輯
    }
}

 

在上述代碼中:

  • topics:指定要監(jiān)聽的 Kafka 主題。
  • groupId:指定消費者組 ID。

listen 方法:當(dāng)有新消息發(fā)布到指定主題時,該方法會被調(diào)用,message 參數(shù)包含消息的內(nèi)容。

批量消費消息

如果希望一次處理多條消息,可以啟用批量監(jiān)聽。

首先,需要配置一個支持批量消費的 KafkaListenerContainerFactory

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;

@EnableKafka
@Configuration
public class KafkaConsumerConfig {

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
            ConsumerFactory<String, String> consumerFactory) {
        ConcurrentKafkaListenerContainerFactory<String, String> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);
        factory.setBatchListener(true); // 啟用批量監(jiān)聽
        return factory;
    }
}

然后,在消費者服務(wù)中使用 @KafkaListener 注解,并指定使用上述配置的工廠:

import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Service;
import java.util.List;

@Service
public class KafkaBatchConsumerService {

    @KafkaListener(
        topics = "your_topic",
        groupId = "your_group_id",
        containerFactory = "kafkaListenerContainerFactory"
    )
    public void listen(List<String> messages) {
        System.out.println("Received batch messages: " + messages);
        // 在此處添加批量處理邏輯
    }
}

在上述代碼中:

  • containerFactory:指定使用支持批量消費的工廠。

listen 方法的參數(shù)類型為 List<String>,用于接收一批消息。

控制消費者的啟動和停止

在某些情況下,可能需要在運行時控制 Kafka 消費者的啟動和停止。

可以通過 KafkaListenerEndpointRegistry 來實現(xiàn):

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.listener.KafkaListenerEndpointRegistry;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;

@Service
public class KafkaListenerManager {

    @Autowired
    private KafkaListenerEndpointRegistry registry;

    // 啟動監(jiān)聽器
    public void startListener(String listenerId) {
        MessageListenerContainer listenerContainer = registry.getListenerContainer(listenerId);
        if (listenerContainer != null && !listenerContainer.isRunning()) {
            listenerContainer.start();
        }
    }

    // 停止監(jiān)聽器
    public void stopListener(String listenerId) {
        MessageListenerContainer listenerContainer = registry.getListenerContainer(listenerId);
        if (listenerContainer != null && listenerContainer.isRunning()) {
            listenerContainer.stop();
        }
    }
}

在上述代碼中:

  • startListener 方法用于啟動指定的監(jiān)聽器。
  • stopListener 方法用于停止指定的監(jiān)聽器。
  • listenerId 對應(yīng)于 @KafkaListener 注解中的 id 屬性。

通過這種方式,可以在應(yīng)用運行時根據(jù)需要動態(tài)地控制 Kafka 消費者的行為。

通過上述配置和代碼示例,可以在 Spring Boot 項目中有效地集成 Kafka,實現(xiàn)消息的生產(chǎn)和消費功能。

總結(jié)

以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • Java動態(tài)代理靜態(tài)代理實例分析

    Java動態(tài)代理靜態(tài)代理實例分析

    這篇文章主要介紹了Java動態(tài)代理靜態(tài)代理實例分析,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2020-03-03
  • Spring Boot與前端配合與Idea配置部署操作過程

    Spring Boot與前端配合與Idea配置部署操作過程

    這篇文章主要介紹了Spring Boot與前端配合與Idea配置部署的操作過程,本文圖文并茂給大家介紹的非常詳細,需要的朋友可以參考下
    2018-02-02
  • Java常用占位符方法簡單代碼實例

    Java常用占位符方法簡單代碼實例

    占位符是Java中常用的技術(shù),用于在字符串中插入變量值或動態(tài)生成字符串,這篇文章主要給大家介紹了關(guān)于Java常用占位符方法的相關(guān)資料,文中介紹的非常詳細,需要的朋友可以參考下
    2024-01-01
  • SpringMVC中Model與Session的區(qū)別說明

    SpringMVC中Model與Session的區(qū)別說明

    這篇文章主要介紹了SpringMVC中Model與Session的區(qū)別說明,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-12-12
  • 基于ScheduledExecutorService的兩種方法(詳解)

    基于ScheduledExecutorService的兩種方法(詳解)

    下面小編就為大家?guī)硪黄赟cheduledExecutorService的兩種方法(詳解)。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-10-10
  • 深入理解Java設(shè)計模式之中介者模式

    深入理解Java設(shè)計模式之中介者模式

    這篇文章主要介紹了JAVA設(shè)計模式之中介者模式的的相關(guān)資料,文中示例代碼非常詳細,供大家參考和學(xué)習(xí),感興趣的朋友可以了解
    2021-11-11
  • springboot整合RabbitMQ中死信隊列的實現(xiàn)

    springboot整合RabbitMQ中死信隊列的實現(xiàn)

    死信是無法被消費的消息,產(chǎn)生原因包括消息TTL過期、隊列最大長度達到以及消息被拒絕且不重新排隊,RabbitMQ的死信隊列機制能夠有效防止消息數(shù)據(jù)丟失,適用于訂單業(yè)務(wù)等場景,本文就來介紹一下
    2024-10-10
  • SpringSecurity的@EnableWebSecurity注解詳解

    SpringSecurity的@EnableWebSecurity注解詳解

    這篇文章主要介紹了SpringSecurity的@EnableWebSecurity注解詳解,@EnableWebSecurity是開啟SpringSecurity的默認行為,它的上面有一個Import注解導(dǎo)入了WebSecurityConfiguration類,就是往IOC容器中注入了WebSecurityConfiguration這個類,需要的朋友可以參考下
    2023-11-11
  • IDEA EasyCode 一鍵幫你生成所需代碼

    IDEA EasyCode 一鍵幫你生成所需代碼

    這篇文章主要介紹了IDEA EasyCode 一鍵幫你生成所需代碼,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-08-08
  • maven導(dǎo)入本地倉庫jar包,報:Could?not?find?artifact的解決

    maven導(dǎo)入本地倉庫jar包,報:Could?not?find?artifact的解決

    這篇文章主要介紹了maven導(dǎo)入本地倉庫jar包,報:Could?not?find?artifact的解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-03-03

最新評論

林芝县| 宜昌市| 通州区| 黔江区| 安宁市| 沽源县| 泉州市| 兴宁市| 紫云| 乐业县| 洛隆县| 巴楚县| 曲周县| 荆州市| 永兴县| 惠水县| 武宣县| 东宁县| 类乌齐县| 航空| 筠连县| 南康市| 建瓯市| 富顺县| 讷河市| 唐山市| 来安县| 贵州省| 娄烦县| 嘉兴市| 南皮县| 孟津县| 宜宾县| 和政县| 娱乐| 平遥县| 鄂尔多斯市| 凉城县| 固阳县| 临桂县| 乌拉特中旗|