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

普通java項(xiàng)目集成kafka方式

 更新時(shí)間:2024年11月28日 11:05:03   作者:西柚感覺日了狗  
文章介紹了如何在非Spring Cloud或Spring Boot項(xiàng)目中配置和使用Kafka,提供了一個(gè)簡(jiǎn)單的Kafka配置讀取類,可以靈活地從不同配置中讀取屬性,并提供默認(rèn)值

現(xiàn)在假設(shè)一種需求,我方業(yè)務(wù)系統(tǒng)要與某服務(wù)平臺(tái)通過(guò)kafka交互,異步獲取服務(wù),而系統(tǒng)架構(gòu)可能老舊,不是spring cloud桶,不是spring boot,只是java普通項(xiàng)目或者 java web項(xiàng)目

依賴

		<dependency>
            <groupId>org.apache.kafka</groupId>
            <artifactId>kafka-clients</artifactId>
            <version>3.1.0</version>
        </dependency>

        <dependency>
            <groupId>org.apache.kafka</groupId>
            <artifactId>kafka_2.11</artifactId>
            <version>2.4.1</version>
        </dependency>

Kafka配置讀取類

本文后邊沒用到,直接填配置了,簡(jiǎn)單點(diǎn)

但如果生產(chǎn)需要,還是有這個(gè)類比較好,可以從不同配置中讀取,同時(shí)給個(gè)默認(rèn)值

import org.apache.commons.lang3.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.File;
import java.io.FileInputStream;
import java.io.IOException;
import java.util.Properties;

/**
 * kafka配置讀取類
 *
 * @author zzy
 */
public class KafkaProperties {

    private static final Logger LOG = LoggerFactory.getLogger(KafkaProperties.class);

    private static Properties serverProps = new Properties();

    private static Properties clientProps = new Properties();

    private static Properties producerProps = new Properties();

    private static Properties consumerProps = new Properties();

    private static KafkaProperties instance = null;

    private KafkaProperties() {

        String filePath = System.getProperty("user.dir") + File.separator
                + "kafkaConf" + File.separator;

        File file;
        FileInputStream fis = null;
        try {
            file = new File(filePath + "producer.properties");
            if (file.exists()) {
                fis = new FileInputStream(filePath + "producer.properties");
                producerProps.load(fis);
            }

            file = new File(filePath + "consumer.properties");
            if (file.exists()) {
                fis = new FileInputStream(filePath + "consumer.properties");
                consumerProps.load(fis);
            }

            file = new File(filePath + "server.properties");
            if (file.exists()) {
                fis = new FileInputStream(filePath + "server.properties");
                serverProps.load(fis);
            }

            file = new File(filePath + "client.properties");
            if (file.exists()) {
                fis = new FileInputStream(filePath + "client.properties");
                clientProps.load(fis);
            }

        } catch (Exception e) {

            LOG.error("init kafka props error." + e.getMessage());

        } finally {

            if (fis != null) {
                try {
                    fis.close();
                } catch (IOException e) {
                    LOG.error("close kafka properties fis error." + e);
                }
            }

        }

    }

    /**
     * 獲取懶漢式單例
     */
    public static synchronized KafkaProperties getInstance() {
        if (instance == null) {
            instance = new KafkaProperties();
        }

        return instance;
    }


    /**
     * 獲取配置,獲取不到時(shí)使用參數(shù)的默認(rèn)配置
     */
    public String getValue(String key, String defaultValue) {
        String value;

        if (StringUtils.isEmpty(key)) {
            LOG.error("key is null or empty");
        }
        value = getPropsValue(key);

        if (value == null) {
            LOG.warn("kafka property getValue return null, the key is " + key);
            value = defaultValue;
        }
        LOG.info("kafka property getValue, key:" + key + ", value:" + value);

        return value;
    }


    private String getPropsValue(String key) {
        String value = serverProps.getProperty(key);

        if (value == null) {
            value = producerProps.getProperty(key);
        }

        if (value == null) {
            value = consumerProps.getProperty(key);
        }

        if (value == null) {
            value = clientProps.getProperty(key);
        }

        return value;
    }
}

producer

import org.apache.kafka.clients.producer.Callback;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Properties;
import java.util.concurrent.ExecutionException;

/**
 * kafka producer
 * @author zzy
 */
public class KafkaProFormal {

    public static final Logger LOG = LoggerFactory.getLogger(KafkaProFormal.class);

    private Properties properties = new Properties();

