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

kafka springBoot配置的實現(xiàn)

 更新時間:2023年11月16日 09:45:48   作者:weixin_41827053  
本文主要介紹了kafka springBoot配置的實現(xiàn),通過詳細解析Spring Boot for Apache Kafka的配置選項,以及如何優(yōu)化Kafka生產者和消費者的屬性設置,感興趣的可以了解一下

1、properties 配置

control.command.kafka.enabled=true
control.command.kafka.bootstrap-servers=172.0.0.1:9092
control.command.kafka.command-topics=lastTopic
control.command.kafka.consumer.group-id=consumer-eslink-iwater-control-command
control.command.kafka.consumer.properties.session.timeout.ms=30000
control.command.kafka.consumer.properties.request.timeout.ms=90000
control.command.kafka.consumer.fetch-min-size=10KB
control.command.kafka.consumer.fetch-max-wait=500
control.command.kafka.consumer.max-poll-records=1000
control.command.kafka.consumer.auto-offset-reset=earliest
control.command.kafka.listener.ack-mode=MANUAL_IMMEDIATE
control.command.kafka.listener.concurrency=1
control.command.kafka.listener.type=SINGLE
control.command.kafka.producer.acks=all
control.command.kafka.producer.batchSize=4096
control.command.kafka.producer.bufferMemory=40960
control.command.kafka.producer.linger=10
control.command.kafka.producer.retries=3

2、Config 配置

@Configuration
@ConditionalOnExpression("${control.command.kafka.enabled:false}")
@EnableKafka
public class ControlCommandKafkaConfig {

    @Bean("controlCommandKafkaProperties")
    @ConfigurationProperties("control.command.kafka")
    @Primary
    public KafkaProperties kafkaProperties() {
        return new KafkaProperties();
    }

    @Bean("controlCommandKafkaConsumerFactory")
    public ConsumerFactory<Object, Object> kafkaConsumerFactory(
            @Qualifier("controlCommandKafkaProperties") KafkaProperties inProps) {
        DefaultKafkaConsumerFactory<Object, Object> consumerFactory = new DefaultKafkaConsumerFactory<>(inProps.buildConsumerProperties());
        return consumerFactory;
    }

    @Bean("controlCommandBatchFactory")
    @DependsOn("controlCommandKafkaProperties")
    public KafkaListenerContainerFactory<?> egBatchFactory(@Qualifier("controlCommandKafkaProperties") KafkaProperties inProps,
                                                           @Qualifier("controlCommandKafkaConsumerFactory") ConsumerFactory consumerFactory) {
        ConcurrentKafkaListenerContainerFactory<Object, Object> listenerFactory =
                new ConcurrentKafkaListenerContainerFactory<>();
        listenerFactory.setConsumerFactory(consumerFactory);
        configureListenerFactory(listenerFactory, inProps);
        configureContainer(listenerFactory.getContainerProperties(), inProps);
        return listenerFactory;
    }

    @Bean("controlCommandKafkaTemplate")
    public KafkaTemplate<String, String> kafkaTemplate(@Qualifier("controlCommandKafkaProperties") KafkaProperties inProps) {
        return new KafkaTemplate(new DefaultKafkaProducerFactory(inProps.buildProducerProperties()));
    }

    private void configureListenerFactory(ConcurrentKafkaListenerContainerFactory<Object, Object> factory, KafkaProperties inProps) {
        PropertyMapper map = PropertyMapper.get().alwaysApplyingWhenNonNull();
        KafkaProperties.Listener properties = inProps.getListener();
        // 設置監(jiān)聽線程并發(fā)數(shù) control.command.kafka.listener.concurrency
        map.from(properties::getConcurrency).to(factory::setConcurrency);
        // control.kafka.listener.type=batch 時 批量監(jiān)聽
        if (properties.getType().equals(KafkaProperties.Listener.Type.BATCH)) {
            factory.setBatchListener(true);
        }
    }

