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 86482
공지 [SQL컨셉] 서적 "SQL컨셉"의 샘플 데이타 베이스 SAMPLE DATABASE of ORACLE 가을의 곰을... 2013.02.10 78912
공지 [G_SQL] Sample Database 가을의 곰을... 2012.05.20 95666
14 [SQLite]10만 TPS와 10억 행 처리: SQLite의 놀라운 효율성 졸리운_곰 2025.12.05 1184
13 [sqlite] SQlite source code analysis-architecture file 졸리운_곰 2021.04.12 1809
12 [SQLite] SQLite 사용자 함수 추가 졸리운_곰 2021.04.12 1665
11 {SQLite] SQLite 페이지 핸들링(3) - 레코드 포맷 졸리운_곰 2021.04.12 1670
10 [SQLite] SQLite 페이지 핸들링(2) - SQLite의 페이지 포맷 file 졸리운_곰 2021.04.12 1496
9 [SQLite] SQLite 페이지 핸들링(1) - SQLite의 구조 file 졸리운_곰 2021.04.12 1487
8 [C/C++ 자료구조] SQLite 의 모든 것 (4부) - Java 에서 사용하기 Database/SQLite file 졸리운_곰 2021.04.12 1663
7 [C/C++] SQLite 의 모든 것 (3부) - C++ 에서 사용하기 Database/SQLite 졸리운_곰 2021.04.12 1722
6 [C/C++ 자료구조] SQLite 의 모든 것 (2부) - Download & Build Database/SQLite file 졸리운_곰 2021.04.12 1668
5 [C/C++ 자료구조] SQLite 의 모든 것 (1부) - 소개 및 FAQ Database/SQLite file 졸리운_곰 2021.04.12 1332
4 SQLite-sync.com version 3 file 졸리운_곰 2017.07.03 2019
3 DB Browser for SQLite file 졸리운_곰 2017.04.26 1653
2 Getting started with SQLite in C# file 졸리운_곰 2017.04.24 1828
1 SQLite C/C++ Tutorial 졸리운_곰 2017.04.24 1498
대표 김성준 주소 : 경기 용인 분당수지 U타워 등록번호 : 142-07-27414
통신판매업 신고 : 제2012-용인수지-0185호 출판업 신고 : 수지구청 제 123호 개인정보보호최고책임자 : 김성준 sjkim70@stechstar.com
대표전화 : 010-4589-2193 [fax] 02-6280-1294 COPYRIGHT(C) stechstar.com ALL RIGHTS RESERVED