## 消息队列

## zookeeper+kafka

工作原理：

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130946.png)

#### 

**Producer：** 生产者，发送消息的一方。生产者负责创建消息，然后将其发送到 Kafka。
**Consumer：** 消费者，接受消息的一方。消费者连接到 Kafka 上并接收消息，进而进行相应的业务逻辑处理。
**Consumer Group：** 一个消费者组可以包含一个或多个消费者。使用多分区 + 多消费者方式可以极大提高数据下游的处理速度，同一消费组中的消费者不会重复消费消息，同样的，不同消费组中的消费者消息消息时互不影响。Kafka 就是通过消费组的方式来实现消息 P2P 模式和广播模式。
**Broker：** 服务代理节点。Broker 是 Kafka 的服务节点，即 Kafka 的服务器。
**Topic：** Kafka 中的消息以 Topic 为单位进行划分，生产者将消息发送到特定的 Topic，而消费者负责订阅 Topic 的消息并进行消费。
**Partition：** Topic 是一个逻辑的概念，它可以细分为多个分区，每个分区只属于单个主题。同一个主题下不同分区包含的消息是不同的，分区在存储层面可以看作一个可追加的日志（Log）文件，消息在被追加到分区日志文件的时候都会分配一个特定的偏移量（offset）。
**Offset：** offset 是消息在分区中的唯一标识，Kafka 通过它来保证消息在分区内的顺序性，不过 offset 并不跨越分区，也就是说，Kafka 保证的是分区有序性而不是主题有序性。
**Replication：** 副本，是 Kafka 保证数据高可用的方式，Kafka 同一 Partition 的数据可以在多 Broker 上存在多个副本，通常只有主副本对外提供读写服务，当主副本所在 broker 崩溃或发生网络异常，Kafka 会在 Controller 的管理下会重新选择新的 Leader 副本对外提供读写服务。
**Record：** 实际写入 Kafka 中并可以被读取的消息记录。每个 record 包含了 key、value 和 timestamp。



#### 生产者-消费者

`生产者`-`消费者`是一种设计模式，`生产者`和`消费者`之间通过添加一个`中间组件`来达到解耦。`生产者`向`中间组件`生成数据，`消费者`消费数据。

就像 65 哥读书时给小芳写情书，这里 65 哥就是`生产者`，情书就是`消息`，小芳就是`消费者`。但有时候小芳不在，或者比较忙，65 哥也比较害羞，不敢直接将情书塞小芳手里，于是将情书塞在小芳抽屉中。所以抽屉就是这个`中间组件`。

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130079.png)

在程序中我们通常使用`Queue`来作为这个`中间组件`。可以使用多线程向队列中写入数据，另外的消费者线程依次读取队列中的数据进行消费。模型如下图所示：

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130065.png)

`生产者`-`消费者`模式通过添加一个中间层，不仅可以解耦生产者和消费者，使其易于扩展，还可以异步化调用、缓冲消息等。

#### 分布式队列

后来 65 哥和小芳异地了，65 哥在`卷都`奋斗，小芳在`魔都`逛街。于是只能通过`邮局`寄暧昧信了。这样 65 哥、邮局和小芳就成了`分布式`的了。65 哥将信件发给邮局，小芳从邮局拿到 65 哥写的信，再回去慢慢看。

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130071.png)

Kafka 的消息`生产者`就是`Producer`，上游消费者进程添加 Kafka Client 创建 Kafka Producer，向 Broker 发送消息，Broker 是集群部署在远程服务器上的 Kafka Server 进程，下游消费者进程引入 Kafka Consumer API 持续消费队列中消息。

因为 Kafka Consumer 使用 Poll 的模式，需要 Consumer 主动拉去消息。所有小芳只能定期去邮局拿信件了(呃，果然主动权都在小芳手上啊)。

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130019.png)

#### 主题

邮局不能只为 65 哥服务，虽然 65 哥一天写好几封信。但也无法挽回邮局的损失。所以邮局是可以供任何人寄信。只需要寄信人写好地址(主题)，邮局建有两地的通道就可以发收信件了。

