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

springKafka生產(chǎn)者、消費(fèi)者,數(shù)據(jù)發(fā)送/接收方式

 更新時(shí)間:2025年09月19日 08:50:42   作者:小風(fēng)010766  
這篇文章主要介紹了springKafka生產(chǎn)者、消費(fèi)者,數(shù)據(jù)發(fā)送/接收方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教

配置文件

# 指定kafka server的地址,集群配多個(gè),中間,逗號(hào)隔開(kāi)
spring.kafka.bootstrap-servers=
# 配置kafka 授權(quán)認(rèn)證 --開(kāi)始
spring.kafka.properties.sasl.mechanism=SCRAM-SHA-256
spring.kafka.properties.security.protocol=SASL_PLAINTEXT
spring.kafka.producer.properties.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="xxxx" password="xxxxx";
spring.kafka.consumer.properties.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="xxxx" password="xxxx";
# 配置kafka 授權(quán)認(rèn)證 --結(jié)束
#重試次數(shù)
spring.kafka.producer.retries=3
# 重試間隔時(shí)間 (毫秒)
spring.kafka.producer.retry.backoff.ms=10000
#批量發(fā)送的消息數(shù)量
spring.kafka.producer.batch-size=1000
#32MB的批處理緩沖區(qū)
spring.kafka.producer.buffer-memory=335544320

#默認(rèn)消費(fèi)者組
spring.kafka.consumer.group-id=isms-group
#最早未被消費(fèi)的offset
spring.kafka.consumer.auto-offset-reset=earliest
#批量一次最大拉取數(shù)據(jù)量
spring.kafka.consumer.max-poll-records=500
#自動(dòng)提交時(shí)間間隔,單位ms
spring.kafka.consumer.auto-commit-interval=1000
#批消費(fèi)并發(fā)量,小于或等于Topic的分區(qū)數(shù)
spring.kafka.consumer.batch.concurrency=2
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.config.KafkaListenerContainerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaAdmin;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.listener.ContainerProperties;

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


@Configuration
public class KafkaConfiguration {

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Value("${spring.kafka.producer.retries}")
    private Integer retries;
    @Value("${spring.kafka.producer.retry.backoff.ms}")
    private Integer retryBackoff;

    @Value("${spring.kafka.producer.batch-size}")
    private Integer batchSize;

    @Value("${spring.kafka.producer.buffer-memory}")
    private Integer bufferMemory;

    @Value("${spring.kafka.consumer.group-id}")
    private String groupId;

    @Value("${spring.kafka.consumer.auto-offset-reset}")
    private String autoOffsetReset;

    @Value("${spring.kafka.consumer.max-poll-records}")
    private Integer maxPollRecords;

    @Value("${spring.kafka.consumer.batch.concurrency}")
    private Integer batchConcurrency;

    @Value("${spring.kafka.consumer.auto-commit-interval}")
    private Integer autoCommitInterval;
    @Value("${spring.kafka.properties.security.protocol}")
    private String securityProtocol;
    @Value("${spring.kafka.properties.sasl.mechanism}")
    private String saslMechanism;
    @Value("${spring.kafka.producer.properties.sasl.jaas.config}")
    private String producerSaslJaasConfig;
    @Value("${spring.kafka.consumer.properties.sasl.jaas.config}")
    private String consumerSaslJaasConfig;

    /**
     *  生產(chǎn)者配置信息
     */
    @Bean
    public Map<String, Object> producerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.ACKS_CONFIG, "all");
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ProducerConfig.RETRIES_CONFIG, retries);
        props.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, retryBackoff);
        props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize);
        props.put(ProducerConfig.LINGER_MS_CONFIG, 1);
        props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, bufferMemory);
        props.put(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG, 900000);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        // 消息冪等開(kāi)關(guān),事務(wù)開(kāi)啟必須打開(kāi)消息冪等,但是冪等可以單獨(dú)使用
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
        props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol);
        props.put(SaslConfigs.SASL_MECHANISM, saslMechanism);
        props.put(SaslConfigs.SASL_JAAS_CONFIG, producerSaslJaasConfig);
        return props;
    }

    /**
     *  生產(chǎn)者工廠
     */
    @Bean
    public ProducerFactory<String, String> producerFactory() {
        return new DefaultKafkaProducerFactory<>(producerConfigs());
    }


    /**
     *  生產(chǎn)者模板
     */
    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return  new KafkaTemplate<>(producerFactory());
    }


    @Bean
    public KafkaAdmin kafkaAdmin() {
        return new KafkaAdmin(kafkaTemplate().getProducerFactory().getConfigurationProperties());
    }

    /**
     *  消費(fèi)者配置信息
     */
    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset);
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, maxPollRecords);
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 600000);
        props.put(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG, 300000);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // 設(shè)置自動(dòng)提交改成false
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        props.put(ConsumerConfig.ALLOW_AUTO_CREATE_TOPICS_CONFIG,false);
        props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG,60000);
        props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol);
        props.put(SaslConfigs.SASL_MECHANISM, saslMechanism);
        props.put(SaslConfigs.SASL_JAAS_CONFIG, consumerSaslJaasConfig);
        return props;
    }

    /**
     *  消費(fèi)者批量工廠
     */
