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

Kafka producer端開發(fā)代碼實(shí)例

 更新時(shí)間:2020年11月11日 15:52:21   作者:碼農(nóng)大衛(wèi)  
這篇文章主要介紹了Kafka producer端開發(fā)代碼實(shí)例,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下

一、producer工作流程

  producer使用用戶啟動(dòng)producer的線程,將待發(fā)送的消息封裝到一個(gè)ProducerRecord類實(shí)例,然后將其序列化之后發(fā)送給partitioner,再由后者確定目標(biāo)分區(qū)后一同發(fā)送到位于producer程序中的一塊內(nèi)存緩沖區(qū)中。而producer的另外一個(gè)線程(Sender線程)則負(fù)責(zé)實(shí)時(shí)從該緩沖區(qū)中提取出準(zhǔn)備就緒的消息封裝進(jìn)一個(gè)批次(batch),統(tǒng)一發(fā)送給對(duì)應(yīng)的broker,具體流程如下圖:

二、producer示例程序開發(fā)

  首先引入kafka相關(guān)依賴,在pom.xml文件中加入如下依賴:

<!--kafka-->
  <dependency>
   <groupId>org.apache.kafka</groupId>
   <artifactId>kafka_2.12</artifactId>
   <version>2.2.0</version>
  </dependency>

  在resources下面創(chuàng)建kafka-producer.properties配置文件,用于設(shè)置kafka參數(shù),內(nèi)容如下:

bootstrap.servers=192.168.184.128:9092,192.168.184.128:9093,192.168.184.128:9094
key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer
acks=-1
retries=3
batch.size=323840
linger.ms=10
buffer.memory=33554432
max.block.ms=3000

  其中,前三個(gè)參數(shù)必須明確指定,因?yàn)檫@三個(gè)參數(shù)沒有默認(rèn)值(注:kafka的producer參數(shù)配置可以參考http://kafka.apache.org/documentation/),然后編寫producer發(fā)送消息的代碼:

/**
   * Kafka發(fā)送消息測(cè)試
   * @throws IOException
   */
  public void sendMsg() throws IOException {
    //1.構(gòu)造properties對(duì)象
    Properties properties = new Properties();
    FileInputStream fileInputStream = new FileInputStream("F:\\javaCode\\jvmdemo\\src\\main\\resources\\kafka-producer.properties");
    properties.load(fileInputStream);
    fileInputStream.close();
    //2.構(gòu)造kafkaProducer對(duì)象
    KafkaProducer producer = new KafkaProducer(properties);
    for (int i = 0; i < 100; i++) {
      //3.構(gòu)造待發(fā)送消息的producerRecord對(duì)象,并指定消息要發(fā)送到哪個(gè)topic,消息的key和value
      ProducerRecord testTopic = new ProducerRecord("testTopic", Integer.toString(i), Integer.toString(i));
      //4.調(diào)用kafkaProducer對(duì)象的send方法發(fā)送消息
      producer.send(testTopic);
    }
    //5.關(guān)閉kafkaProducer
    producer.close();
  }

  然后登陸kafka所在服務(wù)器,執(zhí)行以下命令監(jiān)聽消息: 

cd /usr/local/kafka/bin
./kafka-console-consumer.sh --bootstrap-server 192.168.184.128:9092,192.168.184.128:9093,192.168.184.128:9094 --topic testTopic --from-beginning

  運(yùn)行sendMsg方法,注意觀察消費(fèi)端,

  

  可以看到有0-99之間的數(shù)字依次被消費(fèi)到,說明消息發(fā)送成功。

三、異步和同步發(fā)送消息

  上面發(fā)送消息的示例程序中,沒有對(duì)發(fā)送結(jié)果進(jìn)行處理,如果消息發(fā)送失敗我們也是無法得知的,這種方法在實(shí)際應(yīng)用中是不推薦的。在實(shí)際使用場(chǎng)景中,一般使用異步和同步兩種常見發(fā)送方式。Java版本producer的send方法會(huì)返回一個(gè)Future對(duì)象,如果調(diào)用Future.get()方法就會(huì)無限等待返回結(jié)果,實(shí)現(xiàn)同步發(fā)送的效果,否則就是異步發(fā)送。

  1.異步發(fā)送消息

  Java版本producer的send()方法提供了回調(diào)類參數(shù)來實(shí)現(xiàn)異步發(fā)送以及對(duì)發(fā)送結(jié)果進(jìn)行的響應(yīng),具體代碼如下:

/**
   * 異步發(fā)送消息
   *
   * @throws IOException
   */
  public void sendMsg() throws IOException {
    //1.構(gòu)造properties對(duì)象
    Properties properties = new Properties();
    FileInputStream fileInputStream = new FileInputStream("F:\\javaCode\\jvmdemo\\src\\main\\resources\\kafka-producer.properties");
    properties.load(fileInputStream);
    fileInputStream.close();
    //2.構(gòu)造kafkaProducer對(duì)象
    KafkaProducer producer = new KafkaProducer(properties);
    for (int i = 0; i < 100; i++) {
      //3.構(gòu)造待發(fā)送消息的producerRecord對(duì)象,并指定消息要發(fā)送到哪個(gè)topic,消息的key和value
      ProducerRecord testTopic = new ProducerRecord("testTopic", Integer.toString(i), Integer.toString(i));
      //4.調(diào)用kafkaProducer對(duì)象的send方法發(fā)送消息,傳入Callback回調(diào)參數(shù)
      producer.send(testTopic, new Callback() {
        @Override
        public void onCompletion(RecordMetadata recordMetadata, Exception exception) {
          if (null == exception) {
            //消息發(fā)送成功后的處理
            System.out.println("消息發(fā)送成功");
          } else {
            //消息發(fā)送失敗后的處理
            System.out.println("消息發(fā)送失敗");
          }
        }
      });
    }
    //5.關(guān)閉kafkaProducer
    producer.close();
  }

  以上代碼中,send方法第二個(gè)參數(shù)傳入一個(gè)匿名內(nèi)部類對(duì)象,也可以傳入實(shí)現(xiàn)了org.apache.kafka.clients.producer.Callback接口的類對(duì)象。同時(shí)onCompletion方法的兩個(gè)入?yún)ecordMetadata和exception不會(huì)同時(shí)為空,當(dāng)消息發(fā)送成功后,exception為null,消息發(fā)送失敗后recordMetadata為null。因此可以按照兩個(gè)入?yún)⑦M(jìn)行成功和失敗邏輯的處理。

  其次,Kafka發(fā)送消息失敗的類型包含兩類,可重試異常和不可重試異常。所有的可重試異常都繼承自org.apache.kafka.common.errors.RetriableException抽象類,理論上所有沒有繼承RetriableException 類的其他異常都屬于不可重試異常,鑒于此,可以在消息發(fā)送失敗后,按照是否可以重試,來進(jìn)行不同的處理邏輯處理:

