初識Kafka----------Centos上單機部署、服務啟動、JAVA客戶端調用


  作為Apach下一個優秀的開源消息隊列框架,Kafka已經成為很多互聯網廠商日志采集處理的第一選擇。后面在實際應用場景中可能會應用到,因此就先了解了一下。經過兩個晚上的努力,總算是能夠基本使用。

操作系統:虛擬機Centos 6.5 

1、下載Kafka安裝文件,首先進入官網,找到最新的穩定版本

      wget http://mirrors.hust.edu.cn/apache/kafka/0.10.2.0/kafka_2.12-0.10.2.0.tgz

2、解壓並拷貝到 需要的目錄下,我的設定為 /usr/下

       先 cp   然后解壓 tar -xzvf 

3、由於我本機已經安裝了zookeeper,因此直接修改server.properties 文件

4、啟動服務  bin/kafka-server-start.sh config/server.properties ,問題來了 :

[root@localhost kafka_2.12-0.10.2.0]# Exception in thread "main" java.lang.UnsupportedClassVersionError: kafka/Kafka : Unsupported major.minor version 52.0
at java.lang.ClassLoader.defineClass1(Native Method)
at java.lang.ClassLoader.defineClassCond(ClassLoader.java:631)
at java.lang.ClassLoader.defineClass(ClassLoader.java:615)
 
啟動報錯,看提示,由於使用最新的kafka版本,需要1.8的JDK。而我本機查看目前是1.6。
 
5、更換JDK版本
使用wget下載1.8JDK,vim /etc/profile 更改JAVA_HOME的路徑,source /etc/profile 后,執行 java -version 仍然顯示為 1.6.
重啟依然無效,后網上找到解決辦法:
which  Java
/usr/bin/java
which javac
/usr/bin/javac
 
1.先將usr/bin目錄下的先刪除
rm -rf  java
rm -rf  javac
2.先將jdk1.8
ln -s   $JAVA_HOME/bin/java  /usr/bin/java
ln -s   $JAVA_HOME/bin/javac  /usr/bin/javac
 
6、啟動服務  bin/kafka-server-start.sh config/server.properties  ,未報錯
啟動一個生產者 ,topic為test:
bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test
啟動兩個消費者
bin/kafka-console-consumer.sh --zookeeper 192.168.118.131:3181 --topic test --from-beginning 
bin/kafka-console-consumer.sh --zookeeper 192.168.118.131:3181 --topic test --from-beginning 

兩個消費端均能收到服務端的信息。

7、嘗試用java代碼來接收信息,首先建立一個MAVEN工程,POM.xml文件加入依賴:

<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>0.8.2.1</version>
</dependency>

<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka_2.11</artifactId>
<version>0.8.2.1</version>
</dependency>

然后運行安裝代碼構建完成。具體調用代碼如下

import java.util.List;
import java.util.Properties;
import java.util.concurrent.TimeUnit;

import kafka.consumer.Consumer;
import kafka.consumer.ConsumerConfig;
import kafka.consumer.ConsumerIterator;
import kafka.consumer.KafkaStream;
import kafka.consumer.Whitelist;
import kafka.javaapi.consumer.ConsumerConnector;
import kafka.message.MessageAndMetadata;

import org.apache.kafka.common.utils.CollectionUtils;

/**
* File Name:KafkaConsumer.java
* Package_Name:kafkaTest
* date:2017-3-25上午10:00:56
* Author : cao.zhi10
*
*/
public class KafkaConsumer {
public static void main(String[] args) throws Exception {
Properties properties = new Properties();
properties.put("zookeeper.connect", "192.168.118.131:3181");
properties.put("auto.commit.enable", "true");
properties.put("auto.commit.interval.ms", "60000");
properties.put("group.id", "test");

ConsumerConfig consumerConfig = new ConsumerConfig(properties);

ConsumerConnector javaConsumerConnector = Consumer.createJavaConsumerConnector(consumerConfig);

//topic的過濾器
Whitelist whitelist = new Whitelist("test");
List<KafkaStream<byte[], byte[]>> partitions = javaConsumerConnector.createMessageStreamsByFilter(whitelist);

if (partitions==null) {
System.out.println("empty!");
TimeUnit.SECONDS.sleep(1);
}

//消費消息
for (KafkaStream<byte[], byte[]> partition : partitions) {

ConsumerIterator<byte[], byte[]> iterator = partition.iterator();
while (iterator.hasNext()) {
MessageAndMetadata<byte[], byte[]> next = iterator.next();
System.out.println("partiton:" + next.partition());
System.out.println("offset:" + next.offset());
System.out.println("接收到message:" + new String(next.message(), "utf-8"));
}
}
}
}

執行main方法,一直報DISCONNECT異常,仔細分析啟動日志發現一直在去嘗試調用 localhost:9092 。搜索了下該問題,網上給的解決方案是修改 service.properties 里面的hostname 信息,但是目前最新版本已經取消了該節點,目前必須使用 listeners.

 

 8、重新啟動服務,生產者,消費者(注意必須用IP地址啟動,否則會報錯

bin/kafka-console-producer.sh --broker-list 192.168.118.131:9092 --topic test

 

 
 
 
 
9、java代碼啟動后,在生產者crt界面輸入可以正常接收信息,但是出現了亂碼:

 

 字符轉換方式 GBK,UTF-8 都不行。后來想到原來實時查看服務器端日志也出現亂碼的解決方式,於是修改了下 crt會話里面的編碼方式,如下

問題得到解決

 

 

 

 

 

 

 

 

 
 

       

    

 


免責聲明!

本站轉載的文章為個人學習借鑒使用,本站對版權不負任何法律責任。如果侵犯了您的隱私權益,請聯系本站郵箱yoyou2525@163.com刪除。



 
粵ICP備18138465號   © 2018-2025 CODEPRJ.COM