/*    @Bean
    public KafkaListenerContainerFactory<?> batchFactory() {
        ConcurrentKafkaListenerContainerFactory<Integer, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfigs()));
        //設(shè)置并發(fā)量,小于或等于Topic的分區(qū)數(shù)
        factory.setConcurrency(batchConcurrency);
        factory.getContainerProperties().setPollTimeout(3000);
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
        //設(shè)置為批量消費(fèi),每個(gè)批次數(shù)量在Kafka配置參數(shù)中設(shè)置ConsumerConfig.MAX_POLL_RECORDS_CONFIG
        factory.setBatchListener(true);

        return factory;
    }*/

    /**
     * 單個(gè)消費(fèi)者
     * @return
     */
    @Bean
    public KafkaListenerContainerFactory<?> singleConsumerFactory() {
        ConcurrentKafkaListenerContainerFactory<Integer, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfigs()));
        factory.setConcurrency(1);
        factory.getContainerProperties().setPollTimeout(3000);
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
        return factory;
    }

    @Bean
    public KafkaConsumer<String,String> customKafkaConsumer(){
        return new KafkaConsumer<String, String>(consumerConfigs());
    }
}

生產(chǎn)者

@Component
public class SendMessageToKafka {

    private static final Logger log = LoggerFactory.getLogger(SendMessageToKafka.class);

    @Resource
    private KafkaTemplate<String,String> kafkaTemplate;


    public void sendToKafka(String message,String topicName){
        try {
            String topicDesc = TopicNameEnum.getDescByTopicName(topicName);
            SendResult<String, String> sendResult = kafkaTemplate.send(topicName, null, message).get(60, TimeUnit.SECONDS);
            if(sendResult.getRecordMetadata() != null){
                log.info("topic對(duì)象【{}】發(fā)送kafka成功,:{},message sent,topic:{}, partition={}, offset={}",topicDesc,sendResult.getProducerRecord().topic(), sendResult.getRecordMetadata().partition(),
                        sendResult.getRecordMetadata().offset());
            }else {
                log.warn("topic對(duì)象【{}】消息發(fā)送kafka失敗,topic:{},進(jìn)行重試。。",topicDesc, topicName);
                // 當(dāng)kafka發(fā)送失敗后,進(jìn)行重試3次。
                FunctionResponse<Boolean> response = FunctionUtil.retryExecute(() -> {
                    try {
                        SendResult<String, String> kafkaResult = kafkaTemplate.send(topicName, null, message).get(60, TimeUnit.SECONDS);
                        if(kafkaResult.getRecordMetadata() != null){
                            return new FunctionResponse<>(true);
                        }else {
                            return new FunctionResponse<>(false);
                        }
                    } catch (Exception e) {
                        return new FunctionResponse<>(false);
                    }
                }, 3, 2000);

                if(!response.getResult()) {
                    log.error("topic對(duì)象【{}】 message sent,topic:{},最終重試發(fā)送失?。簕}", topicDesc,topicName, message);
                }
            }
        } catch (Exception e) {
            log.error("topic【{}】寫(xiě)入kafka失敗,{}",topicName,e);
        }
    }

}

消費(fèi)者

@Component
public class CoreCommandListener {

    private static final Logger log = LoggerFactory.getLogger(CoreCommandListener.class);


    @KafkaListener(topics = {"topic名稱"},containerFactory="singleConsumerFactory")
    private void commandConsumer(ConsumerRecord<Object,String> consumerRecord, Acknowledgment ack){
        DcpMessage message = JSONObject.parseObject(consumerRecord.value(),DcpMessage.class);
        log.info("offset:{},partition:{},消費(fèi)到指令數(shù)據(jù):{}",consumerRecord.offset(),consumerRecord.partition(),message);

        //手動(dòng)提交
        ack.acknowledge();
    }
}

總結(jié)

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