    private final String bootstrapServers = "bootstrap.servers";
    private final String clientId = "client.id";
    private final String keySerializer = "key.serializer";
    private final String valueSerializer = "value.serializer";
    //private final String securityProtocol = "security.protocol";
    //private final String saslKerberosServiceName = "sasl.kerberos.service.name";
    //private final String kerberosDomainName = "kerberos.domain.name";
    private final String maxRequestSize = "max.request.size";

    private KafkaProducer<String, String> producer;

    private volatile static KafkaProFormal kafkaProFormal;

    private KafkaProFormal(String servers) {
        properties.put(bootstrapServers, servers);
        properties.put(keySerializer, "org.apache.kafka.common.serialization.StringSerializer");
        properties.put(valueSerializer, "org.apache.kafka.common.serialization.StringSerializer");

        producer = new KafkaProducer<String, String>(properties);
    }

    public static KafkaProFormal getInstance(String servers) {
        if(kafkaProFormal == null) {
            synchronized(KafkaProFormal.class) {
                if(kafkaProFormal == null) {
                    kafkaProFormal = new KafkaProFormal(servers);
                }
            }
        }

        return kafkaProFormal;
    }

    public void sendStringWithCallBack(String topic, String message, boolean asyncFlag) {
        ProducerRecord<String, String> record = new ProducerRecord<String, String>(topic, message);
        long startTime = System.currentTimeMillis();
        if(asyncFlag) {
            //異步發(fā)送
            producer.send(record, new KafkaCallBack(startTime, message));
        } else {
            //同步發(fā)送
            try {
                producer.send(record, new KafkaCallBack(startTime, message)).get();
            } catch (InterruptedException e) {
                LOG.error("InterruptedException occured : {0}", e);
            } catch (ExecutionException e) {
                LOG.error("ExecutionException occured : {0}", e);
            }
        }
    }
}

class KafkaCallBack implements Callback {
    private static Logger LOG = LoggerFactory.getLogger(KafkaCallBack.class);

    private String key;

    private long startTime;

    private String message;

    KafkaCallBack(long startTime, String message) {
        this.startTime = startTime;
        this.message = message;
    }

    @Override
    public void onCompletion(RecordMetadata metadata, Exception exception) {
        long elapsedTime = System.currentTimeMillis() - startTime;

        if(metadata != null) {
            LOG.info("Record(" + key + "," + message + ") sent to partition(" + metadata.partition()
                    + "), offset(" + metadata.offset() + ") in " + elapsedTime + " ms.");
        } else {
            LOG.error("metadata is null." + "Record(" + key + "," + message + ")", exception);
        }
    }
}

consumer

import kafka.utils.ShutdownableThread;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.time.Duration;
import java.util.Properties;
import java.util.Set;

/**
 * kafka consumer
 * @author zzy
 */
public abstract class KafkaConFormal extends ShutdownableThread {

    private static final Logger LOG = LoggerFactory.getLogger(KafkaConFormal.class);

    private Set<String> topics;

    private final String bootstrapServers = "bootstrap.servers";
    private final String groupId = "group.id";
    private final String keyDeserializer = "key.deserializer";
    private final String valueDeserializer = "value.deserializer";
    private final String enableAutoCommit = "enable.auto.commit";
    private final String autoCommitIntervalMs = "auto.commit.interval.ms";
    private final String sessionTimeoutMs = "session.timeout.ms";

    private KafkaConsumer<String, String> consumer;

    public KafkaConFormal(String topic) {
        super("KafkaConsumerExample", false);

        topics.add(topic);

        Properties props = new Properties();
        props.put(bootstrapServers, "your servers");
        props.put(groupId, "TestGroup");
        props.put(enableAutoCommit, "true");
        props.put(autoCommitIntervalMs, "1000");
        props.put(sessionTimeoutMs, "30000");
        props.put(keyDeserializer, "org.apache.kafka.common.serialization.StringDeserializer");
        props.put(valueDeserializer, "org.apache.kafka.common.serialization.StringDeserializer");

        consumer = new KafkaConsumer<>(props);
    }

    /**
     * subscribe and handle the msg
     */
    @Override
    public void doWork() {
        consumer.subscribe(topics);
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));

        dealRecords(records);
    }

    /**
     * 實(shí)例化consumer時(shí),進(jìn)行對(duì)消費(fèi)信息的處理
     * @param records records
     */
    public abstract void dealRecords(ConsumerRecords<String, String> records);

    public void setTopics(Set<String> topics) {
        this.topics = topics;
    }
}

