初识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 ,问题来了 :

[[email protected] 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会话里面的编码方式,如下

问题得到解决

时间: 2024-11-06 21:17:19

初识Kafka----------Centos上单机部署、服务启动、JAVA客户端调用的相关文章

CentOS上安装GitBlit服务

简单介绍 在上一篇文章中,已经简单的介绍了如何在CentOS的服务器上搭建git服务器.但是这种方式实现的服务器功能比较弱,操作起来也比较繁琐.在网上搜索了一圈,感觉Gitblit比较符合我的需求.接下来我就简单地介绍下,如何在CentOS上搭建GitBlit服务吧. GitBlit是一款纯Java库实现用来管理.查看和处理Git资料库,相当于Git的Java管理工具.该管理软件支持Windows和Linux平台.可以有效的对项目.用户权限进行控制和管理.比较适合小型团队进行管理控制. 看上面的

(转)解决:本地计算机 上的 OracleOraDb10g_home1TNSListener服务启动后停止

原文地址:http://justsee.iteye.com/blog/1320059 手动启动一个问题:本地计算机 上的 OracleOraDb10g_home1TNSListener服务启动后停止.某些服务在未由其他服务或程序使用时将自动停止. 在网上找解决方案的时候,发现很多人都遇到了这个问题,但都没有解决.下面自己记录一下,留个备份,方便下次查阅方便 问题1:首先查阅你的[NETWORK\ADMIN]目录下的[tnsnames.ora]和[listener.ora]这两个文件,我的路径是:

Centos 6&7下服务启动方法及添加到开机启动

在linux系统中,安装完一个软件或应用后,有时候需要手动启动该应用,也需要收到将该应用添加到开机启动项中,让其可以能够在linux一开机后就加载该应用 启动应用的方法 CentOS 6 : service SERVICE start|stop|restart|reload|status CentOS 7 : systemctl start|stop|restart|reload|status SERVICE 添加到开机启动项的方法 CentOS 6 : chkconfig SERVICE on

本地计算机 上的 OracleOraDb11g_home1TNSListener 服务启动后停止

今天玩oracle的时候突然遇到一个问题:本地计算机 上的 OracleOraDb11g_home1TNSListener 服务启动后停止.某些服务在未由其他服务或程序使用时将自动停止. 在网上找解决方案的时候,发现很多人都遇到了这个问题,第一个方案没有解决我的问题,下面自己记录一下,留个备份,方便下次查阅方便 第一步:首先查阅你的[NETWORK\ADMIN]目录下的[tnsnames.ora]和[listener.ora]这两个文件,我的路径是:D:\app\Oracle11g\dbhome

MySQL 安装和启动服务,“本地计算机 上的 MySQL 服务启动后停止。某些服务在未由其他服务或程序使用时将自动停止。”

MySQL 安装和启动服务,以及遇到的问题 MySQL版本: mysql-5.7.13-winx64.zip (免安装,解压放到程序文件夹即可,比如 C:\Program Files\mysql-5.7.13-winx64) 下载地址:http://dev.mysql.com/get/Downloads/MySQL-5.7/mysql-5.7.13-winx64.zip 遇到的问题: 1. MySQL service 已经安装成功,创建了空的data文件夹,也填了初始化ini文件,但是无法启动

阿里云CentOS 7.2 MySQL服务启动失败的解决思路

阿里云 CentOS 7.2 MySQL服务启动失败的解决思路 前言 : 昨天刚刚搭建好的MySQL让老大看了一下,经过测试已经完成任务.但是今天早晨来的时候发现服务器被关了,此时我的心情崩溃的,但是我非常冷静的解决了MySQL问题.如下: 启动MySQL服务器失败,如下所示: [[email protected] ~]# /etc/init.d/mysqld start Starting mysqld (via systemctl):  Job for mysqld.service faile

淘宝分布式 key/value 存储引擎Tair安装部署过程及Java客户端测试一例

目录 1. 简介 2. 安装步骤及问题小记 3. 部署配置 4. Java客户端测试 5. 参考资料 声明 1. 下面的安装部署基于Linux系统环境:centos 6(64位),其它Linux版本可能有所差异. 2. 网上有人说tair安装失败可能是因为gcc版本问题,高版本的gcc可能不支持某些特性导致安装失败,经过实验证明,该说法是错误的,tair安装失败有各种可能的原因但绝对与gcc版本无关,比如我的gcc开始版本为4.4.7,后来tair安装失败,我重新编译低版本的gcc(gcc4.1

转载——Java与WCF交互(一):Java客户端调用WCF服务

最近开始了解WCF,写了个最简单的Helloworld,想通过java客户端实现通信.没想到以我的基础,居然花了整整两天(当然是工作以外的时间,呵呵),整个过程大费周折,特写下此文,以供有需要的朋友参考: 第一步:生成WCF服务 新建WCF解决方案,分别添加三个项目,HelloTimeService(类库),HelloTimehost(控制台程序),HelloTimeClient(控制台程序),项目结构如图:各个项目的主要代码:service: Host: Client: 编译通过后,测试Hos

搭建基于asp.net的wcf服务,ios客户端调用的实现记录

一.写wcf 问题: 1.特定的格式 2.数据绑定 3.加密解密 二.发布到iis 问题: 1.访问权限问题,添加everyone权限 访问网站时:http://localhost/WebbUploadSample/ZipUpload.aspx “/WebbUploadSample”应用程序中的服务器错误. -------------------------------------------------------------------------------- 访问被拒绝. 说明: 访问服