• linkedu视频
  • 平面设计
  • 电脑入门
  • 操作系统
  • 办公应用
  • 电脑硬件
  • 动画设计
  • 3D设计
  • 网页设计
  • CAD设计
  • 影音处理
  • 数据库
  • 程序设计
  • 认证考试
  • 信息管理
  • 信息安全
菜单
linkedu.com
导航菜单
  • 网页制作
  • 数据库
  • 程序设计
  • 操作系统
  • CMS教程
  • 游戏攻略
  • 脚本语言
  • 平面设计
  • 软件教程
  • 网络安全
  • 电脑知识
  • 服务器
  • 视频教程
  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx
您的位置:首页 > 服务器 >云计算 > kafka集群搭建,kafka集群

kafka集群搭建,kafka集群

作者:网友 字体:[增加 减小] 来源:互联网

本文主要包含kafka集群搭建,kafka集群,kafka集群安装,kafka集群配置,kafka集群部署等服务器相关知识,网友希望可以进行参考

kafka集群搭建,kafka集群


第一步

先去官网下载 kafka_2.9.2-0.8.1.1.tgz 并解压再进入到安装目录(也可以自己配置路径,方法跟配置java、hadoop等路径是一样的).
> tar -xzf kafka_2.9.2-0.8.1.1.tgz 
> cd kafka_2.9.2-0.8.1.1


第二步
zeekeeper集群搭建(用的是kafka自带的zeekeeper,一共准备了三台机器)
1、关闭各台机器的防火墙(一定要切记,我搭建的时候以为能ping通就ok了,就没关心防火墙的问题了,最后白白浪费了一天的时间)
命令 /ect/init.d/iptables stop

2、进入到打开/ect下的hosts文件
修改为
127.0.0.1 localhost
10.61.5.66 host1
10.61.5.67 host2
10.61.5.68 host3
(ip和机器名根据个人实际情况修改)

3、修改zeekeeper 配置文件
进入到kafka安装目录下的config文件,打开zookeeper.properties
修改dataDir={kafka安装目录}/zookeeper/logs/
注释掉maxClientCnxns=0
在文件末尾添加如下语句
tickTime=2000
initLimit=5
syncLimit=2
#host1、2、3为主机名,可以根据实际情况更改,端口号也可以更改
server.1=host1:2888:3888
server.2=host2:2888:3888
server.3=host3:2888:3888

4、在dataDir目录下的建立一个myid文件
命令   echo 1 >myid
另外两台机子分别设置为2、3,依次类推。


第三步
启动zookeeper服务(每台机子的zeekeeper都要启)
> bin/zookeeper-server-start.sh config/zookeeper.properties
在三台机子的zeekeeper都启动好之前,先启动的机子会有错误日志,这是正常的

 

第四步
配置kafka
1、在kafka安装目录下的config目录下打开server.properties文件
修改
zookeeper.connect=host1:2181,host2:2181,host3:2181    (2181为端口号,可以根据自己的实际情况更改)
其他两台机子的server.properties文件中的broker.id也要改,反正三台机子的broker.id不能有重复
2、修改producer.properties文件
修改
metadata.broker.list=host1:9092,host2:9092,host3:9092
prodeucer.type=async

3、修改consumer.properties文件
修改
zeekeeper.connect=host1:2181,host2:2181,host3:2181

4、在每台机子启动kafka服务
> bin/kafka-server-start.sh config/server.properties

 

第四步:建立一个主题
> bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 3 --partitions 6 --topic my-replicated-test
factor大小不能超过broker数

通过以下命令查看主题
> bin/kafka-topics.sh --list --zookeeper host1:2181 (也可以是host2:2181等)
my-replicated-test

通过下述命令可以看到该主题详情
> bin/kafka-topics.sh --describe --zookeeper host1:2181 --topic my-replicated-test


第五步:发送消息
在host2上建立生产者角色,并发送消息(其实可以是三台机子中的任何一台)
> bin/kafka-console-producer.sh --broker-list host1:9092 --topic my-replicated-test 
This is a message
This is another message

在host3上建立消费者角色(在该终端窗口内可以看到生产者发布这消息)
> bin/kafka-console-consumer.sh --zookeeper host1:2181 --topic my-replicated-test --from-beginning
This is a message
This is another message

