您的位置:首页 > 其它

我的卡夫卡(kafka)只是自己的一个笔记

2017-03-29 15:50 274 查看
kafka 的准备工作

第一步 下载卡夫卡 
http://kafka.apache.org/downloads
kafka_2.11-0.9.0.0(我的)

第二步下载jar包

 kafka_2.10-0.8.2.0.jar
        kafka-clients-0.8.2.0.jar
        metrics-core-2.2.0.jar
        scala-library-2.10.4.jar
        zkclient-0.3.jar
        zookeeper-3.4.6.jar

第三步启动 zookeeper-server-start  和  kafka-server-start

     启动方式 比如解压后的文件放在D:\kafuka\下

打开cmd 找到该目录 D:\kafuka\kafka_2.11-0.9.0.0\bin\windows 回车

D:\kafuka\kafka_2.11-0.9.0.0\bin\windows>zookeeper-server-start D:\kafuka\kafka_2.11-0.9.0.0\config\zookeeper.properties

将红色部分粘贴,(D:\kafuka\kafka_2.11-0.9.0.0\config\zookeeper.properties)是因为找到所对应的文件,可以不必绝对路径

打开cmd 找到该目录 D:\kafuka\kafka_2.11-0.9.0.0\bin\windows 回车

D:\kafuka\kafka_2.11-0.9.0.0\bin\windows>kafka-server-start D:\kafuka\kafka_2.11-0.9.0.0\config\server.properties

将红色部分粘贴,(D:\kafuka\kafka_2.11-0.9.0.0\config\server.properties)是因为找到所对应的文件,可以不必绝对路径

(切记,启动项中的配置需要论情况而定 eg:set KAFKA_HEAP_OPTS=-Xmx512M -Xms512M)太大的话,内存小的电脑是启动不了

第四步 编码

配置文件

KafkaProperties.java

public interface KafkaProperties

{

    final static String zkConnect = "localhost:2181";

    final static String groupId = "group1";

    final static String topic = "topic1";

    final static String kafkaServerURL = "localhost";

    final static int kafkaServerPort = 9092;

    final static int kafkaProducerBufferSize = 64 * 1024;

    final static int connectionTimeOut = 20000;

    final static int reconnectInterval = 10000;

    final static String topic2 = "topic2";

    final static String topic3 = "topic3";

    final static String clientId = "SimpleConsumerDemoClient";

}

KafkaProducer
.java 发送方(生产方)

import java.util.Properties;

import kafka.producer.KeyedMessage;

import kafka.producer.ProducerConfig;

public class KafkaProducer extends Thread

{

    private final kafka.javaapi.producer.Producer<Integer, String> producer;

    private final String topic;

    private final Properties props = new Properties();

    public KafkaProducer(String topic)

    {

        props.put("serializer.class", "kafka.serializer.StringEncoder");

        props.put("metadata.broker.list", "localhost:9092");

        producer = new kafka.javaapi.producer.Producer<Integer, String>(new ProducerConfig(props));

        this.topic = topic;

    }

    @Override

    public void run() {

        int messageNo = 1;

        while (true)

        {

            String messageStr = new String("Message_" + messageNo);

            System.out.println("Send:" + messageStr);

            producer.send(new KeyedMessage<Integer, String>(topic, messageStr));

            messageNo++;

            try {

                sleep(3000);

            } catch (InterruptedException e) {

                // TODO Auto-generated catch block

                e.printStackTrace();

            }

        }

    }

}

KafkaConsumer.java 接收方(消费者)

import java.util.HashMap;

import java.util.List;

import java.util.Map;

import java.util.Properties;

import kafka.consumer.ConsumerConfig;

import kafka.consumer.ConsumerIterator;

import kafka.consumer.KafkaStream;

import kafka.javaapi.consumer.ConsumerConnector;

public class KafkaConsumer extends Thread

{

    private final ConsumerConnector consumer;

    private final String topic;

    public KafkaConsumer(String topic)

    {

        consumer = kafka.consumer.Consumer.createJavaConsumerConnector(

                createConsumerConfig());

        this.topic = topic;

    }

    private static ConsumerConfig createConsumerConfig()

    {

        Properties props = new Properties();

        props.put("zookeeper.connect", KafkaProperties.zkConnect);

        props.put("group.id", KafkaProperties.groupId);

        props.put("zookeeper.session.timeout.ms", "40000");

        props.put("zookeeper.sync.time.ms", "200");

        props.put("auto.commit.interval.ms", "1000");

        return new ConsumerConfig(props);

    }

    @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 = consumer.createMessageStreams(topicCountMap);

        KafkaStream<byte[], byte[]> stream = consumerMap.get(topic).get(0);

        ConsumerIterator<byte[], byte[]> it = stream.iterator();

        while (it.hasNext()) {

            System.out.println("receive:" + new String(it.next().message()));

            try {

                sleep(3000);

            } catch (InterruptedException e) {

                e.printStackTrace();

            }

        }

    }

}

KafkaConsumerProducerDemo.java (启动)

public class KafkaConsumerProducerDemo
{
    public static void main(String[] args)
    {

        KafkaProducer producerThread = new KafkaProducer(KafkaProperties.topic);

        producerThread.start();

        KafkaConsumer consumerThread = new KafkaConsumer(KafkaProperties.topic);
        consumerThread.start();
    }
}
内容来自用户分享和网络整理,不保证内容的准确性,如有侵权内容,可联系管理员处理 点击这里给我发消息
相关文章推荐