    private void configureContainer(ContainerProperties container, KafkaProperties inProps) {
        PropertyMapper map = PropertyMapper.get().alwaysApplyingWhenNonNull();
        KafkaProperties.Listener properties = inProps.getListener();
        map.from(properties::getAckMode).to(container::setAckMode);
        map.from(properties::getClientId).to(container::setClientId);
        map.from(properties::getAckCount).to(container::setAckCount);
        map.from(properties::getAckTime).as(Duration::toMillis).to(container::setAckTime);
        map.from(properties::getPollTimeout).as(Duration::toMillis).to(container::setPollTimeout);
        map.from(properties::getNoPollThreshold).to(container::setNoPollThreshold);
        map.from(properties::getIdleEventInterval).as(Duration::toMillis).to(container::setIdleEventInterval);
        map.from(properties::getMonitorInterval).as(Duration::getSeconds).as(Number::intValue)
                .to(container::setMonitorInterval);
        map.from(properties::getLogContainerConfig).to(container::setLogContainerConfig);
        map.from(properties::isMissingTopicsFatal).to(container::setMissingTopicsFatal);
    }
}

2.1 @EnableKafka 注解

‘@EnableKafka’ 是用于在 Spring Boot 應用程序中啟用 Apache Kafka 的注解。當你在 Spring Boot 應用程序的配置類上添加 @EnableKafka 注解時,它會激活 Kafka 基礎設施,使你能夠在應用程序中使用 Kafka 相關的組件。

以下是對這個注解的簡要解釋:

基礎設施激活:@EnableKafka 注解告訴 Spring 在應用程序中設置所需的 Kafka 基礎設施。這包括創(chuàng)建必要的 Kafka bean 和配置。

Kafka Template:在啟用 Kafka 后,你可以使用 Spring Kafka 提供的 KafkaTemplate 類輕松地向 Kafka 主題發(fā)送消息。KafkaTemplate 抽象了底層的 Kafka 生產者 API,簡化了消息發(fā)送的過程。

消息監(jiān)聽器:通過 @EnableKafka,你還可以使用 Spring 的 @KafkaListener 注解設置 Kafka 消息消費者。@KafkaListener 注解應用于 Spring 組件中的方法,使該方法可以作為 Kafka 消息監(jiān)聽器。當 Kafka 主題中有可用的消息時,帶有 @KafkaListener 注解的方法將自動被調用,并處理消息內容。

下面是在 Spring Boot 應用程序中使用 @EnableKafka 的示例代碼:

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.kafka.annotation.EnableKafka;

@SpringBootApplication
@EnableKafka
public class MyKafkaApplication {
    public static void main(String[] args) {
        SpringApplication.run(MyKafkaApplication.class, args);
    }
}

通過這樣的配置,你的 Spring Boot 應用程序將啟用 Kafka 支持,你可以使用 KafkaTemplate 進行消息發(fā)送,使用 @KafkaListener 進行消息消費。請確保在 application.properties 或 application.yml 文件中配置 Kafka 相關屬性,例如 Kafka 代理地址和其他必要的配置。

請注意,為了使 @EnableKafka 正常工作,你需要在項目的構建配置中包含所需的 Kafka 依賴項。對于 Maven,你可以添加以下依賴項:

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

一旦 Kafka 正在運行,并且應用程序正確配置了 @EnableKafka,你就可以在 Spring Boot 應用程序中構建基于 Kafka 的消息傳遞功能。

2.2、@ConditionalOnExpression 注解

@ConditionalOnExpression 是 Spring Boot 中的一個條件注解之一。它允許你根據(jù)給定的 SpEL 表達式來決定是否啟用或禁用某個 Bean 或配置。
條件注解可以用于在 Spring Boot 應用程序中根據(jù)特定條件來動態(tài)創(chuàng)建 Bean 或配置,從而根據(jù)不同的配置或環(huán)境來靈活地管理應用程序的行為。
@ConditionalOnExpression 的工作方式是:它在配置類或 Bean 上進行標記,然后在應用程序啟動過程中解析 SpEL 表達式。如果 SpEL 表達式的結果為 true,則相關的 Bean 或配置將被啟用,否則將被禁用。
以下是 @ConditionalOnExpression 的示例使用方式:

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Conditional;
import org.springframework.context.annotation.Configuration;

@Configuration
public class MyConfiguration {

    @Bean
    @ConditionalOnExpression("${myapp.feature.enabled:true}")
    public MyBean myBean() {
        // 返回需要創(chuàng)建的 Bean 實例
        return new MyBean();
    }
}

2.3、 @RefreshScope 注解

@RefreshScope 是 Spring Cloud 中的一個注解,用于實現(xiàn)動態(tài)刷新 Spring Bean 的配置信息。它通常與 Spring Cloud Config 配合使用,能夠在運行時更新應用程序配置,而無需重啟應用。
在微服務架構中,使用 Spring Cloud Config 可以將配置信息集中管理,然后通過 @RefreshScope 注解實現(xiàn)配置的動態(tài)刷新。這使得應用程序在運行時可以獲取最新的配置信息,而不需要停止和啟動應用。

