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

Kafka Java Producer代碼實例詳解

 更新時間:2020年06月04日 10:05:31   作者:liuming_1992  
這篇文章主要介紹了Kafka Java Producer代碼實例詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下

根據(jù)業(yè)務(wù)需要可以使用Kafka提供的Java Producer API進(jìn)行產(chǎn)生數(shù)據(jù),并將產(chǎn)生的數(shù)據(jù)發(fā)送到Kafka對應(yīng)Topic的對應(yīng)分區(qū)中,入口類為:Producer

Kafka的Producer API主要提供下列三個方法:

  •   public void send(KeyedMessage<K,V> message) 發(fā)送單條數(shù)據(jù)到Kafka集群
  •   public void send(List<KeyedMessage<K,V>> messages) 發(fā)送多條數(shù)據(jù)(數(shù)據(jù)集)到Kafka集群
  •   public void close() 關(guān)閉Kafka連接資源

一、JavaKafkaProducerPartitioner:自定義的數(shù)據(jù)分區(qū)器,功能是:決定輸入的key/value鍵值對的message發(fā)送到Topic的那個分區(qū)中,返回分區(qū)id,范圍:[0,分區(qū)數(shù)量); 這里的實現(xiàn)比較簡單,根據(jù)key中的數(shù)字決定分區(qū)的值。具體代碼如下:

import kafka.producer.Partitioner;
import kafka.utils.VerifiableProperties;

/**
 * Created by gerry on 12/21.
 */
public class JavaKafkaProducerPartitioner implements Partitioner {

  /**
   * 無參構(gòu)造函數(shù)
   */
  public JavaKafkaProducerPartitioner() {
    this(new VerifiableProperties());
  }

  /**
   * 構(gòu)造函數(shù),必須給定
   *
   * @param properties 上下文
   */
  public JavaKafkaProducerPartitioner(VerifiableProperties properties) {
    // nothings
  }

  @Override
  public int partition(Object key, int numPartitions) {
    int num = Integer.valueOf(((String) key).replaceAll("key_", "").trim());
    return num % numPartitions;
  }
}

二、 JavaKafkaProducer:通過Kafka提供的API進(jìn)行數(shù)據(jù)產(chǎn)生操作的測試類;具體代碼如下:

import kafka.javaapi.producer.Producer;
import kafka.producer.KeyedMessage;
import kafka.producer.ProducerConfig;
import org.apache.log4j.Logger;

import java.util.Properties;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.ThreadLocalRandom;

/**
 * Created by gerry on 12/21.
 */
public class JavaKafkaProducer {
  private Logger logger = Logger.getLogger(JavaKafkaProducer.class);
  public static final String TOPIC_NAME = "test";
  public static final char[] charts = "qazwsxedcrfvtgbyhnujmikolp1234567890".toCharArray();
  public static final int chartsLength = charts.length;