Kafka 的 Topic 才相当于一个队列，Broker 是所有队列部署的机器。可以按业务创建不同的 Topic，Producer 向所属业务的 Topic 发送消息，相应的 Consumer 可以消费并处理消息。

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130096.png)

#### 分区

由于 65 哥写的信太多，一个邮局已经无法满足 65 哥的需求，邮政公司只能多建几个邮局了，65 哥将信件按私密度分类(分区策略)，从不同的邮局寄送。

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130384.png)

同一个 Topic 可以创建多个分区。理论上分区越多并发度越高，Kafka 会根据分区策略将分区尽可能均衡的分布在不同的 Broker 节点上，以避免消息倾斜，不同的 Broker 负载差异太大。分区也不是越多越好哦，毕竟太多邮政公司也管理不过来

#### 副本

为防止由于邮局的问题，比如交通断啦，邮车没油啦。导致 65 哥的暧昧信无法寄到小芳手上，使得 65 哥晚上远程跪键盘。邮局决定将 65 哥的信件复制几份发到多个正常的邮局，这样只要有一个邮局还在，小芳就可以收到 65 哥的信了。

Kafka 采用分区副本的方式来保证数据的高可用，每个分区都将建立指定数量的副本数，kakfa 保证同一分区副本尽量分布在不同的 Broker 节点上，以防止 Broker 宕机导致所有副本不可用。Kafka 会为分区的多个副本选举一个作为主副本(Leader)，主副本对外提供读写服务，从副本(Follower)实时同步 Leader 的数据。

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130487.png)

#### 多消费者

哎，65 哥的信件满天飞，小芳天天跑邮局，还要一一拆开看，65 哥写的信又臭又长，让小芳忙得满身大汗。于是小芳啪的一下，很快啊，变出多个分身去不同的邮局取信，这样小芳终于可以挤出额外的时间逛街了。

#### 广播消息

邮局最近提供了定制明信片业务，每个人都可以设计明信片，同一个身份只能领取一种明信片。65 哥设计了一堆，广播给所有漂亮的小妹妹都可以来领取，美女啪变出的分身也可以来领取，但是同一个身份的多个分身只能取一种明信片。

Kafka 通过 Consumer Group 来实现广播模式消息订阅，即不同 group 下的 consumer 可以重复消费消息，相互不影响，同一个 group 下的 consumer 构成一个整体。

**最后我们完成了 Kafka 的整体架构，如下：**

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130416.png)

## Zookeeper

Zookeeper 是一个成熟的分布式协调服务，它可以为分布式服务提供分布式配置服务、同步服务和命名注册等能力.。对于任何分布式系统，都需要一种协调任务的方法。Kafka 是使用 ZooKeeper 而构建的分布式系统。但是也有一些其他技术（例如 Elasticsearch 和 MongoDB）具有其自己的内置任务协调机制。

Kafka 将 Broker、Topic 和 Partition 的元数据信息存储在 Zookeeper 上。通过在 Zookeeper 上建立相应的数据节点，并监听节点的变化，Kafka 使用 Zookeeper 完成以下功能：

- Kafka Controller 的 Leader 选举
- Kafka 集群成员管理
- Topic 配置管理
- 分区副本管理

我们看一看 Zookeeper 下 Kafka 创建的节点，即可一目了然的看出这些相关的功能。

![图片](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120130535.png)







###  实验环境 

| 主机名 | ip              |
| ------ | --------------- |
| kafka1 | 192.168.110.175 |
| kafka2 | 192.168.110.176 |
| kafka3 | 192.168.110.177 |



###  修改host 

```
[root@kafka1 opt]# vim /etc/hosts
127.0.0.1   localhost localhost.localdomain localhost4 localhost4.localdomain4
::1         localhost localhost.localdomain localhost6 localhost6.localdomain6
192.168.110.175 kafka1
192.168.110.176 kafka2
192.168.110.177 kafka3

[root@kafka1 opt]# scp /etc/hosts 192.168.110.176:/etc/hosts
[root@kafka1 opt]# scp /etc/hosts 192.168.110.177:/etc/hosts
```