至此,一个kafka集群就搭好了,可以作为kafka服务器了

 

 测试程序(在win系统上)

切记要去C:\Windows\system32\drivers\etc\hosts作如下配置,否则测试程序无法访问kafka服务器!

10.61.5.66 host1
10.61.5.67 host2
10.61.5.68 host3

记得将kafka安装目录下libs里的所有包导入项目里去

//生产者测试程序

public class ProducerTest {
 public static void main(String[] args) throws FileNotFoundException {  
        Properties props = new Properties();  
        props.put("zookeeper.connect", "slaves7:2182,slaves8:2182,slaves9:2182");  
        props.put("serializer.class", "kafka.serializer.StringEncoder");  
        props.put("metadata.broker.list","slaves7:9092,slaves8:9092,slaves9:9092");
      
        ProducerConfig config = new ProducerConfig(props);  
        Producer<String, String> producer = new Producer<String, String>(config);  
         File file=new File("E:/test","test.txt");
         BufferedReader readtxt=new BufferedReader(new FileReader(file));
          String line=null;
          byte[] item=null;
   try {
      while((line=readtxt.readLine())!=null){
      item=line.getBytes();
      String str = new String(item);
      System.out.println(str);
      producer.send(new KeyedMessage<String, String>("my-replicated-topic",str));
      }
   } catch (IOException e) {
    e.printStackTrace();
   }
       }  
}
//消费者测试程序

public class ConsumerTest extends Thread {
 private final ConsumerConnector consumer;  
    private final String topic;  
      
    public static void main(String[] args) {  
        ConsumerTest consumerThread = new ConsumerTest("my-replicated-topic");  
        consumerThread.start();  
    }  
      
    public ConsumerTest(String topic) {  
     System.out.println(topic);
        consumer = kafka.consumer.Consumer  
                .createJavaConsumerConnector(createConsumerConfig());  
        this.topic = topic;  
    }  
      
    private static ConsumerConfig createConsumerConfig() {  
        Properties props = new Properties();  
        props.put("zookeeper.connect", "slaves7:2182,slaves8,slaves9:2182");  
        props.put("group.id", "0");  
        props.put("zookeeper.session.timeout.ms", "400000");  
        props.put("zookeeper.sync.time.ms", "200");  
        props.put("auto.commit.interval.ms", "1000");  
      
        return new ConsumerConfig(props);  
      
    }  
      
    public void run() {  
        Map<String, Integer> topicCountMap = new HashMap<String, Integer>();  
     

分享到:QQ空间新浪微博腾讯微博微信百度贴吧QQ好友复制网址打印

您可能想查找下面的文章:

  • kafka集群搭建,kafka集群
  • Kafka集群安装,kafka集群

相关文章

  • 开源微内核seL4 microkernel,sel4microkernel
  • Hadoop自测题,hadoop
  • 智能助手相伴,享受人机温情,智能助手人机
  • Map/Reduce原理
  • NOSQL(五)版本戳,nosql版本戳
  • MapReduce编程之数据去重,mapreduce编程数据
  • Zoho产品经理谈中小企业面临的挑战及应对之策,zoho产品经理
  • MapReduce处理输出多文件格式(MultipleOutputs),mapreduce多文件输出
  • openstack中Nova组件images的所有python API 汇总,openstacknova
  • 与Greenplum度过的三个星期,greenplum三个星期

文章分类

  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx

最近更新的内容

    • 二分Kmeans的java实现,二分kmeansjava
    • 现代数学的引路人,现代数学引路人
    • KVM 【SNAT/DNAT2种配置实现以及扁平化网络模式(flat)实现/virsh2种动态迁移实现】,kvmsnat
    • A Note on Distributed Computing,anoteondialectic
    • hadoop学习笔记(三)——WIN7+eclipse+hadoop2.5.2部署,hadoopwin7
    • OpenStack之swift安装笔记,openstackswift
    • Storm计算结果是如何存放的,Storm计算结果存放
    • Hadoop常见重要命令行操作及命令作用,hadoop命令行命令
    • libvirt 部分API 介绍,libvirtapi
    • Storm On Yarn部署,stormonyarn部署

关于我们 - 联系我们 - 免责声明 - 网站地图

©2020-2025 All Rights Reserved. linkedu.com 版权所有