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 86334
공지 [SQL컨셉] 서적 "SQL컨셉"의 샘플 데이타 베이스 SAMPLE DATABASE of ORACLE 가을의 곰을... 2013.02.10 78792
공지 [G_SQL] Sample Database 가을의 곰을... 2012.05.20 95542
42 블록체인 기반 플랫폼 비즈니스를 이해하자 졸리운_곰 2017.09.10 1554
41 다운타임 없는 서비스 구현 패턴 file 졸리운_곰 2017.09.10 1229
40 테이블의 수직분할과 수평분할에 대한 이해 file 졸리운_곰 2017.09.10 3910
39 정규화와 응집도에 대한 고찰 file 졸리운_곰 2017.05.28 1782
38 머신러닝 새 도전…“클라우드를 벗어나라" file 졸리운_곰 2017.05.28 1711
37 보안성 높이는 공공 거래장부 블록체인 file 졸리운_곰 2017.05.28 1757
36 디지털, 속도의 전쟁 VS 데이터, 품질의 전쟁 file 졸리운_곰 2017.05.05 1516
35 04. 데이터 모델링의 3단계 진행 file 졸리운_곰 2016.03.15 1659
34 데이터베이스 설계의 기본 원리.pdf file 졸리운_곰 2016.03.15 2354
33 실체유형(Entity Type) 정의 사항 및 도출 file 졸리운_곰 2015.05.21 2199
32 마농의 SQL 백문백답: 단순하고 쉽게 작성하는 SQL 노하우 [1회] file 졸리운_곰 2015.05.21 1830
31 sql개발자-sql전문가자격시험 시험 예제.pdf file 졸리운_곰 2015.02.15 2303
30 데이터베이스 선정에는 비밀이 있다 - 4부 졸리운_곰 2015.01.15 2099
29 데이터베이스 선정에는 비밀이 있다 - 3부 졸리운_곰 2015.01.15 2044
28 데이터베이스 선정에는 비밀이 있다 - 2부 졸리운_곰 2015.01.15 1828
27 데이터베이스 선정에는 비밀이 있다 - 1부 졸리운_곰 2015.01.15 2392
26 지금 우리에게 필요한 것은 데이터베이스 성능 최적화이다 (2부) 졸리운_곰 2015.01.15 1488
25 우리에게 필요한 것은 데이터베이스 성능 (1부) 졸리운_곰 2015.01.15 1782
24 21회 결과 secret 졸리운_곰 2014.11.10 0
23 [데이터아키텍쳐준전문가] 시험 fail 자료 secret 졸리운_곰 2014.08.31 0
대표 김성준 주소 : 경기 용인 분당수지 U타워 등록번호 : 142-07-27414
통신판매업 신고 : 제2012-용인수지-0185호 출판업 신고 : 수지구청 제 123호 개인정보보호최고책임자 : 김성준 sjkim70@stechstar.com
대표전화 : 010-4589-2193 [fax] 02-6280-1294 COPYRIGHT(C) stechstar.com ALL RIGHTS RESERVED