//4.調(diào)用kafkaProducer對(duì)象的send方法發(fā)送消息
producer.send(testTopic, new Callback() {
  @Override
  public void onCompletion(RecordMetadata recordMetadata, Exception exception) {
    if (null == exception) {
      //消息發(fā)送成功后的處理
      System.out.println("消息發(fā)送成功");
    } else {
      if(exception instanceof RetriableException){
        // 可重試異常
        System.out.println("可重試異常");
      }else{
        // 不可重試異常
        System.out.println("不可重試異常");
      }
    }
  }
});

  2.同步發(fā)送消息

  同步發(fā)送和異步發(fā)送是通過Java的Futrue來區(qū)分的,調(diào)用Future.get()無限等待結(jié)果返回,即實(shí)現(xiàn)了同步發(fā)送的結(jié)果,具體代碼如下:

// 發(fā)送消息
 Future future = producer.send(testTopic);
 try {
   // 調(diào)用get方法等待結(jié)果返回,發(fā)送失敗則會(huì)拋出異常
   future.get();
 } catch (Exception e) {
   System.out.println("消息發(fā)送失敗");
 }

四、其他高級(jí)特性

1.消息分區(qū)機(jī)制

  kafka producer提供了分區(qū)策略以及分區(qū)器(partitioner)用于確定將消息發(fā)送到指定topic的哪個(gè)分區(qū)中。默認(rèn)分區(qū)器根據(jù)murmur2算法計(jì)算消息key的哈希值,然后對(duì)總分區(qū)數(shù)求模確認(rèn)消息要被發(fā)送的目標(biāo)分區(qū)號(hào)(這點(diǎn)讓我想起了redis集群中key值的分配方法),這樣就確保了相同key的消息被發(fā)送到相同分區(qū)。若消息沒有key值,將采用輪詢的方式確保消息在topic的所有分區(qū)上均勻分配。

  除了使用kafka默認(rèn)的分區(qū)機(jī)制,也可以通過實(shí)現(xiàn)org.apache.kafka.clients.producer.Partitioner接口來自定義分區(qū)器,此時(shí)需要在構(gòu)造KafkaProducer的 properties中增加partitioner.class來指明分區(qū)器實(shí)現(xiàn)類,如:partitioner.class=com.demo.service.CustomerPartitioner。

2.消息序列化

  在本篇開始的producer示例程序中,在構(gòu)造KafkaProducer對(duì)象的時(shí)候,有兩個(gè)配置項(xiàng)

  • key.serializer=org.apache.kafka.common.serialization.StringSerializer
  • value.serializer=org.apache.kafka.common.serialization.StringSerializer

