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 86873
공지 [SQL컨셉] 서적 "SQL컨셉"의 샘플 데이타 베이스 SAMPLE DATABASE of ORACLE 가을의 곰을... 2013.02.10 79171
공지 [G_SQL] Sample Database 가을의 곰을... 2012.05.20 95908
18 TensorFlow Lite 101 - MoblieNet 맛보기 file 졸리운_곰 2018.05.30 1093
17 Apache MXNet에서 사전 트레이닝된 모델을 사용해 보세요. 졸리운_곰 2018.05.30 993
16 MXNet을 활용한 이미지 분류 앱 개발하기 file 졸리운_곰 2018.05.30 1428
15 세상에 있는 (거의) 모든 머신러닝 문제 공략법 file 졸리운_곰 2018.05.30 1121
14 sklearn 내부의 pickle lib 를 통해 모델을 저장하고 다시 로드하여 재사용할 수 있다. 졸리운_곰 2018.05.30 1640
13 텐서플로우 기반 딥러닝 훈련 모델 파일 저장, 로딩 및 재활용 file 졸리운_곰 2018.05.30 1773
12 TensorFlow 모델을 저장하고 불러오기 (save and restore) 졸리운_곰 2018.05.30 1640
11 텐서플로우(TensorFlow)를 이용해서 글자 생성(Text Generation) 해보기 – Recurrent Neural Networks(RNNs) 예제 – Char-RNN file 졸리운_곰 2018.05.13 1125
10 외장형 그래픽카드로 우분투에서 텐서플로우 사용 How to setup an eGPU on Ubuntu for TensorFlow file 졸리운_곰 2018.05.09 1428
9 Keras and NLTK 케라스를 이용한 NLTK 자연어처리 졸리운_곰 2018.05.08 1607
8 Windows7에서 "처음 심층 학습 프로그램」을 사경 보면 (1-2) 제 1 장 후반 file 졸리운_곰 2018.05.08 1143
7 CSLAIER CSLAIER에 의한 LSTM file 졸리운_곰 2018.05.08 1666
6 Tensorboard 사용하기 1 file 졸리운_곰 2018.05.08 1417
5 텐서보드 사용법 file 졸리운_곰 2018.05.08 1415
4 Ubuntu 18.04 Settings for TensorFlow 설치 file 졸리운_곰 2018.05.08 1184
3 딥러닝용 서버 설치기 file 졸리운_곰 2018.05.07 1396
2 Installing Tensorflow GPU on Ubuntu 18.04 LTS file 졸리운_곰 2018.05.06 1217
1 TensorFlow Lite 101 - MoblieNet 맛보기 file 졸리운_곰 2018.04.07 1275
대표 김성준 주소 : 경기 용인 분당수지 U타워 등록번호 : 142-07-27414
통신판매업 신고 : 제2012-용인수지-0185호 출판업 신고 : 수지구청 제 123호 개인정보보호최고책임자 : 김성준 sjkim70@stechstar.com
대표전화 : 010-4589-2193 [fax] 02-6280-1294 COPYRIGHT(C) stechstar.com ALL RIGHTS RESERVED