Kafka 学习笔记
环境搭建
首先下载 Kafka 的 release 文件,在这里我使用的版本是 kafka_2.13-2.4.0。
在这里我们准备部署一个三个节点的集群,首先解压 release 文件,之后进入文件夹 kafka_2.13-2.4.0/config,复制文件 server.properties 得到 server-1.properties 和 server-2.properties,这样一来我们就拥有三份配置文件了,之后我们会使用这三份配置文件来启动三个进程构建 Kafka 集群。分别对三份配置文件中的相关属性进行修改:
server.properties
# broker 的序号
broker.id=0
# 服务监听的端口号
listeners=PLAINTEXT://localhost:9092
# 数据文件的存放目录
log.dirs=~/kafka/logs/kafka-logs-0
# 复制因子
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=3
# ZK 的地址
zookeeper.connect=127.0.0.1:2181
server-1.properties
# broker 的序号
broker.id=1
# 服务监听的端口号
listeners=PLAINTEXT://localhost:9093
# 数据文件的存放目录
log.dirs=~/kafka/logs/kafka-logs-1
# 复制因子
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=3
# ZK 的地址
zookeeper.connect=127.0.0.1:2181
server-2.properties
# broker 的序号
broker.id=2
# 服务监听的端口号
listeners=PLAINTEXT://localhost:9094
# 数据文件的存放目录
log.dirs=~/kafka/logs/kafka-logs-2
# 复制因子
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=3
# ZK 的地址
zookeeper.connect=127.0.0.1:2181
上面之所以要把复制因子修改为 3(默认是 1)是因为 Kafka 默认会使用
__consumer_offsets这个 Topic 保存 Consumer Group 的消费状况,而这个 topic 会在集群启动时自动的创建。因此我们没有办法在集群启动后控制它的 replicas 数量,所以只能在配置文件中进行设置。如果一个__consumer_offsets 的 replicas 只有一个,那么集群中任意一个 broker 挂掉都可能会导致消费出现异常。
修改好配置文件之后就可以根据配置文件使用如下命令启动三个 Kafka 进程了
# 进程 1
JMX_PORT=9192 ./bin/kafka-server-start.sh config/server.properties
# 进程 2
JMX_PORT=9193 ./bin/kafka-server-start.sh config/server-1.properties
# 进程 3
JMX_PORT=9194 ./bin/kafka-server-start.sh config/server-2.properties
启动完毕之后此时 Kafka 集群就已经构建完毕了,我们还可以安装 kafka-manager 这个工具来方便对 Kafka 集群进行管理。kafka-manager 的安装和使用不再赘述,直接参考相关文章即可。
简单使用
首先我们使用 kafka-manager 创建一个 topic,并设置相关属性如下
| properties | value |
|---|---|
| Topic | test |
| Partitions | 3 |
| Replication Factor | 2 |
之后我们构建一个 maven 项目并添加 Kafka Client 的依赖
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.4.0</version>
</dependency>
常量属性
1 | final class Constants { |
Producer
1 | public class TestProducer { |
Consumer
1 | public class TestConsumer { |
在上面我们创建了一个 producer 和一个 consumer,之后我们可以启动一个 producer 和多个 consumer,具体命令如下。
java TestProducer
java TestConsumer
java TestConsumer group_0
java TestConsumer group_0
java TestConsumer group_0
这样一来我们就有了一个 producer 和四个 consumer,其中 consumer 分为两个 group,groupId 分别为 group_0 和 group_test。不同的 group 会消费同一条消息,而同一个 group 内的 consumer 不会消费同一条消息。即:
- 多个 consumer group 可以同时消费一个 topic 中的消息;(topic)
- 同一个 consumer group 中的 consumers 会共享消费的 offset;(queue)
一些原理
Kafka 使用 ZooKeeper 的节点注册来实现 Broker Controller 的选举。
一个 topic 分为多个 partition,每个 partition 又会有多个 replicas,replicas 由 leader 和 follower 组成,其中 replicas 的数量是包含了 leader 的。
producer 负责 push 消息到 partition,consumer 负责从 partition 里 pull 消息。
具体发送到哪一个分区是由 producer 决定的,producer 也可以使用自定义的分区器修改默认的分区配置。
consumer 与 partition 之间的对应关系如下:
- partition > consumer:一个消费者消费多个 partition
- partition = consumer:一个消费者消费一个 partition
- partition < consumer:一个消费者消费一个 partition,多余的 consumer 会处于空闲状态
一个 partition 由多个 segment 组成,每个 segment 文件的名称为其第一条消息的索引值,文件分为日志文件和索引文件。
request.required.acks:
- 1:producer 发送消息到 leader,leader 写入本地日志成功返回(默认)
- 0:~,leader 立即返回(速度快但是有可能丢失消息)
- -1:~,leader 等待所有的 follower 同步完成才返回(速度慢但是强一致)
确定一个 Consumer Group 的 GroupCoordinator 的位置:
- Consumer Group
- GroupId
- abs(GroupId.hashCode) % NumPartition,NumPartition 就是__consumer_offsets 的分区数
- 计算结果表示了__consumer_offsets 的一个 partition
- 找到该 partition 的 leader 所在的 broker
- GroupCoordinator 就在这个 broker 上面
_consumer_offsets 保存了消费者消费的 offset 信息。
后记
Kafka 的内容太多了,而且感觉都比较琐碎,所以记录的东西也都是相当琐碎的。后面会继续学习,遇到相关的内容再进行补充。