分別用于配置消息key和value的序列化方式為String類型,除此之外,Kafka中還提供了如下默認(rèn)的序列化器:

  ByteArraySerializer:本質(zhì)上什么也不做,因?yàn)榫W(wǎng)絡(luò)中傳輸就是以字節(jié)傳輸?shù)模?/p>

  ByteBufferSerializer:序列化ByteBuffer消息;

  BytesSerializer:序列化kafka自定義的Bytes類型;

  IntegerSerializer:序列化Integer類型;

  DoubleSerializer:序列化Double類型;

  LongSerializer:序列化Long類型;

  如果要自定義序列化器,則需要實(shí)現(xiàn)org.apache.kafka.common.serialization.Serializer接口,并且將key.serializer和value.serializer配置為自定義的序列化器。

3.消息壓縮

  消息壓縮可以顯著降低磁盤占用以及帶寬占用,從而有效提升I/O密集型應(yīng)用性能,但是引入壓縮同時(shí)會(huì)消耗額外的CPU,因此壓縮是I/O性能和CPU資源的平衡。kafka目前支持3種壓縮算法:CZIP,Snappy和LZ4,性能測(cè)試結(jié)果顯示三種壓縮算法的性能如下:LZ4>>Snappy>GZIP,目前啟用LZ4進(jìn)行消息壓縮的producer的吞吐量是最高的。

  默認(rèn)情況下Kafka是不壓縮消息的,但是可以通過在創(chuàng)建KafkaProducer 對(duì)象的時(shí)候設(shè)置producer端參數(shù)compression.type來開啟消息壓縮,如配置compression.type=LZ4。那么什么時(shí)候開啟壓縮呢?首先判斷是否啟用壓縮的依據(jù)是I/O資源消耗與CPU資源消耗的對(duì)比,如果環(huán)境上I/O資源非常緊張,比如producer程序占用了大量的網(wǎng)絡(luò)帶寬或broker端的磁盤占用率很高,而producer端的CPU資源非常富裕,那么就可以考慮為producer開啟壓縮。

4.無消息丟失配置

  在使用KafkaProducer.send()方法發(fā)送消息的時(shí)候,其實(shí)是把消息放入緩沖區(qū)中,再由一個(gè)專屬I/O線程負(fù)責(zé)從緩沖區(qū)提取消息并封裝消息到batch中,然后再發(fā)送出去。如果在I/O線程將消息發(fā)送出去之前,producer奔潰了,那么所有的消息都將丟失。同時(shí),存在多消息發(fā)送時(shí)候由于網(wǎng)絡(luò)抖動(dòng)導(dǎo)致消息亂序的問題,為了解決這兩個(gè)問題,可以通過在producer端以及broker端進(jìn)行配置進(jìn)行避免。

4.1 producer端配置

  max.block.ms=3000:設(shè)置block的時(shí)長(zhǎng),當(dāng)緩沖區(qū)被填滿或者metadata丟失時(shí)產(chǎn)生block,停止接收新的消息;

  acks=all:等待所有follower都響應(yīng)了發(fā)送消息認(rèn)為消息發(fā)送成功;

  retries=Integer.MAX_VALUE:設(shè)置重試次數(shù),設(shè)置一個(gè)比較大的值可以保證消息不丟失;

  max.in.flight.requests.per.connection=1:限制producer在單個(gè)broker連接上能夠發(fā)送的未響應(yīng)請(qǐng)求的數(shù)量,從而防止同topic統(tǒng)一分區(qū)下消息亂序問題;

  除了設(shè)置以上參數(shù)之外,在發(fā)送消息的時(shí)候,應(yīng)該盡量使用帶有回調(diào)參數(shù)的send方法來處理發(fā)送結(jié)果,如果數(shù)據(jù)發(fā)送失敗,則顯示調(diào)用KafkaProducer.close(0)方法來立即關(guān)閉producer,防止消息亂序。

4.2 broker端配置

  unclean.leader.election.enable=false:關(guān)閉unclean leader選舉,即不允許非ISR中的副本被選舉為leader;

  replication.factor>=3:至少使用3個(gè)副本保存數(shù)據(jù);

  min.issync.replicas>1:控制某條消息至少被寫入到ISR中多少個(gè)副本才算成功,當(dāng)且僅當(dāng)producer端acks參數(shù)設(shè)置為all或者-1時(shí),該參數(shù)才有效。

  最后,確保replication.factor>min.issync.replicas,如果兩者相等,那么只要有一個(gè)副本掛掉,分區(qū)就無法工作,推薦配置replication.factor=min.issync.replicas+1。

  關(guān)于producer端的開發(fā)就介紹到這兒,下一篇將介紹consumer端的開發(fā)。

以上就是本文的全部?jī)?nèi)容,希望對(duì)大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。

