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 86349
공지 [SQL컨셉] 서적 "SQL컨셉"의 샘플 데이타 베이스 SAMPLE DATABASE of ORACLE 가을의 곰을... 2013.02.10 78808
공지 [G_SQL] Sample Database 가을의 곰을... 2012.05.20 95551
22 UNION과 UNION ALL 의 차이 및 주의 사항 졸리운_곰 2017.08.27 1282
21 컬럼내 특정 문자를 다른문자로 변경하고자 할때 졸리운_곰 2017.06.10 1500
20 [MySQL] 레코드 데이터 치환하기 (REPLACE) 졸리운_곰 2017.06.10 1419
19 MySQL / 테이블에서 특정 문자열 바꾸기 졸리운_곰 2017.06.10 1305
18 MySQL Redis Plugin file 졸리운_곰 2017.05.30 1722
17 [DB] MySQL Check, Repair, Optimize(개별/전체 테이블 포함) 졸리운_곰 2017.05.21 1249
16 Transfer from sqlite to MySQL/ sqlite에서 mysql로 변환 졸리운_곰 2017.03.18 1115
15 [MySQL] 힌트설정 / 쿼리캐시 졸리운_곰 2017.03.15 1506
14 [MySQL힌트 정리] 졸리운_곰 2017.03.15 1715
13 [mysql]Hint 사용방법 졸리운_곰 2017.03.15 1502
12 MySQL Ver. 5.1 힌트를 이용한 실행계획 제어 file 졸리운_곰 2017.03.15 1183
11 MySQL 덤프 / 임포트 dump / import 졸리운_곰 2017.01.05 1520
10 [질문] 두개의 컬럼에 대해 group by 적용 할 수 있을까요? 졸리운_곰 2016.12.17 841
9 MySQL 중복 키 관리 방법 (INSERT 시 중복 키 관리 방법 (INSERT IGNORE, REPLACE INTO, ON DUPLICATE UPDATE) 졸리운_곰 2016.12.17 1106
8 [mysql] 쿼리값이 NULL 일때 0으로 바꾸기 졸리운_곰 2016.12.14 1194
7 MySql] Insert Select 문 졸리운_곰 2016.12.06 1410
6 [MySQL] substring_index , substring ( split, explode ) 졸리운_곰 2016.12.06 1366
5 MySQL 테이블 이름변경, 테이블 복사 졸리운_곰 2016.11.23 1534
4 [MySQL] MySQL 테이블 수정 졸리운_곰 2016.11.16 1541
3 generate days from date range mysql 날짜검색시 between 안에 포함되는 날짜전체 출력 졸리운_곰 2016.11.02 1258
대표 김성준 주소 : 경기 용인 분당수지 U타워 등록번호 : 142-07-27414
통신판매업 신고 : 제2012-용인수지-0185호 출판업 신고 : 수지구청 제 123호 개인정보보호최고책임자 : 김성준 sjkim70@stechstar.com
대표전화 : 010-4589-2193 [fax] 02-6280-1294 COPYRIGHT(C) stechstar.com ALL RIGHTS RESERVED