• 01-Kafaka


    1、Kafka 2 的安装与配置

    1、上传kafka_2.12-1.0.2.tgz到服务器并解压:

            tar -zxf kafka_2.12-1.0.2.tgz  -C /opt

    2、配置环境变量并更新:

    编辑profile配置文件:  vim /etc/profile

    1. #设置kafka的环境变量
    2. export KAFKA_HOME=/opt/kafka_2.12-1.0.2
    3. export PATH=$PATH:$KAFKA_HOME/bin

    重新加载profile文件:        source /etc/profile

    3、在/opt/kafka_2.12-1.0.2目录中输入kafka-按住tab键,如果能调出其他的指令说明我们配置

    profile成功。

    4、配置/opt/kafka_2.12-1.0.2/config中的server.properties文件:

    > Kafka连接Zookeeper的地址:49.234.5.32:2181,后面的 myKafka 是Kafka在Zookeeper中的根节点路径。

    zookeeper.connect=49.234.5.32:2181/mykafka

    > 发消息到kafka,kafka会给你进行一个持久化,存储的目录。

    Log.dir=/var/niko/kafka/kafka-logs

    我们创建这个目录/var/niko/kafka/kafka-logs

    mkdir -p /var/niko/kafka/kafka-logs

    5、启动zookeeper

    进入到/opt/zookeeper-3.4.14/bin目录

    cd /opt/zookeeper-3.4.14/bin

    启动

    zkServer.sh start

    6、验证zookeeper:

    zkServer.sh status

    ZooKeeper JMX enabled by default

    Using config: /opt/zookeeper-3.4.14/bin/../conf/zoo.cfg

    Mode: standalone 说明成功了

    7、启动Kafka:

    进入Kafka安装的bin目录,执行如下命令:

    1. cd  /opt/kafka_2.12-1.0.2/bin
    2. kafka-server-start.sh  ../config/server.properties

    启动成功,可以看到控制台输出的最后一行的started状态:

    [2019-07-31 21:18:53,199] INFO [KafkaServer id=0] started (kafka.server.KafkaServer)

    8、查看Zookeeper的节点

    进入到Zookeeper安装目录的bin目录下

    执行 zkCli.sh

    执行命令ls /    查看所有的子节点

    [mykafka, zookeeper]

    9、此时Kafka是前台模式启动,要停止,使用Ctrl+C。

    10、如果要后台启动 

    进入Kafka安装的bin目录,执行如下命令:

    cd /opt/kafka_2.12-1.0.2/bin

    执行:kafka-server-start.sh -daemon  ../config/server.properties

    11、查看Kafka的后台进程:ps aux | grep kafka

    注意 kafka端口号9092

    2、查看kafka是否启动

    1、输入指令jps,查看kafka是否启动。

    2、在任意目录下以后台的方式启动kafka

    kafka-server-start.sh  -daemon  /opt/kafka_2.12-1.0.2/config/server.properties

    3、再次输入jps,查看kafka是否启动成功。

    3、在Linux使用命令生产与消费(了解)

    3.1、kafka-topics.sh 用于管理主题

    # 列出现有的主题(主题是放在zookeeper的节点上的)

    kafka-topics.sh --list --zookeeper localhost:2181/mykafka

    # 创建主题,该主题包含一个分区,该分区为Leader分区,它没有Follower分区副本

    --partitions    创建的分区个数

    --replication-factor       创建的副本个数,用来实现高可用

    kafka-topics.sh --zookeeper 49.234.5.32:2181/mykafka --create --topic topic_1 --partitions 1  --replication-factor 1

    # 查看指定主题的详细信息

    kafka-topics.sh --zookeeper 49.234.5.32:2181/mykafka --describe --topic topic_1

    输出结果:

    1. Topic:topic_1   PartitionCount:1        ReplicationFactor:1     Configs:
    2. Topic: topic_1  Partition: 0    Leader: 0       Replicas: 0     Isr: 0

    #删除指定主题

    kafka-topics.sh --zookeeper 49.234.5.32:2181/mykafka --delete --topic topic_1

    3.2、kafka-console-producer.sh用于生产消息

    kafka-console-producer.sh --topic topic_1 --broker-list 49.234.5.32:9092

    3.3、kafka-console-consumer.sh用于消费消息

    kafka-console-consumer.sh  --bootstrap-server 49.234.5.32:9092 --topic topic_1

    开启消费者方式二,从头消费,不按照偏移量消费

    kafka-console-consumer.sh --bootstrap-server 49.234.5.32:9092 --topic topic_1 --from-beginning

    注意:先开启消费消息,在开启生成消息,这样有生成消息的时候就可以直接消费了。

    3.4、查看Kafka所有持久化的数据

    进入到我们创建的用来保存持久化数据的目录:

    cd /var/niko/kafka/kafka-logs

    ls  

     有下面的偏移量,说明我们使用kafka成功。

    1.   __consumer_offsets-22  __consumer_offsets-35  __consumer_offsets-48
    2. __consumer_offsets-10      __consumer_offsets-23

    4、kafka发送消息的流程

    5、Maven项目中使用Kafka开发实战(了解)

    1、首先创建一个maven工程,我们将src目录删除,然后再pom.xml文件中设置这个工程的打包方式为pom。

    1.     
    2.     pom

    2、设置工程的maven仓库目录和setting文件。

    3、创建子模块producer-consumer-test01 和 producer-product-test01

    4、在模块producer-consumer-test01的pom.xml文件中导入依赖。

    1. <dependencies>
    2. <dependency>
    3. <groupId>org.apache.kafkagroupId>
    4. <artifactId>kafka-clientsartifactId>
    5. <version>1.0.2version>
    6. dependency>
    7. dependencies>

    5.1、生产者

    消费者生产消息后,需要broker端的确认,可以同步确认,也可以异步确认。

    同步确认效率低,异步确认效率高,但是需要设置回调对象。

    1. package com.wei.producer;
    2. import org.apache.kafka.clients.producer.*;
    3. import org.apache.kafka.common.header.Header;
    4. import org.apache.kafka.common.header.internals.RecordHeader;
    5. import org.apache.kafka.common.serialization.IntegerSerializer;
    6. import org.apache.kafka.common.serialization.StringSerializer;
    7. import java.nio.charset.StandardCharsets;
    8. import java.util.ArrayList;
    9. import java.util.HashMap;
    10. import java.util.Map;
    11. import java.util.concurrent.ExecutionException;
    12. import java.util.concurrent.Future;
    13. public class MyProducer1 {
    14. public static void main(String[] args) throws ExecutionException, InterruptedException {
    15. /**
    16. * 1.1、KafkaProducer 的创建需要指定的参数
    17. * server地址 key的序列化 value的序列化 timeout ack retries
    18. */
    19. Map map = new HashMap<>();
    20. map.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"49.234.5.32:9092");
    21. map.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class);
    22. map.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    23. // map.put("request.timeout.ms",55);
    24. // map.put(ProducerConfig.ACKS_CONFIG,"all"); //有默认值可以不设置
    25. // map.put(ProducerConfig.RETRIES_CONFIG,3); //也可以不设置
    26. // 1、创建发送消息的类对象KafkaProducer
    27. KafkaProducer producer = new KafkaProducer<>(map);
    28. /**
    29. * -String topic: 主题
    30. * -Integer partition: 分区
    31. * -Long timestamp: 时间戳
    32. * -K key: key
    33. * -V value: value
    34. * -Iterable
      headers) :用于设置用户自定义的消息头字段
    35. *
    36. */
    37. // 2.1、创建一个数组,里面存放的都是Header
    38. ArrayList
      headers = new ArrayList<>();
    39. // 添加的是Header接口的实现类RecordHeader 通过构造方法实例化一个RecordHeader对象
    40. headers.add(new RecordHeader("wode.name","wode.value".getBytes(StandardCharsets.UTF_8)));
    41. // 2、使用producerRecord用来给kafka发送封装的消息
    42. ProducerRecord producerRecord = new ProducerRecord<>(
    43. "topic_2",
    44. 0,
    45. 0,
    46. "nihao wudi",
    47. headers
    48. );
    49. // 3、发送消息
    50. // 消费者生产消息后,需要broker端的确认,可以同步确认,也可以异步确认。
    51. // 同步确认效率低,异步确认效率高,但是需要设置回调对象。
    52. // 3.1、同步发送
    53. // final Future future = producer.send(producerRecord);
    54. // final RecordMetadata recordMetadata = future.get();
    55. // System.out.println("主题是:"+recordMetadata.topic());
    56. // System.out.println("分区是:"+recordMetadata.partition());
    57. // System.out.println("变异量是:"+recordMetadata.offset());
    58. // 3.2、异步发送
    59. producer.send(producerRecord, new Callback() {
    60. @Override
    61. public void onCompletion(RecordMetadata metadata, Exception exception) {
    62. if (exception == null) {
    63. System.out.println("消息的主题:" + metadata.topic());
    64. System.out.println("消息的分区号:" + metadata.partition());
    65. System.out.println("消息的偏移量:" + metadata.offset());
    66. } else {
    67. System.out.println("异常消息:" + exception.getMessage());
    68. }
    69. }
    70. });
    71. // 4、关闭producer
    72. producer.close();
    73. }
    74. }

    5.2、消费者

    1. package com.lagou.kafka.demo.consumer;
    2. import org.apache.kafka.clients.consumer.ConsumerConfig;
    3. import org.apache.kafka.clients.consumer.ConsumerRecord;
    4. import org.apache.kafka.clients.consumer.ConsumerRecords;
    5. import org.apache.kafka.clients.consumer.KafkaConsumer;
    6. import org.apache.kafka.common.serialization.IntegerDeserializer;
    7. import org.apache.kafka.common.serialization.StringDeserializer;
    8. import java.util.Arrays;
    9. import java.util.HashMap;
    10. import java.util.Map;
    11. import java.util.function.Consumer;
    12. public class MyConsumer2 {
    13. public static void main(String[] args) {
    14. Map configs = new HashMap<>();
    15. // mac的hosts文件中手动配置域名解析
    16. configs.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "49.234.5.32:9092");
    17. // 使用常量代替手写的字符串,配置key的反序列化器
    18. configs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class);
    19. // 配置value的反序列化器
    20. configs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    21. // 配置消费组ID
    22. configs.put(ConsumerConfig.GROUP_ID_CONFIG, "consumer_demo2");
    23. // 如果找不到当前消费者的有效偏移量,则自动重置到最开始
    24. // latest表示直接重置到消息偏移量的最后一个
    25. configs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
    26. KafkaConsumer consumer = new KafkaConsumer(configs);
    27. // 先订阅,再消费
    28. consumer.subscribe(Arrays.asList("topic_1"));
    29. while (true) {
    30. // 如果主题中没有可以消费的消息,则该方法可以放到while循环中,每过3秒重新拉取一次
    31. // 如果还没有拉取到,过3秒再次拉取,防止while循环太密集的poll调用。
    32. // 批量从主题的分区拉取消息
    33. final ConsumerRecords consumerRecords = consumer.poll(3_000);
    34. // 遍历本次从主题的分区拉取的批量消息
    35. consumerRecords.forEach(new Consumer>() {
    36. @Override
    37. public void accept(ConsumerRecord record) {
    38. System.out.println(record.topic() + "\t"
    39. + record.partition() + "\t"
    40. + record.offset() + "\t"
    41. + record.key() + "\t"
    42. + record.value());
    43. }
    44. });
    45. }
    46. // consumer.close();
    47. }
    48. }

    6、SpringBoot整合 Kafka

    1、首先是创建一个springboot-kafka-sum-demo02项目。

    2、通过快速构建的方式添加spring-web 、spring-kafka或者是手动在pom.xml文件中添加spring-web 、spring-kafka的依赖。

    1. <dependency>
    2. <groupId>org.springframework.bootgroupId>
    3. <artifactId>spring-boot-starter-webartifactId>
    4. dependency>
    5. <dependency>
    6. <groupId>org.springframework.kafkagroupId>
    7. <artifactId>spring-kafkaartifactId>
    8. dependency>

    3、resource目录下的application.properties文件:

    1. #1、设置应用程序的名称和端口号
    2. spring.application.name=springboot-kafka-02
    3. server.port=8080
    4. #2、kafka单体或者集群的host和端口号
    5. spring.kafka.bootstrap-servers=49.234.5.32:9092
    6. #3、producer的配置
    7. spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.IntegerSerializer
    8. spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer
    9. #producer生产者每一批次可以放多少条记录
    10. spring.kafka.producer.batch-size=16384
    11. #生产者端 可以用来发送的缓冲区的大小 32MB 单位是字节
    12. spring.kafka.producer.buffer-memory=33554432
    13. #4、consumer的配置
    14. spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.IntegerDeserializer
    15. spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer
    16. spring.kafka.consumer.group-id=springboot-consumer02
    17. #如果在kafka的消费者中找不到当前的偏移量 从最早的偏移量开始获取数据
    18. spring.kafka.consumer.auto-offset-reset=earliest
    19. #消费者的偏移量是自动提交还是手动提交 设置成true表示是自动提交变异量 如果有事务的情况下 我们通常是设置成手动提交
    20. spring.kafka.consumer.enable-auto-commit=true
    21. #消费者设置成自动提交偏移量的一个提交频率
    22. spring.kafka.consumer.auto-commit-interval=1000

    4、生产者同步发送消息

    1. package com.wei.springbootkafka.producer;
    2. import org.apache.kafka.clients.producer.RecordMetadata;
    3. import org.springframework.beans.factory.annotation.Autowired;
    4. import org.springframework.kafka.core.KafkaTemplate;
    5. import org.springframework.kafka.support.SendResult;
    6. import org.springframework.util.concurrent.ListenableFuture;
    7. import org.springframework.web.bind.annotation.PathVariable;
    8. import org.springframework.web.bind.annotation.RequestMapping;
    9. import org.springframework.web.bind.annotation.RestController;
    10. import java.util.concurrent.ExecutionException;
    11. @RestController
    12. public class MyProducer01 {
    13. // 1、自动注入KafkaTemplate
    14. @Autowired
    15. private KafkaTemplate template;
    16. @RequestMapping("/send/sync/{message}")
    17. public String sendSyncMessage(@PathVariable("message") String message){
    18. // 2、使用KafkaTemplate发送消息
    19. ListenableFuture> future =
    20. template.send("spring-topic-01", 0, 0, message);
    21. // 3、同步发送消息 get()
    22. try {
    23. SendResult result = future.get();
    24. RecordMetadata metadata = result.getRecordMetadata();
    25. System.out.println("主题是:"+metadata.topic());
    26. System.out.println("分区是:"+metadata.partition());
    27. System.out.println("偏移量是:"+metadata.offset());
    28. } catch (Exception e) {
    29. e.printStackTrace();
    30. System.out.println("异常了");
    31. }
    32. return "success";
    33. }
    34. }

    5、生产者异步发送消息 

    1. @RestController
    2. public class MyProducer02 {
    3. // 1、自动注入KafkaTemplate
    4. @Autowired
    5. private KafkaTemplate template;
    6. @RequestMapping("/send/sync/{message}")
    7. public String sendSyncMessage(@PathVariable("message") String message){
    8. // 2、使用KafkaTemplate发送消息
    9. ListenableFuture> future =
    10. template.send("spring-topic-01", 0, 0, message);
    11. // 3、异步发送消息
    12. future.addCallback(new ListenableFutureCallback>() {
    13. @Override
    14. public void onFailure(Throwable ex) {
    15. System.out.println("失败了"+ex.getMessage());
    16. }
    17. @Override
    18. public void onSuccess(SendResult result) {
    19. RecordMetadata recordMetadata = result.getRecordMetadata();
    20. System.out.println("发送消息成功:" + metadata.topic() + "\t"
    21. + metadata.partition() + "\t"
    22. + metadata.offset());
    23. }
    24. });
    25. return "success";
    26. }
    27. }

    6、消费者消费消息

    1. package com.wei.springbootkafka.consumer;
    2. import org.apache.kafka.clients.consumer.ConsumerRecord;
    3. import org.apache.kafka.clients.producer.RecordMetadata;
    4. import org.springframework.beans.factory.annotation.Autowired;
    5. import org.springframework.kafka.annotation.KafkaListener;
    6. import org.springframework.stereotype.Component;
    7. //1、交由spring进行管理
    8. @Component
    9. public class MyConsumer01 {
    10. // 2、KafkaListener监听指定的主题
    11. @KafkaListener(topics="spring-topic-01")
    12. // 3、添加了@KafkaListener注解后 方法中就可以使用ConsumerRecord:用来接收kafka的消息
    13. public void getMessage(ConsumerRecord record){
    14. System.out.println("consumer"
    15. + record.topic()+"\t"
    16. + record.partition()+"\t"
    17. + record.offset()+"\t"
    18. + record.key()+"\t"
    19. + record.value()+"\t"
    20. );
    21. }
    22. }

    7、kafka报错 UnknownHostException 解决方案

    运行springboot和kafka整合项目报错java.net.UnknownHostException: VM-4-7-centos

    解决方案:

    1、cd 到/opt/kafka_2.12-1.0.2/config目录下

    2、vim server.properties

    设置成listeners=PLAINTEXT://VM-4-7-centos:9092

    3、通过查看 linux服务器的 /etc/hosts 文件:将VM-4-7-centos指向的就是linux服务器ip。

    127.0.0.1 VM-4-7-centos VM-4-7-centos

    49.234.5.32  VM-4-7-centos VM-4-7-centos

    4、由于我是在本机服务中访问到了linux服务器上的kafka服务,自然就无法解析到 VM-4-7-centos。因此需要在本机的hosts文件中也加入相应的配置!

    5、Mac系统的hosts 文件就在 /etc/hosts 路径里,我们直接是无法编辑的,需要通过下面的方法来修改我们的 hosts 文件。

    进入终端(命令窗口)里,输入 sudo vi /etc/hosts ,回车后再输入密码,再回车就可以打开我们的hosts文件了。

    添加VM-4-7-centos 的服务器的ip地址:

    49.234.5.32 VM-4-7-centos

  • 相关阅读:
    前端软件快捷键集合
    Spring高手之路——深入理解与实现IOC依赖查找与依赖注入
    git 时忽略某个文件或文件夹
    设计模式:状态模式(C#、JAVA、JavaScript、C++、Python、Go、PHP)
    基于STC12C5A60S2系列1T 8051单片机实现一主单片机与一从单片机相互发送数据的RS485通信功能
    橘子学Flink03之Flink的流处理与批处理
    剑指 Offer 28. 对称的二叉树
    [附源码]计算机毕业设计基于Springboot设备运维平台出入库模块APP
    UI案例——登陆系统
    OSI网络模型与TCP/IP协议
  • 原文地址:https://blog.csdn.net/NikoChina/article/details/132710954