  public static void main(String[] args) {
    String brokerList = "192.168.187.149:9092";
    brokerList = "192.168.187.149:9092,192.168.187.149:9093,192.168.187.149:9094,192.168.187.149:9095";
    brokerList = "192.168.187.146:9092";
    Properties props = new Properties();
    props.put("metadata.broker.list", brokerList);
    /**
     * 0表示不等待結(jié)果返回<br/>
     * 1表示等待至少有一個服務(wù)器返回數(shù)據(jù)接收標(biāo)識<br/>
     * -1表示必須接收到所有的服務(wù)器返回標(biāo)識,及同步寫入<br/>
     * */
    props.put("request.required.acks", "0");
    /**
     * 內(nèi)部發(fā)送數(shù)據(jù)是異步還是同步
     * sync:同步, 默認(rèn)
     * async:異步
     */
    props.put("producer.type", "async");
    /**
     * 設(shè)置序列化的類
     * 可選:kafka.serializer.StringEncoder
     * 默認(rèn):kafka.serializer.DefaultEncoder
     */
    props.put("serializer.class", "kafka.serializer.StringEncoder");
    /**
     * 設(shè)置分區(qū)類
     * 根據(jù)key進(jìn)行數(shù)據(jù)分區(qū)
     * 默認(rèn)是:kafka.producer.DefaultPartitioner ==> 安裝key的hash進(jìn)行分區(qū)
     * 可選:kafka.serializer.ByteArrayPartitioner ==> 轉(zhuǎn)換為字節(jié)數(shù)組后進(jìn)行hash分區(qū)
     */
    props.put("partitioner.class", "JavaKafkaProducerPartitioner");

    // 重試次數(shù)
    props.put("message.send.max.retries", "3");

    // 異步提交的時候(async),并發(fā)提交的記錄數(shù)
    props.put("batch.num.messages", "200");

    // 設(shè)置緩沖區(qū)大小,默認(rèn)10KB
    props.put("send.buffer.bytes", "102400");

    // 2. 構(gòu)建Kafka Producer Configuration上下文
    ProducerConfig config = new ProducerConfig(props);

    // 3. 構(gòu)建Producer對象
    final Producer<String, String> producer = new Producer<String, String>(config);

    // 4. 發(fā)送數(shù)據(jù)到服務(wù)器,并發(fā)線程發(fā)送
    final AtomicBoolean flag = new AtomicBoolean(true);
    int numThreads = 50;
    ExecutorService pool = Executors.newFixedThreadPool(numThreads);
    for (int i = 0; i < 5; i++) {
      pool.submit(new Thread(new Runnable() {
        @Override
        public void run() {
          while (flag.get()) {
            // 發(fā)送數(shù)據(jù)
            KeyedMessage message = generateKeyedMessage();
            producer.send(message);
            System.out.println("發(fā)送數(shù)據(jù):" + message);

            // 休眠一下
            try {
              int least = 10;
              int bound = 100;
              Thread.sleep(ThreadLocalRandom.current().nextInt(least, bound));
            } catch (InterruptedException e) {
              e.printStackTrace();
            }
          }

          System.out.println(Thread.currentThread().getName() + " shutdown....");
        }
      }, "Thread-" + i));

    }

    // 5. 等待執(zhí)行完成
    long sleepMillis = 600000;
    try {
      Thread.sleep(sleepMillis);
    } catch (InterruptedException e) {
      e.printStackTrace();
    }
    flag.set(false);

    // 6. 關(guān)閉資源

    pool.shutdown();
    try {
      pool.awaitTermination(6, TimeUnit.SECONDS);
    } catch (InterruptedException e) {
    } finally {
      producer.close(); // 最后之后調(diào)用
    }
  }

  /**
   * 產(chǎn)生一個消息
   *
   * @return
   */
  private static KeyedMessage<String, String> generateKeyedMessage() {
    String key = "key_" + ThreadLocalRandom.current().nextInt(10, 99);
    StringBuilder sb = new StringBuilder();
    int num = ThreadLocalRandom.current().nextInt(1, 5);
    for (int i = 0; i < num; i++) {
      sb.append(generateStringMessage(ThreadLocalRandom.current().nextInt(3, 20))).append(" ");
    }
    String message = sb.toString().trim();
    return new KeyedMessage(TOPIC_NAME, key, message);
  }

  /**
   * 產(chǎn)生一個給定長度的字符串
   *
   * @param numItems
   * @return
   */
  private static String generateStringMessage(int numItems) {
    StringBuilder sb = new StringBuilder();
    for (int i = 0; i < numItems; i++) {
      sb.append(charts[ThreadLocalRandom.current().nextInt(chartsLength)]);
    }
    return sb.toString();
  }
}

三、Pom.xml依賴配置如下

<properties>
  <kafka.version>0.8.2.1</kafka.version>
</properties>

<dependencies>
  <dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka_2.10</artifactId>
    <version>${kafka.version}</version>
  </dependency>
</dependencies>

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