2.4、@DependsOn 注解

@DependsOn 是 Spring Framework 中的一個注解,用于指定 Spring Bean 之間的依賴關系。通過在 Bean 上添加 @DependsOn 注解,你可以確保指定的 Bean 會在其所依賴的其他 Bean 初始化之后再進行初始化。
當一個 Bean 希望在另一個 Bean 初始化完成后再初始化時,可以使用 @DependsOn 注解來定義這種依賴關系。這對于確保 Bean 之間的正確順序初始化非常有用,特別是當某些 Bean 需要依賴其他 Bean 才能正確地進行初始化或工作時。

2.5、 listener.ack-mode 消息的確認模式

在 Spring Kafka 中,AckMode 是用于配置消息消費者的消息確認模式(Acknowledgment Mode)。這個枚舉類型用于決定在消費者處理完 Kafka 消息后如何向 Kafka 服務器發(fā)送確認,告知服務器消息是否已經被成功消費。

Spring Kafka 支持以下幾種 AckMode:

  • AUTO: 這是默認的確認模式。在這種模式下,消費者會自動在處理完消息后向 Kafka 服務器發(fā)送確認。當消費者成功處理消息后,Kafka 服務器會將偏移量(offset)移動到下一條消息,表示該消息已成功消費。但是在這種模式下,如果處理消息時發(fā)生異常,Kafka 服務器會重新發(fā)送相同的消息,可能會導致消息的重復消費。

  • MANUAL: 在這種模式下,消費者需要手動管理消息的確認。當消費者成功處理消息后,需要調用 Acknowledgment 對象的 acknowledge() 方法,手動向 Kafka 服務器發(fā)送確認。這樣可以確保消息的準確處理,避免重復消費。

  • MANUAL_IMMEDIATE: 這是另一種手動確認模式,與 MANUAL 模式類似。區(qū)別在于,當消費者調用 acknowledge() 方法時,它會立即提交確認而不是等待下一次輪詢。這樣可以更快地將確認提交給 Kafka 服務器。

你可以通過在 @KafkaListener 注解中設置 AckMode 來配置消息消費者的確認模式。例如:

import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;

@KafkaListener(topics = "my_topic", groupId = "my_group", ackMode = "MANUAL")
public void processMessage(String message, Acknowledgment acknowledgment) {
    // 處理消息的邏輯
    // 手動確認消息
    acknowledgment.acknowledge();
}

在上面的例子中,我們將 ackMode 設置為 MANUAL,表明消息的確認模式為手動確認。在處理消息后,我們手動調用 acknowledgment.acknowledge() 來確認消息的消費。

選擇合適的 AckMode 取決于你的應用程序需求和消費者的可靠性要求。如果應用程序對消息的重復消費有一定的容忍度,并且希望簡化消費者的處理邏輯,可以選擇 AUTO 模式。如果應用程序對消息的準確性要求較高,并且愿意手動確認消息的處理,可以選擇 MANUAL 或 MANUAL_IMMEDIATE 模式。

2.6、kafka.consumer.auto-offset-reset 偏移量

kafka.consumer.auto-offset-reset 是 Kafka 消費者的一個重要配置屬性,它決定了當一個新的消費者加入消費者組或者消費者在某個分區(qū)上沒有有效的偏移量時,消費者應該從何處開始消費消息。

這個屬性有以下幾個可能的值:

  • rliest: 如果消費者在某個分區(qū)上沒有有效的偏移量,或者消費者組第一次加入時,從最早的可用偏移量開始消費消息。換句話說,從分區(qū)的起始位置開始消費。

  • latest: 如果消費者在某個分區(qū)上沒有有效的偏移量,或者消費者組第一次加入時,從最新的可用偏移量開始消費消息。換句話說,只消費自加入后產生的新消息。

  • none: 如果消費者在某個分區(qū)上沒有有效的偏移量,拋出一個異常。這種情況下,消費者必須手動設置初始偏移量。

  • anything else: 拋出一個異常。

這個屬性通常在 Kafka 消費者的配置中使用,用來控制消費者在特定情況下的起始消費位置。你可以根據(jù)你的業(yè)務需求來選擇適合的值。如果你希望消費者能夠從最早的消息開始消費,以確保不錯過任何消息,可以將它設置為 earliest。如果你只關心新產生的消息,可以將它設置為 latest。

3、KafkaProducer 生產者