### 配置java环境

三台机子都需要

过程略



### zookeeper单点启动（测试）

#### 1.下载

zookeeper官网：https://www.apache.org/

版本 3.7

三台机子都下载

```
cd /opt/ && yum -y install wget && wget --no-check-certificate https://downloads.apache.org/zookeeper/zookeeper-3.7.1/apache-zookeeper-3.7.1-bin.tar.gz
```

#### 2.单点启动

```
[root@kafka1 opt]# tar -xvf apache-zookeeper-3.7.1-bin.tar.gz
[root@kafka1 opt]# cd apache-zookeeper-3.7.1-bin
[root@kafka1 apache-zookeeper-3.7.1-bin]# ll
总用量 36
drwxr-xr-x 2 1000 1000  4096 3月  17 2021 bin
drwxr-xr-x 2 1000 1000    77 3月  17 2021 conf
drwxr-xr-x 5 1000 1000  4096 3月  17 2021 docs
drwxr-xr-x 2 root root  4096 4月  17 17:15 lib
-rw-r--r-- 1 1000 1000 11358 3月  17 2021 LICENSE.txt
-rw-r--r-- 1 1000 1000   432 3月  17 2021 NOTICE.txt
-rw-r--r-- 1 1000 1000  2214 3月  17 2021 README.md
-rw-r--r-- 1 1000 1000  3570 3月  17 2021 README_packaging.md
[root@kafka1 apache-zookeeper-3.7.1-bin]# cd conf/
[root@kafka1 conf]# cp zoo_sample.cfg zoo.cfg
[root@kafka1 conf]# grep -Ev "#|^$" zoo.cfg   # 修改如下，根据业务需求修改
tickTime=2000
initLimit=10
syncLimit=5
dataDir=/opt/zookeeperkeeper    # 存储快照的目录，不要默认
clientPort=2181
maxClientCnxns=128   # 客户最大连接数

[root@kafka1 opt]# mv /opt/apache-zookeeper-3.7.1-bin /usr/local/zookeeper
[root@kafka1 opt]# mkdir /opt/zookeeperkeeper

[root@kafka1 conf]# /usr/local/zookeeper/bin/zkServer.sh start 
ZooKeeper JMX enabled by default
Using config: /opt/apache-zookeeper-3.7.1-bin/bin/../conf/zoo.cfg
Starting zookeeper ... STARTED
#出现STARTED表示启动成功,出现FAILED TO START要检查配置文件是否有问题

# 检查2181端口
[root@kafka1 opt]# netstat -aulntp|grep 2181
```

#### 3.单点测试

```
[root@kafka1 opt]# /usr/local/zookeeper/bin/zkCli.sh -server 192.168.110.175:2181
[zk: 192.168.110.175:2181(CONNECTED) 1] ls /
[zookeeper]
[zk: 192.168.110.175:2181(CONNECTED) 2] ls /zookeeper 
[config, quota]
[zk: 192.168.110.175:2181(CONNECTED) 3] ls /zookeeper/config 
[]
[zk: 192.168.110.175:2181(CONNECTED) 4] ls /zookeeper/
config   quota    
[zk: 192.168.110.175:2181(CONNECTED) 4] ls /zookeeper/quota 
[]
[zk: 192.168.110.175:2181(CONNECTED) 5] create /linux
Created /linux
[zk: 192.168.110.175:2181(CONNECTED) 6] ls /
[linux, zookeeper]
[zk: 192.168.110.175:2181(CONNECTED) 7] create /shy shy514
Created /shy
[zk: 192.168.110.175:2181(CONNECTED) 8] get /shy 
shy514
[zk: 192.168.110.175:2181(CONNECTED) 9] delete /shy 
[zk: 192.168.110.175:2181(CONNECTED) 10] get /shy 
org.apache.zookeeper.KeeperException$NoNodeException: KeeperErrorCode = NoNode for /shy
[zk: 192.168.110.175:2181(CONNECTED) 11] ls
ls [-s] [-w] [-R] path
[zk: 192.168.110.175:2181(CONNECTED) 12] ls /
[linux, zookeeper]
[zk: 192.168.110.175:2181(CONNECTED) 13] delete /linux
[zk: 192.168.110.175:2181(CONNECTED) 14] ls /
[zookeeper]
# zookeeper千万不要删


[zk: 192.168.110.175:2181(CONNECTED) 15] quit
```

