当前位置:首页 > CN2资讯 > 正文内容

【环境部署】windows10构建kafka环境。windows安装kafka调试环境

19小时前CN2资讯


前言:kafka的教程,概念比较多和细,这里不做分享,这里分享下环境搭建和基础demo构建。

【环境】

windows10

JDK1.8

zookeeper 3.5.9

kafka 2.12-2.8.0

注意:如果没有JDK,请先安装JDK

【安装zookeeper】

Zookeeper:

1)    建议下载稳定版。

       下载地址:http:///apache/zookeeper/

2)    下载后解压到一个目录:eg: E:\\env\\zookeeper

3)    在zookeeper目录下,新建文件夹,并命名(eg: data).(路径为:E:\\env\\zookeeper\\data)

4)    进入Zookeeper设置目录,eg: E:\\env\\zookeeper\\conf

       复制“zoo_sample.cfg”副本à并将副本重命名为“zoo.cfg”

       在任意文本编辑器(eg:记事本)中打开zoo.cfg

       找到并编辑dataDir=E:\\env\\zookeeper\\data

5)    添加系统环境变量:

       在系统变量中添加ZOOKEEPER_HOME = E:\\env\\zookeeper

       编辑path系统变量,添加为路径%ZOOKEEPER_HOME%\bin

6)    在zoo.cfg文件中修改默认的Zookeeper端口(默认端口2181),比如修改为12181(我这里避开了hype-v默认的保留端口)

7)    Dos(cmd或者powershell)下运行:zkserver

【安装kafka】

2)    下载后解压缩。eg: E:\\env\\kafka

3)    建立一个空文件夹 logs. eg: E:\\env\\kafka\\logs

4)    进入config目录,编辑 server.properties文件(eg: 用“写字板”打开)。

       找到并编辑log.dirs= E:\\env\\kafka\\logs

       找到并编辑zookeeper.connect=localhost:12181。表示本地运行。 需要配置成和你的zk一样的端口号

       (Kafka会按照默认,在9092端口上运行,并连接zookeeper的默认端口:12181)

运行:请确保在启动Kafka服务器前,Zookeeper实例已经准备好并开始运行。(就是开着Zookeeper窗口不要关)

1)    在 E:\env\kafka(你的kafka路径)下,按住shift+鼠标右键。

       选择“在此处打开Powershell窗口(S)”(如果没有此选项,在此处打开命令窗口)。

2)    运行:.\bin\windows\kafka-server-start.bat .\config\server.properties

3)    可能会报错:“找不到或无法加载主类 ”

4)    解决(3)的办法:

       在kafka安装目录中找到bin\windows目录中的kafka-run-class.bat为%CLASSPATH%加上双引号(可用Matlab打开,并进行搜索)

       修改前:setCOMMAND=%JAVA%%KAFKA_HEAP_OPTS% %KAFKA_JVM_PERFORMANCE_OPTS% %KAFKA_JMX_OPTS%%KAFKA_LOG4J_OPTS% -cp%CLASSPATH% %KAFKA_OPTS% %*   

       修改后:SetCOMMAND=%JAVA%%KAFKA_HEAP_OPTS% %KAFKA_JVM_PERFORMANCE_OPTS% %KAFKA_JMX_OPTS%%KAFKA_LOG4J_OPTS% -cp"%CLASSPATH%"%KAFKA_OPTS% %*

5)    再次运行:.\bin\windows\kafka-server-start.bat.\config\server.properties

【Demo】(这里引用网上的一个简单的例子,没有使用Boot,使用boot的话例子回头开一篇详尽的)

依赖:


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