@Component
@RefreshScope
@ConditionalOnExpression("${control.command.kafka.enabled:false}")
@DependsOn(value = {"controlCommandKafkaTemplate"})
public class ControlCommandKafkaProducer extends BaseLogable {

    @Resource(name = "controlCommandKafkaTemplate")
    private KafkaTemplate<String, String> controlCommandKafkaTemplate;

    public void send(String key, String topic, String msg) {
        bizLogger.info("sent kafka msg, topic:{}, key: {}, msg: {}", topic, key, msg);
        ListenableFuture<SendResult<String, String>> listenableFuture = controlCommandKafkaTemplate.send(topic, key, msg);
        listenableFuture.addCallback(result -> bizLogger.info("sent msg success:{}", msg),
                e -> {
                    bizLogger.error("sent msg failure, msg: {}", msg, e);
                });
    }

}

3.1 org.springframework.kafka.core.KafkaTemplate 使用

3.1.1 發(fā)送消息

import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;

@Service
public class KafkaProducerService {
    
    private final KafkaTemplate<String, String> kafkaTemplate;
    
    public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }
    
    public void sendMessage(String topic, String message) {
        kafkaTemplate.send(topic, message);
    }
}

3.1.2 發(fā)送消息并指定分區(qū)

public void sendMessageToPartition(String topic, int partition, String message) {
    kafkaTemplate.send(topic, partition, null, message);
}

3.1.3 發(fā)送消息并指定鍵

public void sendMessageWithKey(String topic, String key, String message) {
    kafkaTemplate.send(topic, key, message);
}

3.1.4 發(fā)送消息并等待確認

import org.apache.kafka.clients.producer.ProducerRecord;
import org.springframework.kafka.support.SendResult;
import org.springframework.kafka.support.future.ListenableFuture;
import org.springframework.util.concurrent.ListenableFutureCallback;

public void sendMessageAndWaitForConfirmation(String topic, String message) {
    ProducerRecord<String, String> record = new ProducerRecord<>(topic, message);
    ListenableFuture<SendResult<String, String>> future = kafkaTemplate.send(record);
    
    future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
        @Override
        public void onSuccess(SendResult<String, String> result) {
            // 消息發(fā)送成功的處理邏輯
        }

        @Override
        public void onFailure(Throwable ex) {
            // 消息發(fā)送失敗的處理邏輯
        }
    });
}

4、KafkaConsumer 消費者

4.1 單條消費

單條消費時配置

control.command.kafka.listener.type=SINGLE

@Component
@RefreshScope
@ConditionalOnExpression("${control.command.kafka.enabled:false}")
@DependsOn(value = {"controlCommandBatchFactory"})
public class ControlCommandKafkaConsumer extends BaseLogable {


    /**
     * 監(jiān)聽主題,單條消費
     */
    @KafkaListener(id = "${control.command.kafka.consumer.groupId}", topics = "${kafka.command.topic.ecgs-command-up}", containerFactory = "controlCommandBatchFactory")
    public void listen(ConsumerRecord<String, String> record, Acknowledgment ack) {
        bizLogger.info("kafka msg::: " + record.value() + ", key: " + record.key());
        //接收Kafka消息
        String data = record.value();

        //接收消息確認
        ack.acknowledge();
    }

}

4.2 多條消費

多條消費時配置

control.command.kafka.listener.type=BATCH

package cc.eslink.yq.iwater.kafka;

import cc.eslink.common.base.BaseLogable;
import cc.eslink.yq.iwater.dto.schedule.AlarmInfoDTO;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.alibaba.fastjson.TypeReference;
import org.apache.commons.lang3.StringUtils;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.context.annotation.DependsOn;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;

import java.util.List;


@Component
@RefreshScope
@ConditionalOnExpression("${control.command.kafka.enabled:false}")
@DependsOn(value = {"controlCommandBatchFactory"})
public class ControlCommandKafkaConsumer extends BaseLogable {


    @KafkaListener(id = "${control.command.kafka.consumer.groupId}", topics = "${kafka.command.topic.ecgs-command-up}", containerFactory = "controlCommandBatchFactory")
    public void listen(List<ConsumerRecord<String, String>> records, Acknowledgment ack) {
        bizLogger.info("== 接收到 kafka msg 條數(shù):::{} ==", records.size());
        //接收消息確認
        ack.acknowledge();
        //接收Kafka消息
        for (ConsumerRecord<String, String> record : records) {
            try {
                final AlarmInfoDTO alarmInfoDTO = JSON.parseObject(record.value(), new TypeReference<AlarmInfoDTO>() {});
            } catch (Exception e) {
                expLogger.error("alarm.dispose error", record, e);
            }
        }


    }

}

