第一章.kafka概述

1.定义

  1. - Kafka是一个分布式的基于发布/订阅模式的消息队列,主要应用于大数据实时处理领域
  2. - Kafka是一个开源的分布式事件流平台,被数千家公司用于高性能数据管道,流分析.数据集成和关键任务应用

2.消息队列

$01[Kafka概述_入门] - 图1

  1. # 使用消息队列的好处
  2. 1. 解耦:允许你独立的扩展或修改两边的处理过程,只要确保他们遵守同样的接口约束
  3. 2. 可恢复性: 系统的一部分组件失效时,不会影响到整个系统,消息队列降低了进程之间的耦合度,所以即使一个处理消息的进程挂掉,加入队列中的消息仍然可以在系统恢复后被处理
  4. 3. 缓冲: 有助于控制和优化数据流经过系统的速度,解决生产消息和消费消息的处理速度不一致的情况
  5. 4. 灵活性&峰值处理能力: 在访问量剧增的情况下,应用仍然需要继续发挥作用,但是这样的突发流量并不常见,如果为以能处理这类峰值访问为标准来投入资源随时待命无疑是巨大的浪费,使用消息队列能够使关键组件顶住突发的访问能力,而不会因为突发的超负荷的请求而完全崩溃
  6. 5. 异步通信: 很多时候,用户不想也不需要立即处理消息,消息队列提供了异步处理机制,允许用户把一个消息放入队列,但并不立即处理它,想让队列中放入多少消息就放多少,然后在需要的时候再去处理他们
  1. # 消息队列的两种模式
  2. 1. 点对点模式(一对一,消费者主动拉取数据,消息收到后消息清除)
  3. 2. 发布订阅模式(一对多,消费者消费数据之后不会清除消息)

Kafka基础架构

$01[Kafka概述_入门] - 图2

  1. - Producer: 消息生产者,就是向Kafka broker 发消息的客户端
  2. - Consumer: 消息消费者,向kafka broker 取消息的客户端
  3. - Consumer Group(CG) : 消费者组,由多个consumer组成,消费者组内每个消费者负责消费不同分区的数据,一个分区只能由一个组内消费者消费,消费者组之间互不影响,所有的消费者都属于某个消费者组,,即消费者组是逻辑上的一个订阅者
  4. - Broker: 一台Kafka服务器就是一个broker,一个集群由多个broker组成,一个broker可以容纳多个topic
  5. - Topic: 可以理解为一个队列,生产者和消费者面向的都是一个Topic
  6. - Partition: 为了实现扩展性,一个非常大的topic可以分布在多个broker(服务器)上,一个topic可以分为多个partition,每个partition是一个有序的队列
  7. - Replica: 副本,为保证集群中某个节点发生故障时,该节点上的partition数据不丢失,且Kafka仍然能够继续工作,Kafka提供了副本机制,一个topic的每个分区都有若干个副本,一个leader和若干个follower
  8. - leader: 每个分区中多个副本的"主",生产者发送数据的对象,以及消费者消费数据的对象都是leader
  9. - follower: 每个分区多个副本的"从",实时从leader中同步数据,保持和leader数据的同步,leader发生故障时,某个follower会成为新的leader

第二章.Kafka快速入门

1.集群规划

hadoop102 hadoop103 hadoop104
zookeeper zookeeper zookeeper
kafka kafka kafka

Kafka下载:http://kafka.apache.org/downloads.html