```
# 复制会话，这些操作记录在log.1
[root@kafka1 version-2]# ll /opt/zookeeperkeeper/version-2/
总用量 12
-rw-r--r-- 1 root root 67108880 4月  17 18:03 log.1
-rw-r--r-- 1 root root      457 4月  17 17:31 snapshot.0
```





## zookeeper与kafka集群



### 1、配置zookeeper集群

#### 1.修改配置文件

```
[root@kafka1 opt]# grep -Ev "#|^$" /usr/local/zookeeper/conf/zoo.cfg 
tickTime=2000
initLimit=10
syncLimit=5
dataDir=/opt/zookeeperkeeper
clientPort=2181
maxClientCnxns=4096
autopurge.snapRetainCount=128   # /opt/zookeeperkeeper里保存快照的最大数量
autopurge.purgeInterval=1   # 几小时清理一次
# 可以用主机名，因为设置了映射
server.1=192.168.110.175:2888:3888
server.2=192.168.110.176:2888:3888
server.3=192.168.110.177:2888:3888

配置参数解读：

server.A=B:C:D

A是一个数字，表示这个是第几号服务器。myid中的编号就是这个值。zookeeper启动时读取此文件，拿到里面的数据与zoo.cfg里面的配置信息比较从而判断到底是哪个server。
B是这个服务器的地址。
C是这个服务器Follower与集群中的Leader服务器交换信息的端口。
D是万一集群中的leader服务器挂了，需要一个端口来重新进行选举，选举一个新的leader，而这个端口就是用来执行选举时服务器相互通信的端口。
```

```
# 因为之前单点启动过，先关闭
[root@kafka1 opt]# /usr/local/zookeeper/bin/zkServer.sh stop
ZooKeeper JMX enabled by default
Using config: /usr/local/zookeeper/bin/../conf/zoo.cfg
Stopping zookeeper ... STOPPED
```

#### 2.添加集群id

```
[root@kafka1 opt]# cd /usr/local/zookeeper/conf/
[root@kafka1 conf]# scp zoo.cfg 192.168.110.176:`pwd`
root@192.168.110.176's password: 
zoo.cfg                                                                                               100% 1248     2.1MB/s   00:00    
[root@kafka1 conf]# scp zoo.cfg 192.168.110.177:`pwd`
root@192.168.110.177's password: 
zoo.cfg                                                                                               100% 1248     1.4MB/s   00:00    

# 验证是否发到对应机子目录下
```

```
# 在原来的机上清空zoo文件夹
[root@kafka1 conf]# rm -rf //opt/zookeeperkeeper/*
```

```
# 在不同的机子上给予不同的id
[root@kafka1 conf]# echo "1" > /opt/zookeeperkeeper/myid

[root@kafka2 conf]# echo "2" > /opt/zookeeperkeeper/myid

[root@kafka3 conf]# echo "3" > /opt/zookeeperkeeper/myid
```

#### 3.启动zookeeper集群

```
# 三台机子同时操作
/usr/local/zookeeper/bin/zkServer.sh start

# 查询状态
/usr/local/zookeeper/bin/zkServer.sh status

# 两个follower一个leader，我的是kafka2是leader
```

启动时报错如下：

![image-20221012153803366](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210121538462.png)

错误原因

在配置zoo.cfg文件时候，端口号后面还有空格。因为是从其他地方复制的，所以不容易发现