相關(guān)文章

  • Mybatis動(dòng)態(tài)元素if的使用方式

    Mybatis動(dòng)態(tài)元素if的使用方式

    這篇文章主要介紹了Mybatis動(dòng)態(tài)元素if的使用方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-12-12
  • Spring的Aware接口實(shí)現(xiàn)及執(zhí)行順序詳解

    Spring的Aware接口實(shí)現(xiàn)及執(zhí)行順序詳解

    這篇文章主要為大家介紹了Spring的Aware接口實(shí)現(xiàn)及執(zhí)行順序詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-12-12
  • springcloud整合gateway實(shí)現(xiàn)網(wǎng)關(guān)的示例代碼

    springcloud整合gateway實(shí)現(xiàn)網(wǎng)關(guān)的示例代碼

    本文主要介紹了springcloud整合gateway實(shí)現(xiàn)網(wǎng)關(guān)的示例代碼,文中通過示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2022-01-01
  • Java的Swing編程中使用SwingWorker線程模式及頂層容器

    Java的Swing編程中使用SwingWorker線程模式及頂層容器

    這篇文章主要介紹了在Java的Swing編程中使用SwingWorker線程模式及頂層容器的方法,適用于客戶端圖形化界面軟件的開發(fā),需要的朋友可以參考下
    2016-01-01
  • JDK 1.8 安裝配置教程(win7 64bit )

    JDK 1.8 安裝配置教程(win7 64bit )

    這篇文章主要為大家詳細(xì)介紹了win7 64bit下JDK 1.8 安裝配置教程,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-08-08
  • java自帶的四種線程池實(shí)例詳解

    java自帶的四種線程池實(shí)例詳解

    java線程的創(chuàng)建非常昂貴,需要JVM和OS(操作系統(tǒng))互相配合完成大量的工作,下面這篇文章主要給大家介紹了關(guān)于java自帶的四種線程池的相關(guān)資料,文中通過圖文介紹的非常詳細(xì),需要的朋友可以參考下
    2022-04-04
  • javaweb 國(guó)際化:DateFormat,NumberFormat,MessageFormat,ResourceBundle的使用

    javaweb 國(guó)際化:DateFormat,NumberFormat,MessageFormat,ResourceBu

    本文主要介紹javaWEB國(guó)際化的知識(shí),這里整理了詳細(xì)的資料及實(shí)現(xiàn)代碼,有興趣的小伙伴可以參考下
    2016-09-09
  • MyBatis-Plus實(shí)現(xiàn)字段自動(dòng)填充功能的示例

    MyBatis-Plus實(shí)現(xiàn)字段自動(dòng)填充功能的示例

    本文主要介紹了MyBatis-Plus實(shí)現(xiàn)字段自動(dòng)填充功能的示例,文中通過示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2021-11-11
  • Java實(shí)現(xiàn)度分秒坐標(biāo)轉(zhuǎn)十進(jìn)制度

    Java實(shí)現(xiàn)度分秒坐標(biāo)轉(zhuǎn)十進(jìn)制度

    隨著技術(shù)的發(fā)展,十進(jìn)制度因其精確性和便捷性在現(xiàn)代應(yīng)用中越來越受到青睞,下面我們就來看看如何使用Java實(shí)現(xiàn)度分秒坐標(biāo)轉(zhuǎn)十進(jìn)制度吧
    2024-12-12
  • 使用MDC實(shí)現(xiàn)日志鏈路跟蹤

    使用MDC實(shí)現(xiàn)日志鏈路跟蹤

    這篇文章主要介紹了使用MDC實(shí)現(xiàn)日志鏈路跟蹤,在微服務(wù)環(huán)境中,我們經(jīng)常使用Skywalking、CAT等去實(shí)現(xiàn)整體請(qǐng)求鏈路的追蹤,但是這個(gè)整體運(yùn)維成本高,架構(gòu)復(fù)雜,我們來使用MDC通過Log來實(shí)現(xiàn)一個(gè)輕量級(jí)的會(huì)話事務(wù)跟蹤功能,下面就來看看具體的過程吧,需要的朋友可以參考一下
    2022-01-01

最新評(píng)論

南平市| 灵石县| 锡林浩特市| 临清市| 安丘市| 明光市| 随州市| 新安县| 济阳县| 城市| 乌兰察布市| 响水县| 井冈山市| 漳平市| 邯郸市| 黄骅市| 瓮安县| 蒙城县| 平遥县| 宣威市| 华容县| 从化市| 邵阳市| 静宁县| 广汉市| 桦川县| 林西县| 岳池县| 阳朔县| 铅山县| 墨江| 江津市| 宿州市| 永修县| 纳雍县| 拉萨市| 丽水市| 阿拉善左旗| 灌阳县| 阿克陶县| 修武县|