相關(guān)文章

  • Java解析變量公式的簡單示例

    Java解析變量公式的簡單示例

    在Java編程中,經(jīng)常會遇到需要解析表達(dá)式或公式的情況,特別是涉及到動態(tài)計算或配置項的場景,在本篇文章中,我將介紹如何在Java中解析變量公式,并給出一個簡單的實現(xiàn)示例,需要的朋友可以參考下
    2024-10-10
  • Java縮小文件內(nèi)存占用的方法技巧分享

    Java縮小文件內(nèi)存占用的方法技巧分享

    在Java應(yīng)用程序中,處理大文件時經(jīng)常會遇到內(nèi)存占用過高的問題,為了縮小文件的內(nèi)存占用,我們可以采取一些有效的方法來優(yōu)化和管理內(nèi)存的使用,本文將介紹一些在Java中縮小文件內(nèi)存占用的技巧,需要的朋友可以參考下
    2024-10-10
  • 如何解決shardingsphere報錯Missing?the?data?source?name:‘null‘

    如何解決shardingsphere報錯Missing?the?data?source?name:‘null‘

    使用ShardingSphere進(jìn)行分庫操作時,如果遇到“Missing?the?datasource?name:?‘null’”的錯誤,通常是因為所操作的表沒有配置相關(guān)的路由信息,例如,如果在properties中僅配置了health_record和health_task的路由規(guī)則
    2024-11-11
  • java swing實現(xiàn)貪吃蛇雙人游戲

    java swing實現(xiàn)貪吃蛇雙人游戲

    這篇文章主要為大家詳細(xì)介紹了java swing實現(xiàn)貪吃蛇雙人小游戲,文中示例代碼介紹的非常詳細(xì),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-01-01
  • springBoot項目如何實現(xiàn)啟動多個實例

    springBoot項目如何實現(xiàn)啟動多個實例

    這篇文章主要介紹了springBoot項目如何實現(xiàn)啟動多個實例的操作,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • Java編寫時間工具類ZTDateTimeUtil的示例代碼

    Java編寫時間工具類ZTDateTimeUtil的示例代碼

    這篇文章主要為大家詳細(xì)介紹了如何利用Java編寫時間工具類ZTDateTimeUtil,文中的示例代碼講解詳細(xì),有需要的小伙伴可以跟隨小編一起學(xué)習(xí)一下
    2023-11-11
  • Java獲取XML節(jié)點總結(jié)之讀取XML文檔節(jié)點的方法

    Java獲取XML節(jié)點總結(jié)之讀取XML文檔節(jié)點的方法

    下面小編就為大家?guī)硪黄狫ava獲取XML節(jié)點總結(jié)之讀取XML文檔節(jié)點的方法。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2016-10-10
  • java中redissonClient 分布式鎖的使用

    java中redissonClient 分布式鎖的使用

    在集群的情況下,用戶多次請求接口時,存入的內(nèi)容可能會導(dǎo)致重復(fù),這時候就可以使用分布式鎖來限制,本文就來介紹一下java中redissonClient 分布式鎖的使用,感興趣的可以了解一下
    2024-03-03
  • spring boot admin 搭建詳解

    spring boot admin 搭建詳解

    本篇文章主要介紹了spring boot admin 搭建詳解,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-04-04
  • 95%的Java程序員人都用不好Synchronized詳解

    95%的Java程序員人都用不好Synchronized詳解

    這篇文章主要為大家介紹了95%的Java程序員人都用不好Synchronized詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-03-03

最新評論

任丘市| 万山特区| 博客| 古浪县| 巫山县| 定州市| 宁远县| 尖扎县| 高淳县| 扶余县| 定南县| 江西省| 宝丰县| 连山| 邢台市| 察雅县| 斗六市| 富源县| 德兴市| 娱乐| 平江县| 吉首市| 延庆县| 山东省| 开远市| 徐水县| 安义县| 嵩明县| 天气| 衢州市| 丹江口市| 塔河县| 博野县| 新昌县| 太保市| 广西| 凤城市| 汪清县| 太仆寺旗| 五常市| 铁岭市|