```
server.0=spark1:2888:3888(此处有空格)
server.1=spark2:2888:3888
server.2=spark3:2888:3888
```

解决办法：

删掉端口号后面的空格即可





### 2、配置kafka集群

kafka官网https://kafka.apache.org/

#### 1.下载安装kafka

```
cd /opt && wget https://archive.apache.org/dist/kafka/2.6.0/kafka_2.12-2.6.0.tgz
tar -zxvf kafka_2.12-2.6.0.tgz
mv kafka_2.12-2.6.0 /usr/local/kafka
```

#### 2.修改kafka配置文件

```
[root@kafka1 ~]# cd /usr/local/kafka/config/
[root@kafka1 config]# cp server.properties server.properties.bak
[root@kafka1 config]# grep -Ev "#|^$" server.properties
broker.id=0   # 节点id
listeners=PLAINTEXT://192.168.110.175:9092   # 改成本机ip
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
log.dirs=/opt/kafka   # 储存日志文件
num.partitions=1	#这个参数用于设置新创建的topic有多少个分区，可以根据消费者实际情况配置
num.recovery.threads.per.data.dir=1
offsets.topic.replication.factor=1
transaction.state.log.replication.factor=1
transaction.state.log.min.isr=1
log.retention.hours=168		#这个参数用于配置kafka中消息保存的时间，还支持log.retention.minutes和log.retention.ms配置项
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000
zookeeper.connect=kafka1:2181,kafka2:2181,kafka3:2181   # 所有节点的ip或者映射
zookeeper.connection.timeout.ms=18000
group.initial.rebalance.delay.ms=0
```

复制配置文件

```
[root@kafka1 config]# scp -r server.properties 192.168.110.176:`pwd`
root@192.168.110.176's password: 
server.properties                                     100% 6879     8.6MB/s   00:00    
[root@kafka1 config]# scp -r server.properties 192.168.110.177:`pwd`
root@192.168.110.177's password: 
server.properties                                     100% 6879     7.4MB/s   00:00    

```

```
# 修改其他机子的id和ip
[root@kafka2 ~]# cd /usr/local/kafka/config/
[root@kafka2 config]# vim server.properties
# 修改如下
21 broker.id=1
31 listeners=PLAINTEXT://192.168.110.176:9092


[root@kafka3 config]# cd /usr/local/kafka/config/
[root@kafka3 config]# vim server.properties
# 修改如下
21 broker.id=2
31 listeners=PLAINTEXT://192.168.110.177:9092
```

```
# 全部机子都操作
mkdir -p /opt/kafka
```

#### 2.启动kafka

```
# 全部机子都操作
/usr/local/kafka/bin/kafka-server-start.sh -daemon /usr/local/kafka/config/server.properties
```

```
# 验证
# 主机会有5个端口，两个为自身机子，其他三个是集群
[root@kafka1 config]# netstat -aulntp | grep 9092
tcp6       0      0 192.168.110.175:9092    :::*                    LISTEN      18641/java          
tcp6       0      0 192.168.110.175:40070   192.168.110.175:9092    ESTABLISHED 18641/java          
tcp6       0      0 192.168.110.175:9092    192.168.110.175:40070   ESTABLISHED 18641/java          
tcp6       0      0 192.168.110.175:34324   192.168.110.176:9092    ESTABLISHED 18641/java          
tcp6       0      0 192.168.110.175:39458   192.168.110.177:9092    ESTABLISHED 18641/java


# 从机的
[root@kafka2 config]# netstat -aulntp | grep 9092
tcp6       0      0 192.168.110.176:9092    :::*                    LISTEN      18068/java          
tcp6       0      0 192.168.110.176:9092    192.168.110.175:34324   ESTABLISHED 18068/java
```

### 8.kafka集群基本命令操作

kefka提供了多个命令用于查看、创建、修改、删除topic信息，也可以通过命令测试如何生产消息、消费消息等，这些命令位于kafka安装目录的bin目录下，这里是/usr/local/kafka/bin。登录任意一台kafka集群节点，切换到此目录下，即可进行命令操作。下面列举kafka的一些常用命令的使用方法。