import java.util.Properties; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; public class ProducerSend { public static void main(String args[]) { //1.参数配置:端口、缓冲内存、最大连接数、key序列化、value序列化等等(不是每一个非要配置) Properties props=new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("acks", "all"); props.put("retries", 0); props.put("batch.size", 16384); props.put("linger.ms", 1); props.put("buffer.memory", 33554432); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); //2.创建生产者对象,并建立连接 Producer<String, String> producer = new KafkaProducer<String,String>(props); try { //3.在my-topic主题下,发送消息 for (int i = 0; i < 10000; i++) { System.out.println(Integer.toString(i)); producer.send(new ProducerRecord<String, String>("my-topic", Integer.toString(i), Integer.toString(i))); Thread.sleep(500); } } catch (Exception e) { System.out.println("ERROR"); } //4.关闭 producer.close(); } }

 

package com.arcvideo.kafka.service; import java.util.Arrays; import java.util.Properties; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; /** * @author martin * @version 1.0 * @name: ConsumerReceive * @date: 2021/5/23 22:55 * @description * @comepony **/ public class ConsumerReceive { public static void main(String args[]) { //1.参数配置:不是每一非得配置 Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("", "1000"); //因为每一个消费者必须属于某一个消费者组,所以必须还设置group.id props.put("group.id", "test1"); props.put("enable.auto.commit", "true"); props.put("session.timeout.ms", "30000"); props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); //2.创建消费者对象,并建立连接 KafkaConsumer<String, String> consumer = new KafkaConsumer<String,String>(props); //3.设置从"my-topic"主题下拿取数据 consumer.subscribe(Arrays.asList("my-topic")); //4.消费数据 while (true) { //阻塞时间,从kafka中取出100毫秒的数据,有可能一次性去除0-n条 ConsumerRecords<String, String> records = consumer.poll(100); //遍历 for (ConsumerRecord<String, String> record : records) //打印结果 //System.out.printf("offset = %d, key = %s, value = %s", record.offset(), record.key(), record.value()); System.out.println("消费者消费的数据为:"+record.value()); } } }

     

 

    你可能想看:

    扫描二维码推送至手机访问。

    版权声明:本文由皇冠云发布,如需转载请注明出处。

    本文链接:https://www.idchg.com/info/32530.html

    分享给朋友:

    “【环境部署】windows10构建kafka环境。windows安装kafka调试环境” 的相关文章

    AWS VPS Free: 如何利用AWS Free Tier免费服务轻松构建云计算项目

    当我第一次接触AWS (亚马逊网络服务) 的时候,最吸引我的就是他们提供的各种免费的VPS服务。AWS的VPS免费服务实际上是一种叫做AWS Free Tier的计划,它允许用户在一定条件下使用AWS的多种服务而无需支付费用。这项计划的意义在于,它为刚入门的开发者和小型企业提供了一个绝佳的机会,能够...

    如何找到便宜的域名并有效管理

    在了解便宜域名之前,首先我们要对“域名”这个概念有个清晰的认识。域名其实就是互联网上某个特定网站的地址。当我在浏览器中输入一个域名,比如“example.com”,就能顺利地找到对应的网站。简单来说,域名是你在网上的身份标识。而它的作用不仅是让用户更容易记住和访问你的网站,还能提升你品牌的可信度。...

    如何实现Windows链接服务器的应用与配置

    在现代工作和生活中,远程连接的重要性日益凸显。Windows链接服务器作为一种强大的工具,帮助用户在不同的设备之间实现无缝的远程访问。它的定义其实就是这样一款可以让用户通过网络访问和管理远程Windows服务器的技术。这意味着无论是在办公室还是在家中,只要有网络连接,我都能方便地使用和维护我的服务器...

    Fiberstate: 让你的健康饮食更简单轻松

    每当提起“Fiberstate”这个词,我的脑海中便浮现出一个全新的健康理念。Fiberstate并不仅仅是一个产品,它代表着对身体健康全面理解的汇聚,是我提升生活质量的伴侣。从简单的定义开始探讨,这个概念到底蕴含着什么深意。 Fiberstate的定义实际上很贴近我们的日常生活。简单来说,它是指通...

    解决Hostodo官网无法打开的问题的有效方法

    在使用 Hostodo 官网时,偶尔会遇到无法打开的情况。这种情况可能让人感到无助,尤其是当你迫切需要访问相关信息时。让我来分享一些常见原因,帮助你更好地理解。 首先,服务器的维护或故障是一个普遍的原因。当网站进行定期更新或修复时,服务器可能会暂时不可用。通常,官方会提前通知用户,然而,有时我们无法...

    Hostwinds评测:全面解析优秀的网络托管服务

    Hostwinds概述 在了解Hostwinds之前,首先想分享一下我对这家公司的印象。Hostwinds成立于2010年,作为一家相对年轻的网络托管服务提供商,虽然起步不久,但它的发展速度却让我感到惊叹。起初,Hostwinds仅是一家提供基本虚拟主机服务的小公司,随着需求的不断增长,他们逐步扩展...