2.集群部署

  1. 上传压缩包到linux上并解压
  1. tar -zxvf kafka_2.11-2.4.1.tgz -C /opt/module
  1. 修改解压后的文件名称
  1. mv kafka_2.11-2.4.1 kafka-2.4.1
  1. 修改配置文件
  1. cd config/
  2. sudo vim server.properties
  3. # 修改以下内容
  4. #broker的全局唯一编号,不能重复
  5. broker.id=2
  6. #kafka运行日志(数据)存放的路径
  7. log.dirs=/opt/module/kafka-2.4.1/datas
  8. #配置连接Zookeeper集群地址
  9. zookeeper.connect=hadoop102:2181,hadoop103:2181,hadoop104:2181/kafka
  1. 配置环境变量
  1. sudo vim /etc/profile.d/my_env.sh
  1. # KAFKA_HOME
  2. export KAFKA_HOME=/opt/module/kafka-2.4.1
  3. export PATH=$PATH:$KAFKA_HOME/bin
  1. xsync kafka-2.4.1
  1. 分别在hadoop103和hadoop104上修改配置文件 server.properties
  1. # hadoop103
  2. broker.id = 3
  3. # hadoop104
  4. broker.id = 4
  1. 启动集群
  1. 1. 先启动Zookeeper集群(脚本启动)
  2. zkCluster.sh start
  3. 2. 依次在hadoop102,hadoop103,hadoop104节点上启动Kafka
  4. [atguigu@hadoop102 kafka-2.4.1]$ bin/kafka-server-start.sh -daemon config/server.properties
  5. [atguigu@hadoop103 kafka-2.4.1]$ bin/kafka-server-start.sh -daemon config/server.properties
  6. [atguigu@hadoop104 kafka-2.4.1]$ bin/kafka-server-start.sh -daemon config/server.properties
  7. 3. 关闭集群
  8. [atguigu@hadoop102 kafka-2.4.1]$ bin/kafka-server-stop.sh stop
  9. [atguigu@hadoop103 kafka-2.4.1]$ bin/kafka-server-stop.sh stop
  10. [atguigu@hadoop104 kafka-2.4.1]$ bin/kafka-server-stop.sh stop
  1. Kafka群起脚本(kafka.sh)
  1. #! /bin/bash
  2. if (($#==0)); then
  3. echo -e "请输入参数:\n start 启动kafka集群;\n stop 停止kafka集群;\n" && exit
  4. fi
  5. case $1 in
  6. "start")
  7. for host in hadoop103 hadoop102 hadoop104
  8. do
  9. echo "---------- $1 $host 的kafka ----------"
  10. ssh $host "/opt/module/kafka-2.4.1/bin/kafka-server-start.sh -daemon /opt/module/kafka-2.4.1/config/server.properties"
  11. done
  12. ;;
  13. "stop")
  14. for host in hadoop103 hadoop102 hadoop104
  15. do
  16. echo "---------- $1 $host 的kafka ----------"
  17. ssh $host "/opt/module/kafka-2.4.1/bin/kafka-server-stop.sh /opt/module/kafka-2.4.1/config/server.properties"
  18. done
  19. ;;
  20. *)
  21. echo -e "---------- 请输入正确的参数 ----------\n"
  22. echo -e "start 启动kafka集群;\n stop 停止kafka集群;\n" && exit
  23. ;;
  24. esac

添加可执行权限

  1. chmod +x kafka.sh

3.Kafka命令行操作

  1. # 查看当前服务器中的所有topic
  2. kafka-topics.sh --zookeeper hadoop102:2181/kafka --list
  3. kafka-topics.sh --bootstrap-server hadoop102:9092 --list
  4. # 创建topic
  5. kafka-topics.sh --bootstrap-server hadoop102:9092 --create --replication-factor 2 --partitions 1 --topic atguigu
  6. # 删除topic
  7. kafka-topics.sh --bootstrap-server hadoop102:9092 --delete --topic atguigu
  8. # 发送消息
  9. kafka-console-producer.sh --broker-list hadoop102:9092 --topic first
  10. # 消费消息
  11. kafka-console-consumer.sh --bootstrap-server hadoop102:9092 --from-beginning --topic first
  12. # 查看某个Topic的详情
  13. kafka-topics.sh --bootstrap-server hadoop102:9092 --describe --topic first
  14. # 修改分区数
  15. kafka-topics.sh --bootstrap-server hadoop102:9092 --alter –-topic first --partitions 3

第三章.Kafka架构深入

1.kafka工作流程及文件存储机制