**创建主题与测试**

```
[root@kafka1 config]# /usr/local/kafka/bin/kafka-topics.sh --create --zookeeper kafka1:2181,kafka2:2181,kafka3:2181 --partitions 3 --replication-factor 3 --topic test
Created topic test.



# 验证topic
[root@kafka1 config]# /usr/local/kafka/bin/kafka-topics.sh --describe --zookeeper kafka1:2181,kafka2:2181,kafka3:2181 --topic test
Topic: test	PartitionCount: 3	ReplicationFactor: 3	Configs: 
	Topic: test	Partition: 0	Leader: 2	Replicas: 2,0,1	Isr: 2,0,1
	Topic: test	Partition: 1	Leader: 0	Replicas: 0,1,2	Isr: 0,1,2
	Topic: test	Partition: 2	Leader: 1	Replicas: 1,2,0	Isr: 1,2,0

# 获取所有topic
[root@kafka1 config]# /usr/local/kafka/bin/kafka-topics.sh --list --zookeeper kafka1:2181,kafka2:2181,kafka3:2181
test
```

**生产消息**

```
[root@kafka1 config]# /usr/local/kafka/bin/kafka-console-producer.sh --broker-list kafka1:9092,kafka2:9092,kafka3:9092 --topic test
>this is a test
```

**消费消息**

```
[root@kafka1 config]# /usr/local/kafka/bin/kafka-console-consumer.sh --bootstrap-server kafka1:9092,kafka2:9092,kafka3:9092 --topic test --from-beginning
this is a test

# 在1号即发消息，这里能收到消息
```

**删除消息**

```
[root@kafka1 config]# /usr/local/kafka/bin/kafka-topics.sh --zookeeper kafka1:2181,kafka2:2181,kafka3:2181 --delete --topic test
```

kafka可视化工具

https://www.freesion.com/article/72661130783/



### ELK+Filebeat+Kafka+ZooKeeper构建大数据日志分析平台案例

1 日志分析平台架构图



![img](https://note-1308251438.cos.ap-guangzhou.myqcloud.com/typora/202210120055447.png)

修改filebeat配置文件

```
#---------------------------------------kafka+zookeeper------------------------------------------
filebeat.inputs:
- type: log
  enabled: true
  paths:
   - /var/log/messages
   - /var/log/secure
  fields:
    log_topic: osmessages
name: "192.168.110.173"		#设置filebeat收集的日志中对应主机的ip或主机名
output.kafka:			
  enabled: true
  hosts: ["192.168.110.175:9092","192.168.110.176:9092","192.168.110.177:9092"]
# version: "2.2.2"
  topic: '%{[fields][log_topic]}'	#指定要发送数据给kafka集群的哪个topic，若指定的topic不存在，则会自动创建此topic。
  partition.round_robin:
    reachable_only: true
  worker: 2
  required_acks: 1
  compression: gzip
  max_message_bytes: 10000000
logging.level: debug
```

修改logstash配置文件

```
input {
        kafka {
        bootstrap_servers => "192.168.110.175:9092,192.168.110.176:9092,192.168.110.177:9092"
        topics => ["osmessages"]
        }
}
output {
        elasticsearch {
        hosts => ["192.168.110.170:9200","]
        index => "%{topics}-%{+YYYY-MM-dd}"
        }
}
```

重启服务

```
# 重启filebeat
[root@elk4 opt]# systemctl restart filebeat

# 重启logstash
[root@elk4 conf.d]# logstash -f /etc/logstash/conf.d/kafka.conf 

# 去到es集群的head页面(port:9100)刷新出来，如下图
```

需截图

![chrome_3iEIFYBmBB](https://ye5201314-1312898079.cos.ap-guangzhou.myqcloud.com/202210122221691.png)

 index => "%{topics}-%{+YYYY-MM-dd}"
