Java Client for publishing and consuming messages from Apache Kafka

 
I wanted to learn how to use Apache Kafka for publishing and consuming messages from Apache Kafka using Java client, so i followed these steps.
  1. Download the Kafka binaries from Kafka download page
  2. Unzip the kafka tar file by executing tar -xzf kafka_2.9.2-0.8.1.1.tgz. Then go to kafka directory by executing cd kafka_2.9.2-0.8.1.1
  3. Next start the Zookeeper server by executing following command
    
    bin/zookeeper-server-start.sh config/zookeeper.properties
    
  4. Start the Kafka server by executing following command
    
    bin/kafka-server-start.sh config/server.properties
    
  5. Now your Zookeeper and Kafka server are ready and you can download the source code for sample project from here
  6. This is how a Java Client that publishes messages to Kafka looks like, execute it couple of times to publish couple of messages
      package com.spnotes.kafka;
       
      import java.text.SimpleDateFormat;
      import java.util.Date;
      import java.util.Properties;
       
      import kafka.producer.KeyedMessage;
      import kafka.producer.ProducerConfig;
       
      /**
      * Created by user on 8/4/14.
      */
      public class HelloKafkaProducer {
      final static String TOPIC = "javatest";
       
       
      public static void main(String[] argv){
      Properties properties = new Properties();
      properties.put("metadata.broker.list","localhost:9092");
      properties.put("serializer.class","kafka.serializer.StringEncoder");
      ProducerConfig producerConfig = new ProducerConfig(properties);
      kafka.javaapi.producer.Producer<String,String> producer = new kafka.javaapi.producer.Producer<String, String>(producerConfig);
      SimpleDateFormat sdf = new SimpleDateFormat();
      KeyedMessage<String, String> message =new KeyedMessage<String, String>(TOPIC,"Test message from java program " + sdf.format(new Date()));
      producer.send(message);
      producer.close();
      }
      }
    First thing that you have to do while developing a producer is connect to the Kafka server, for that you will set value of metadata.broker.list property to point to the port on which kafka server is listening (You can find value of port and host name from server.properties that you used in step 4. Once you have Producer object you can use it for publishing messages by creating object of kafka.producer.KeyedMessage, you will have to pass name of the topic and message as argument
  7. This is how the Java client for consumer of messages from Kafka looks like, run it and it will start a thread that will keep listening to messages on topic and every time there is a message it will print it to console
      package com.spnotes.kafka;
       
      import java.io.UnsupportedEncodingException;
      import java.nio.ByteBuffer;
      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.ConsumerIterator;
      import kafka.consumer.KafkaStream;
      import kafka.javaapi.consumer.ConsumerConnector;
      import kafka.javaapi.message.ByteBufferMessageSet;
      import kafka.message.MessageAndOffset;
       
       
      /**
      * Created by user on 8/4/14.
      */
      public class HelloKafkaConsumer extends Thread {
      final static String clientId = "SimpleConsumerDemoClient";
      final static String TOPIC = "pythontest";
      ConsumerConnector consumerConnector;
       
       
      public static void main(String[] argv) throws UnsupportedEncodingException {
      HelloKafkaConsumer helloKafkaConsumer = new HelloKafkaConsumer();
      helloKafkaConsumer.start();
      }
       
      public HelloKafkaConsumer(){
      Properties properties = new Properties();
      properties.put("zookeeper.connect","localhost:2181");
      properties.put("group.id","test-group");
      ConsumerConfig consumerConfig = new ConsumerConfig(properties);
      consumerConnector = Consumer.createJavaConsumerConnector(consumerConfig);
      }
       
      @Override
      public void run() {
      Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
      topicCountMap.put(TOPIC, new Integer(1));
      Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumerConnector.createMessageStreams(topicCountMap);
      KafkaStream<byte[], byte[]> stream = consumerMap.get(TOPIC).get(0);
      ConsumerIterator<byte[], byte[]> it = stream.iterator();
      while(it.hasNext())
      System.out.println(new String(it.next().message()));
       
      }
       
      private static void printMessages(ByteBufferMessageSet messageSet) throws UnsupportedEncodingException {
      for(MessageAndOffset messageAndOffset: messageSet) {
      ByteBuffer payload = messageAndOffset.message().payload();
      byte[] bytes = new byte[payload.limit()];
      payload.get(bytes);
      System.out.println(new String(bytes, "UTF-8"));
      }
      }
      }
    The HelloKafkaConsumer class extends Thread class. In the constructor of this class first i am creating Properties class with value of zookeeper.connect property equal to the port on which zookeeper server is listening on. In the constructor i am creating object of kafka.javaapi.consumer.ConsumerConnector
    Once the ConsumerConnector is ready in the run() method i am passing it name of the topic on which i want to listen (You can pass multiple topic names here). Everytime there is a new message i am reading it and printing it to console.

 

 

경축! 아무것도 안하여 에스천사게임즈가 새로운 모습으로 재오픈 하였습니다.
어린이용이며, 설치가 필요없는 브라우저 게임입니다.
https://s1004games.com

 

[출처] http://wpcertification.blogspot.kr/2014/08/java-client-for-publishing-and.html

본 웹사이트는 광고를 포함하고 있습니다.
광고 클릭에서 발생하는 수익금은 모두 웹사이트 서버의 유지 및 관리, 그리고 기술 콘텐츠 향상을 위해 쓰여집니다.
번호 제목 글쓴이 날짜 조회 수
공지 오라클 기본 샘플 데이터베이스 졸리운_곰 2014.01.02 87093
공지 [SQL컨셉] 서적 "SQL컨셉"의 샘플 데이타 베이스 SAMPLE DATABASE of ORACLE 가을의 곰을... 2013.02.10 79310
공지 [G_SQL] Sample Database 가을의 곰을... 2012.05.20 96042
38 [책보다낫다] kafka 활용문서 Kafka 0.10.0 Documentation file 졸리운_곰 2016.06.05 3036
37 [책보다 낫다] flume 사용자 가이드 졸리운_곰 2016.06.05 1311
» Java Client for publishing and consuming messages from Apache Kafka 졸리운_곰 2016.06.05 1280
35 kafka create message : 0.8.0 Producer Example 졸리운_곰 2016.06.05 971
34 Create a topic - Apache Kafka 졸리운_곰 2016.06.05 1310
33 [flume] pollable source 플럼 주기적 실행 custom sources in flume file 졸리운_곰 2016.06.05 1182
32 Flume 메트릭 커스텀 리포터 구현 file 졸리운_곰 2016.06.05 1323
31 flume-ng를 윈도에서 구동하려면 졸리운_곰 2016.06.05 1323
30 Using Kafka with Flume 졸리운_곰 2016.06.05 1658
29 Flafka: Apache Flume Meets Apache Kafka for Event Processing file 졸리운_곰 2016.06.05 2910
28 6. 주키퍼(zookeeper) 활용 - 분산서버 구현 1편 file 졸리운_곰 2016.05.30 1423
27 주키퍼 (ZooKeeper)란? file 졸리운_곰 2016.05.30 1351
26 ZooKeeper란 무엇인가? file 졸리운_곰 2016.05.30 1917
25 Parse Redis dump.rdb files, Analyze Memory, and Export Data to JSON 졸리운_곰 2016.05.29 1989
24 Flume or Kafka? Try both! 졸리운_곰 2016.05.29 1011
23 Flafka: Apache Flume Meets Apache Kafka for Event Processing file 졸리운_곰 2016.05.29 1567
22 네이버 라인은 왜 카카오톡보다 병목현상이 적을까? file 졸리운_곰 2016.05.29 2542
21 최악의 빅데이터 프랙티스 10가지 졸리운_곰 2016.05.29 1469
20 Redis(레디스)를 어떻게 활용할 수 있을까? 졸리운_곰 2016.05.29 2069
19 [NoSQL & Cache] Redis vs Memcached ( 왜 Redis 를 사용해야 하는가? ) 졸리운_곰 2016.05.29 1178
대표 김성준 주소 : 경기 용인 분당수지 U타워 등록번호 : 142-07-27414
통신판매업 신고 : 제2012-용인수지-0185호 출판업 신고 : 수지구청 제 123호 개인정보보호최고책임자 : 김성준 sjkim70@stechstar.com
대표전화 : 010-4589-2193 [fax] 02-6280-1294 COPYRIGHT(C) stechstar.com ALL RIGHTS RESERVED