使用

KafkaProFormal producer = KafkaProFormal.getInstance("kafka server1.1.1.1:9092,2.2.2.2:9092");

KafkaConFormal consumer = new KafkaConFormal("consume_topic") {
	@Override
	public void dealRecords(ConsumerRecords<String, String> records) {
		for (ConsumerRecord<String, String> record: records) {
			producer.sendStringWithCallBack("target_topic", record.value(), true);
    	}
	}
};

consumer.start();

總結(jié)

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

相關(guān)文章

  • 淺談MultipartFile中transferTo方法的坑

    淺談MultipartFile中transferTo方法的坑

    這篇文章主要介紹了MultipartFile中transferTo方法的坑,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-07-07
  • Java鎖競(jìng)爭(zhēng)導(dǎo)致sql慢日志原因分析

    Java鎖競(jìng)爭(zhēng)導(dǎo)致sql慢日志原因分析

    這篇文章主要介紹了Java鎖競(jìng)爭(zhēng)導(dǎo)致sql慢的日志原因分析,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)吧
    2022-11-11
  • Mybatis mapper標(biāo)簽中配置子標(biāo)簽package的坑及解決

    Mybatis mapper標(biāo)簽中配置子標(biāo)簽package的坑及解決

    這篇文章主要介紹了Mybatis mapper標(biāo)簽中配置子標(biāo)簽package的坑及解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • Spring Boot REST國(guó)際化的實(shí)現(xiàn)代碼

    Spring Boot REST國(guó)際化的實(shí)現(xiàn)代碼

    本文我們將討論如何在現(xiàn)有的Spring Boot項(xiàng)目中添加國(guó)際化。只需幾個(gè)簡(jiǎn)單的步驟即可實(shí)現(xiàn)Spring Boot應(yīng)用的國(guó)際化,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-10-10
  • java分布式面試降級(jí)組件Hystrix的功能特性

    java分布式面試降級(jí)組件Hystrix的功能特性

    這篇文章主要為大家介紹了java分布式面試關(guān)于降級(jí)組件Hystrix的功能特性回答,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步
    2022-03-03
  • 詳細(xì)說(shuō)一說(shuō)Java序列化的幾種方式對(duì)比

    詳細(xì)說(shuō)一說(shuō)Java序列化的幾種方式對(duì)比

    在某種情況下需要考慮一些安全問題和數(shù)據(jù)對(duì)象的使用問題,這時(shí)候就可以用序列化的技術(shù)來(lái)進(jìn)行數(shù)據(jù)存儲(chǔ),這篇文章主要介紹了Java序列化幾種方式對(duì)比的相關(guān)資料,需要的朋友可以參考下
    2025-07-07
  • Java如何解決ArrayList的并發(fā)問題

    Java如何解決ArrayList的并發(fā)問題

    這篇文章主要介紹了Java如何解決ArrayList的并發(fā)問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2025-04-04
  • Java中BIO、NIO和AIO的區(qū)別、原理與用法

    Java中BIO、NIO和AIO的區(qū)別、原理與用法

    這篇文章主要介紹了Java中BIO、NIO和AIO的區(qū)別、原理與用法,文中通過(guò)示例代碼介紹的非常詳細(xì)。對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-12-12
  • Java正則之貪婪匹配、惰性匹配

    Java正則之貪婪匹配、惰性匹配

    這篇文章主要介紹了Java正則之貪婪匹配、惰性匹配的相關(guān)資料,需要的朋友可以參考下
    2015-03-03
  • Java異常跟蹤棧定義與用法示例

    Java異常跟蹤棧定義與用法示例

    這篇文章主要介紹了Java異常跟蹤棧定義與用法,結(jié)合具體實(shí)例形式分析了異常處理?xiàng)5母拍?、原理及相關(guān)使用技巧,需要的朋友可以參考下
    2018-05-05

最新評(píng)論

简阳市| 五指山市| 唐山市| 阿巴嘎旗| 册亨县| 苏州市| 普格县| 罗定市| 卓尼县| 留坝县| 汪清县| 密山市| 云梦县| 淅川县| 屯昌县| 满城县| 红河县| 民县| 台北县| 晋江市| 射阳县| 双鸭山市| 望城县| 临武县| 黑山县| 南皮县| 广河县| 定陶县| 金秀| 江川县| 遂昌县| 封开县| 丽水市| 南江县| 屏东市| 荆州市| 丁青县| 莫力| 鹤山市| 巨鹿县| 印江|