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

Kafka簡單客戶端編程實例

 更新時間:2017年11月08日 11:33:31   作者:liuyazhuang  
這篇文章主要為大家詳細(xì)介紹了Kafka簡單客戶端編程實例,利用Kafka的API進(jìn)行客戶端編程,具有一定的參考價值,感興趣的小伙伴們可以參考一下

今天,我們給大家?guī)硪黄绾卫肒afka的API進(jìn)行客戶端編程的文章,這篇文章很簡單,就是利用Kafka的API創(chuàng)建一個生產(chǎn)者和消費(fèi)者,生產(chǎn)者不斷向Kafka寫入消息,消費(fèi)者則不斷消費(fèi)Kafka的消息。下面是具體的實例代碼。

一、創(chuàng)建配置類Config

這個類很簡單,只是存放了兩個常量,一個是話題TOPIC,一個是線程數(shù)THREADS

package com.lya.kafka; 
 
/** 
 * 配置項 
 * @author liuyazhuang 
 * 
 */ 
public class Config { 
  
 /** 
  * 話題 
  */ 
 public static final String TOPIC = "wordcount"; 
 /** 
  * 線程數(shù) 
  */ 
 public static final Integer THREADS = 1; 
} 

二、編程生產(chǎn)者類ProducerDemo

這個類的主要作用就是向Kafka寫入相應(yīng)的消息,并且將消息寫入wordcount話題。

package com.lya.kafka; 
 
import java.util.Properties; 
 
import kafka.javaapi.producer.Producer; 
import kafka.producer.KeyedMessage; 
import kafka.producer.ProducerConfig; 
 
/** 
 * 生產(chǎn)者實例 
 * @author liuyazhuang 
 * 
 */ 
public class ProducerDemo { 
 public static void main(String[] args) throws Exception { 
  Properties props = new Properties(); 
  props.put("zk.connect", "192.168.209.121:2181"); 
  props.put("metadata.broker.list","192.168.209.121:9092"); 
  props.put("serializer.class", "kafka.serializer.StringEncoder"); 
  props.put("zk.connectiontimeout.ms", "15000"); 
  ProducerConfig config = new ProducerConfig(props); 
  Producer<String, String> producer = new Producer<String, String>(config); 
 
  // 發(fā)送業(yè)務(wù)消息 
  // 讀取文件 讀取內(nèi)存數(shù)據(jù)庫 讀socket端口 
  for (int i = 1; i <= 100; i++) { 
   Thread.sleep(500); 
   producer.send(new KeyedMessage<String, String>(Config.TOPIC, 
     "this number ===>>> " + i)); 
  } 
 
 } 
} 

三、編寫消息者類ConsumerDemo

這個類的主要作用就是消費(fèi)Kafka中wordcount話題的消息。

package com.lya.kafka; 
 
import java.util.HashMap; 
import java.util.List; 
import java.util.Map; 
import java.util.Properties; 
 
import kafka.consumer.Consumer; 
import kafka.consumer.ConsumerConfig; 
import kafka.consumer.KafkaStream; 
import kafka.javaapi.consumer.ConsumerConnector; 
import kafka.message.MessageAndMetadata; 
 
/** 
 * 消費(fèi)者實例 
 * @author liuyazhuang 
 * 
 */ 
public class ConsumerDemo { 
  
 
 public static void main(String[] args) { 
   
  Properties props = new Properties(); 
  props.put("zookeeper.connect", "192.168.209.121:2181"); 
  props.put("group.id", "1111"); 
  props.put("auto.offset.reset", "smallest"); 
  props.put("zk.connectiontimeout.ms", "15000"); 
 
  ConsumerConfig config = new ConsumerConfig(props); 
  ConsumerConnector consumer =Consumer.createJavaConsumerConnector(config); 
  Map<String, Integer> topicCountMap = new HashMap<String, Integer>(); 
  topicCountMap.put(Config.TOPIC, Config.THREADS); 
  Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumer.createMessageStreams(topicCountMap); 
  List<KafkaStream<byte[], byte[]>> streams = consumerMap.get(Config.TOPIC); 
   
  for(final KafkaStream<byte[], byte[]> kafkaStream : streams){ 
   new Thread(new Runnable() { 
    @Override 
    public void run() { 
     for(MessageAndMetadata<byte[], byte[]> mm : kafkaStream){ 
      String msg = new String(mm.message()); 
      System.out.println(msg); 
     } 
    } 
    
   }).start(); 
   
  } 
 } 
} 

四、運(yùn)行實例

首先,運(yùn)行消費(fèi)者類ConsumerDemo
運(yùn)行結(jié)果如下:

沒有打印任何信息。
此時,我們運(yùn)行生產(chǎn)者類ProducerDemo
我們再次打開消費(fèi)者的控制臺查看如下:

打印出了生產(chǎn)者生產(chǎn)的消息。
至此,Kafka簡單客戶端編程實例結(jié)束。

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

相關(guān)文章

  • 解決SpringBoot中使用@Async注解失效的問題

    解決SpringBoot中使用@Async注解失效的問題

    這篇文章主要介紹了解決SpringBoot中使用@Async注解失效的問題,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-09-09
  • 基于Java創(chuàng)建XML(無中文亂碼)過程解析

    基于Java創(chuàng)建XML(無中文亂碼)過程解析

    這篇文章主要介紹了基于Java創(chuàng)建XML(無中文亂碼)過程解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2019-10-10
  • jpa實體@ManyToOne @OneToMany無限遞歸方式

    jpa實體@ManyToOne @OneToMany無限遞歸方式

    這篇文章主要介紹了jpa實體@ManyToOne @OneToMany無限遞歸方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-10-10
  • mybatis中的count()按條件查詢方式

    mybatis中的count()按條件查詢方式

    這篇文章主要介紹了mybatis中的count()按條件查詢方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-01-01
  • Java?數(shù)據(jù)結(jié)構(gòu)與算法系列精講之單向鏈表

    Java?數(shù)據(jù)結(jié)構(gòu)與算法系列精講之單向鏈表

    單向鏈表特點(diǎn)是鏈表的鏈接方向是單向的,訪問要通過順序讀取從頭部開始。鏈表是使用指針構(gòu)造的列表,是由一個個結(jié)點(diǎn)組裝起來的,又稱為結(jié)點(diǎn)列表。其中每個結(jié)點(diǎn)都有指針成員變量指向列表中的下一個結(jié)點(diǎn),head指針指向第一個結(jié)點(diǎn)稱為表頭,而終止于最后一個指向nuLL的指針
    2022-02-02
  • java后臺利用Apache poi 生成excel文檔提供前臺下載示例

    java后臺利用Apache poi 生成excel文檔提供前臺下載示例

    本篇文章主要介紹了java后臺利用Apache poi 生成excel文檔提供前臺下載示例,非常具有實用價值,需要的朋友可以參考下
    2017-05-05
  • 深入了解java中的string對象

    深入了解java中的string對象

    這篇文章主要介紹了java中的string對象,String對象是Java中使用最頻繁的對象之一,所以Java開發(fā)者們也在不斷地對String對象的實現(xiàn)進(jìn)行優(yōu)化,以便提升String對象的性能。對此感興趣的朋友跟隨小編一起看看吧
    2019-11-11
  • Spring Boot 日志概念及使用詳解

    Spring Boot 日志概念及使用詳解

    文章介紹了日志的重要性、日志格式、日志使用方法、日志級別、日志配置、日志持久化以及使用Lombok簡化日志輸出,感興趣的朋友一起看看吧
    2025-03-03
  • Java實現(xiàn)畫圖的詳細(xì)步驟(完整代碼)

    Java實現(xiàn)畫圖的詳細(xì)步驟(完整代碼)

    今天給大家?guī)淼氖顷P(guān)于Java的相關(guān)知識,文章圍繞著Java實現(xiàn)畫圖的詳細(xì)步驟展開,文中有非常詳細(xì)的介紹及代碼示例,需要的朋友可以參考下
    2021-06-06
  • 詳解SpringBoot的jar為什么可以直接運(yùn)行

    詳解SpringBoot的jar為什么可以直接運(yùn)行

    SpringBoot提供了一個插件spring-boot-maven-plugin用于把程序打包成一個可執(zhí)行的jar包,本文給大家介紹了為什么SpringBoot的jar可以直接運(yùn)行,文中有相關(guān)的代碼示例供大家參考,感興趣的朋友可以參考下
    2024-02-02

最新評論

福建省| 孟村| 乐清市| 青川县| 松桃| 顺义区| 石渠县| 封丘县| 天柱县| 洮南市| 来凤县| 贵南县| 忻城县| 肇州县| 苏尼特左旗| 新建县| 靖边县| 长汀县| 琼结县| 洪泽县| 罗甸县| 阿拉善盟| 信丰县| 杨浦区| 内乡县| 鹤山市| 稷山县| 体育| 自治县| 黄陵县| 定日县| 辽阳市| 芦山县| 九江市| 石泉县| 屏南县| 施秉县| 曲阳县| 蒙阴县| 阳原县| 奉新县|