- Download the Kafka binaries from Kafka download page
- Unzip the kafka tar file by executing
tar -xzf kafka_2.9.2-0.8.1.1.tgz. Then go to kafka directory by executingcd kafka_2.9.2-0.8.1.1 - Next start the Zookeeper server by executing following command
bin/zookeeper-server-start.sh config/zookeeper.properties - Start the Kafka server by executing following command
bin/kafka-server-start.sh config/server.properties - Now your Zookeeper and Kafka server are ready and you can download the source code for sample project from here
- This is how a Java Client that publishes messages to Kafka looks like, execute it couple of times to publish couple of messages
First thing that you have to do while developing a producer is connect to the Kafka server, for that you will set value of
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(); } } metadata.broker.listproperty 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 haveProducerobject you can use it for publishing messages by creating object ofkafka.producer.KeyedMessage, you will have to pass name of the topic and message as argument - 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
The
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")); } } } HelloKafkaConsumerclass extends Thread class. In the constructor of this class first i am creating Properties class with value ofzookeeper.connectproperty equal to the port on which zookeeper server is listening on. In the constructor i am creating object ofkafka.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.
- 전체
- Sample DB
- database modeling
- [표준 SQL] Standard SQL
- G-SQL
- 10-Min
- ORACLE
- MS SQLserver
- MySQL
- SQLite
- postgreSQL
- 데이터아키텍처전문가 - 국가공인자격
- 데이터 분석 전문가 [ADP]
- [국가공인] SQL 개발자/전문가
- NoSQL
- hadoop
- hadoop eco system
- big data (빅데이터)
- stat(통계) R 언어
- XML DB & XQuery
- spark
- DataBase Tool
- 데이터분석 & 데이터사이언스
- Engineer Quality Management
- [기계학습] machine learning
- 데이터 수집 및 전처리
- 국가기술자격 빅데이터분석기사
- 암호화폐 (비트코인, cryptocurrency, bitcoin)
big data (빅데이터) Java Client for publishing and consuming messages from Apache Kafka
2016.06.05 20:51
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.
[출처] http://wpcertification.blogspot.kr/2014/08/java-client-for-publishing-and.html
본 웹사이트는 광고를 포함하고 있습니다.
광고 클릭에서 발생하는 수익금은 모두 웹사이트 서버의 유지 및 관리, 그리고 기술 콘텐츠 향상을 위해 쓰여집니다.
광고 클릭에서 발생하는 수익금은 모두 웹사이트 서버의 유지 및 관리, 그리고 기술 콘텐츠 향상을 위해 쓰여집니다.
댓글 0
| 번호 | 제목 | 글쓴이 | 날짜 | 조회 수 |
|---|---|---|---|---|
| 공지 | 오라클 기본 샘플 데이터베이스 | 졸리운_곰 | 2014.01.02 | 86317 |
| 공지 | [SQL컨셉] 서적 "SQL컨셉"의 샘플 데이타 베이스 SAMPLE DATABASE of ORACLE | 가을의 곰을... | 2013.02.10 | 78771 |
| 공지 | [G_SQL] Sample Database | 가을의 곰을... | 2012.05.20 | 95528 |