到此這篇關于kafka springBoot配置的實現(xiàn)的文章就介紹到這了,更多相關kafka springBoot 配置內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家! 

相關文章

  • SpringBoot使用classfinal-maven-plugin插件加密Jar包的示例代碼

    SpringBoot使用classfinal-maven-plugin插件加密Jar包的示例代碼

    這篇文章給大家介紹了SpringBoot使用classfinal-maven-plugin插件加密Jar包的實例,文中通過代碼示例和圖文講解的非常詳細,對大家的學習或工作有一定的幫助,需要的朋友可以參考下
    2024-02-02
  • JAVA中split函數(shù)的常見用法實例

    JAVA中split函數(shù)的常見用法實例

    Java中我們可以利用split把字符串按照指定的分割符進行分割,然后返回字符串數(shù)組,下面這篇文章主要給大家介紹了關于JAVA中split函數(shù)的常見用法,文中通過實例代碼介紹的非常詳細,需要的朋友可以參考下
    2022-07-07
  • java必懂的冷知識點之Base64加密與解密

    java必懂的冷知識點之Base64加密與解密

    這篇文章主要介紹了java必懂的冷知識點之Base64加密與解密的相關資料,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-03-03
  • Springboot詳細講解RocketMQ實現(xiàn)順序消息的發(fā)送與消費流程

    Springboot詳細講解RocketMQ實現(xiàn)順序消息的發(fā)送與消費流程

    RocketMQ作為一款純java、分布式、隊列模型的開源消息中間件,支持事務消息、順序消息、批量消息、定時消息、消息回溯等,本篇我們了解如何實現(xiàn)順序消息的發(fā)送與消費
    2022-06-06
  • SpringBoot請求響應方式示例詳解

    SpringBoot請求響應方式示例詳解

    這篇文章主要介紹了SpringBoot請求響應的相關操作,本文通過示例代碼給大家介紹的非常詳細,感興趣的朋友跟隨小編一起看看吧
    2024-06-06
  • 詳解Java中hashCode的作用

    詳解Java中hashCode的作用

    這篇文章主要介紹了詳解Java中hashCode的作用的相關資料,需要的朋友可以參考下
    2017-03-03
  • 分享Java死鎖的4種排查工具

    分享Java死鎖的4種排查工具

    這篇文章主要介紹了分享Java死鎖的4種排查工具,死鎖指的是兩個或兩個以上的運算單元,都在等待對方停止執(zhí)行,以取得系統(tǒng)資源,但是沒有一方提前退出,就稱為死鎖,下文更多相關內容需要的小伙伴可以參考一下
    2022-05-05
  • Java class文件格式之屬性_動力節(jié)點Java學院整理

    Java class文件格式之屬性_動力節(jié)點Java學院整理

    在本文中, 主要講解了class文件中的一些屬性。 這些屬性可以出現(xiàn)在class文件中的對個地方, 用來描述一些其他信息
    2017-06-06
  • 詳解SpringBoot基礎之banner玩法解析

    詳解SpringBoot基礎之banner玩法解析

    SpringBoot項目啟動時會在控制臺打印一個默認的啟動圖案,這個圖案就是我們要講的banner,這篇文章主要介紹了SpringBoot基礎之banner玩法解析,感興趣的小伙伴們可以參考一下
    2019-04-04
  • java jdk動態(tài)代理詳解

    java jdk動態(tài)代理詳解

    動態(tài)代理類的Class實例是怎么生成的呢,是通過ProxyGenerator類來生成動態(tài)代理類的class字節(jié)流,把它載入方法區(qū)
    2013-09-09

最新評論

江孜县| 大渡口区| 安平县| 裕民县| 油尖旺区| 县级市| 齐河县| 仁怀市| 碌曲县| 会东县| 博乐市| 建德市| 洛南县| 木兰县| 喀喇| 瑞丽市| 和顺县| 竹北市| 岳西县| 南郑县| 吴桥县| 白朗县| 日喀则市| 宜昌市| 定安县| 囊谦县| 大冶市| 成武县| 游戏| 绥芬河市| 周宁县| 右玉县| 崇礼县| 固阳县| 合肥市| 丰都县| 延边| 定陶县| 胶南市| 新安县| 龙井市|