$01[Kafka概述_入门] - 图3

  1. - kafka中消息是以topic进行分类的,生产者生产消息,消费者消费消息,都是面向topic的
  2. - 一个topic下的每一个分区都单独维护一个offset,所以分发到不同分区中的数据是不同的数据,消费者的分区维护是一个消费者组一个主题的一个分区维护一个offset
  3. - topic是逻辑上的概念,而partition是物理上的概念,每个partition对应于一个log文件,该log文件中存储的就是producer生产的数据,Producer生产的数据会被不断追加到该log文件末端,且每条数据都有自己的offect,消费者组中的每一个消费者,都会实时记录自己消费到了哪个offset,以便出错恢复时,从上次的位置继续消费

$01[Kafka概述_入门] - 图4

  1. 由于生产者生产的消息会不断追加到log文件末尾,为防止log文件过大导致数据定位效率低下,Kafka采取了分片和索引机制,将每个partition分为多个segment,每个segent对应两个文件--.index文件和.log文件,这些文件位于一个文件夹下,该文件夹的命名规则是:topic 名称 + 分区序号,例如,first这个topic有三个分区,则其对应的文件夹为first-0,first-1,first-2
  1. 00000000000000000000.index
  2. 00000000000000000000.log
  3. 00000000000000170410.index
  4. 00000000000000170410.log
  5. 00000000000000239430.index
  6. 00000000000000239430.log
  1. index和log文件以当前segment的第一条消息的offset命名,下图为index文件和log文件的结构示意图图

$01[Kafka概述_入门] - 图5

2.kafka生产者

2.1消息发送流程

  1. kafka的Producer发送消息采用的是异步发送的方式,在消息发送的过程中,涉及到了两个线程--main线程和Sender线程,以及一个线程共享变量--RecordAccumulator,main线程将消息发送给RecordAccumulator,Sender线程不断从RecordAccumulator中拉取消息发送到Kafka broker

$01[Kafka概述_入门] - 图6

  1. # 相关参数
  2. batch.size:只有数据积累到batch.size之后,sender才会发送数据
  3. linger.ms:如果数据迟迟未达到batch.size,sender等待linger.time之后就会发送数据

2.2生产者-简单的生产者

  1. 新建一个maven工程,并导入依赖
  1. <dependencies>
  2. <dependency>
  3. <groupId>org.apache.kafka</groupId>
  4. <artifactId>kafka-clients</artifactId>
  5. <version>2.4.1</version>
  6. </dependency>
  7. </dependencies>
  1. 编写代码

需要用到的类:

  • KafkaProducer : 需要创建一个生产者对象,用来发送数据
  • ProducerConfig : 获取所需的一系列配置参数
  • ProducerRecord: 每条数据都要封装成一个ProducerRecord对象
  1. package com.atguigu.kafka.producer;
  2. import org.apache.kafka.clients.producer.KafkaProducer;
  3. import org.apache.kafka.clients.producer.ProducerConfig;
  4. import org.apache.kafka.clients.producer.ProducerRecord;
  5. import java.util.Properties;
  6. import java.util.concurrent.TimeUnit;
  7. public class CustomProducer {
  8. public static void main(String[] args) {
  9. //1.创建kafka生产者的配置对象
  10. Properties properties = new Properties();
  11. //2.给kafka配置对象添加配置信息
  12. properties.put("bootstrap.servers","hadoop102:9092");
  13. properties.put(ProducerConfig.BATCH_SIZE_CONFIG,16384);
  14. properties.put(ProducerConfig.LINGER_MS_CONFIG,10);
  15. properties.put(ProducerConfig.BUFFER_MEMORY_CONFIG,33554432);
  16. //key,value序列化
  17. properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer");
  18. properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
  19. //初始化生产者
  20. KafkaProducer<String, String> producer = new KafkaProducer<>(properties);
  21. ProducerRecord<String,String> producerRecord = null;
  22. for(int i =1;i<10;i++){
  23. //将数据包装为ProducerRecord
  24. producerRecord = new ProducerRecord<String,String>(
  25. "first",
  26. 0,
  27. "aa" + i,
  28. "atguigu" + i
  29. );
  30. //发送数据
  31. producer.send(producerRecord);
  32. }
  33. try{
  34. TimeUnit.SECONDS.sleep(1);
  35. }catch(Exception e){
  36. e.printStackTrace();
  37. }
  38. }
  39. }

$01[Kafka概述_入门] - 图7