相關(guān)文章

  • Javaweb實(shí)現(xiàn)郵件發(fā)送

    Javaweb實(shí)現(xiàn)郵件發(fā)送

    這篇文章主要為大家詳細(xì)介紹了Javaweb實(shí)現(xiàn)郵件發(fā)送,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2022-06-06
  • 解決RestTemplate加@Autowired注入不了的問(wèn)題

    解決RestTemplate加@Autowired注入不了的問(wèn)題

    這篇文章主要介紹了解決RestTemplate加@Autowired注入不了的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。
    2021-08-08
  • springboot常用注釋的講解

    springboot常用注釋的講解

    今天小編就為大家分享一篇關(guān)于springboot常用注釋的講解,小編覺(jué)得內(nèi)容挺不錯(cuò)的,現(xiàn)在分享給大家,具有很好的參考價(jià)值,需要的朋友一起跟隨小編來(lái)看看吧
    2019-04-04
  • 基于銀河麒麟操作系統(tǒng)安裝Java環(huán)境的詳細(xì)步驟

    基于銀河麒麟操作系統(tǒng)安裝Java環(huán)境的詳細(xì)步驟

    銀河麒麟系統(tǒng)基于Linux系統(tǒng)定制開(kāi)發(fā),部分軟件安裝會(huì)根據(jù)硬件芯片架構(gòu)不同,導(dǎo)致安裝版本不同,這篇文章主要介紹了基于銀河麒麟操作系統(tǒng)安裝Java環(huán)境的詳細(xì)步驟,需要的朋友可以參考下
    2025-06-06
  • SpringBoot學(xué)習(xí)之全局異常處理設(shè)置(返回JSON)

    SpringBoot學(xué)習(xí)之全局異常處理設(shè)置(返回JSON)

    本篇文章主要介紹了SpringBoot學(xué)習(xí)之全局異常處理設(shè)置(返回JSON),小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2018-02-02
  • Activiti7整合Springboot使用記錄

    Activiti7整合Springboot使用記錄

    這篇文章主要介紹了Activiti7+Springboot使用整合記錄,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-08-08
  • Java 詳細(xì)講解分治算法如何實(shí)現(xiàn)歸并排序

    Java 詳細(xì)講解分治算法如何實(shí)現(xiàn)歸并排序

    分治算法的基本思想是將一個(gè)規(guī)模為N的問(wèn)題分解為K個(gè)規(guī)模較小的子問(wèn)題,這些子問(wèn)題相互獨(dú)立且與原問(wèn)題性質(zhì)相同。求出子問(wèn)題的解,就可得到原問(wèn)題的解,本篇文章我們就用分治算法來(lái)實(shí)現(xiàn)歸并排序
    2022-04-04
  • SpringBoot中二級(jí)緩存實(shí)現(xiàn)方案總結(jié)

    SpringBoot中二級(jí)緩存實(shí)現(xiàn)方案總結(jié)

    隨著業(yè)務(wù)的發(fā)展,單一的緩存方案往往無(wú)法同時(shí)兼顧性能、可靠性和一致性等多方面需求,此時(shí),二級(jí)緩存架構(gòu)應(yīng)運(yùn)而生,本文將介紹在Spring Boot中實(shí)現(xiàn)二級(jí)緩存的三種方案,大家可以根據(jù)需要進(jìn)行選擇
    2025-06-06
  • Java如何替換字符

    Java如何替換字符

    文章介紹了Java中String類的replace()方法及其變體replaceFirst()的使用,包括如何替換單個(gè)字符、第一次出現(xiàn)的字符以及多個(gè)字符,通過(guò)示例展示了如何處理字符串中的特殊字符和空格
    2024-11-11
  • JavaCV實(shí)現(xiàn)照片馬賽克效果

    JavaCV實(shí)現(xiàn)照片馬賽克效果

    這篇文章主要介紹了如何通過(guò)JavaCV實(shí)現(xiàn)照片馬賽克效果,文中的示例代碼講解詳細(xì),對(duì)我們學(xué)習(xí)JavaCV有一定的幫助,感興趣的小伙伴可以跟隨小編一起動(dòng)手試一試
    2022-01-01

最新評(píng)論

郑州市| 湘潭县| 武川县| 同德县| 西吉县| 芦山县| 蒲江县| 卢湾区| 将乐县| 庐江县| 万州区| 内江市| 炉霍县| 清新县| 竹北市| 白沙| 昌吉市| 大埔区| 磐安县| 张家口市| 枣庄市| 榆中县| 宝丰县| 南宫市| 中牟县| 桑日县| 喀喇沁旗| 孟村| 长沙市| 泾阳县| 靖宇县| 梁平县| 闻喜县| 淄博市| 宝丰县| 阿拉尔市| 沂源县| 巴青县| 乐东| 绥江县| 许昌市|