kafka为什么这么快

  • 磁盘顺序写入,比随机写入快。
  • 批量发送,生产者可以累积几条消息之后一起发送,减少io次数。
  • 使用零拷贝技术,和netty一样,减少了用户空间和内核空间的多次数据拷贝。
  • 压缩技术,通过压缩减少网络带宽。

kafka模型介绍

Kafka 是发布-订阅模型,通过生产者-消费者模式去处理数据;其中生产端 push,消费端 pull。

数据流的逻辑如下:

1
2
3
4
Producer --push--> Broker --pull--> Consumer
^
|
FetchRequest

这里 FetchRequest 是 Consumer 发给 Broker 的,不是 Broker 发给 Consumer。

Kafka顺序读取

要了解这个,首先要了解broker里面的分区结构:

1
2
3
4
5
6
Broker
└── Topic
└── Partition
├── Leader Replica
├── Follower Replica
└── 状态:ISR、Offset、LEO、HW

Leader Replica 是分区的主副本,处理生产写入和消费读取
Follower Replica 是分区的从副本,从 Leader 同步数据;Leader 挂了,可被选为新 Leader
ISR、Offset、LEO、HW,依次是同步副本集合,偏移量,日志末端偏移,消费者只能消费到 HW 之前的消息

Kafka可以保证同一个分区(partition)是有序的,和FIFO一样。在这个基础上,我们就可以通过业务上的生产者和消费者的控制来实现消息的顺序性读取:

  • 生产者要确保消息进入broker的顺序性,可以通过key来区分不同类型的消息进入不同的分区。
  • 消费者要保证顺序消费,也就是要使用一线程消费一个分区。

这是局部顺序性的一种处理方式。如果要全局的消息顺序性的话,可以:

  • 只使用一个分区,但是会丢失并发部分的性能。
  • 业务层面处理,通过添加不同的key,在收到消息之后进行排序处理,但是会比较复杂。

消息积压处理

  • 增加消费者数量,但是不宜超过分区的数量。因为一个分区同一时间只能被一个消费者消费。
  • 增加分区的数量,提高并发能力,之后需要重新平衡消费者和生产者。

高可用

其实所有的高可用服务都是一样的,服务不能挂,数据不能丢。

  • 集群+多副本架构:集群保证了服务不能挂,多个副本备份数据保证了数据不能丢失。但是这里有个复制的问题,是异步还是同步?实际中还是根据业务的敏感点找一个折中点。
  • 自动选主机制(和redis的哨兵很像),主节点(Master)宕机了,其他从节点选取一个出来成为新的主节点。
  • 动态路由感知机制,在主节点跟换之后,生产者和消费者也需要更新节点信息。这个时候就需要再加一层,去存放这个信息,如KRaft。

消息确认机制(消息的可靠性)

Kafka 是为了海量日志处理诞生的,追求的是极致的吞吐量。

  1. acks机制,它是基于分布式集群的副本同步机制来确保消息的可靠性的。

    • acks=0,发出去就算成功
    • acks=1,有节点存储下来就算成功
    • acks=all/-1,必须等所有同步副本ISR都存下来才算成功
  2. Kafka的消费者确认,叫作提交偏移量(Offset Commit)。

    • 就是按偏移量批量读取数据,要是其中有一条读取失败了,要不重来,要不自己去捕获异常。
    • 其实这种颗粒度太大了,不够精细。

docker 安装Kafka

1
2
3
4
mkdir -p /data/kafka-docker/data
chown -R 1001:1001 /data/kafka-docker
cd /data/kafka-docker

在 /data/kafka-docker 目录下创建 docker-compose.yml:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
version: '3'
services:
kafka:
image: bitnami/kafka:3.6.2
container_name: kafka
restart: always
ports:
- "9092:9092" # 容器内部/本机 PLAINTEXT 通信
- "9094:9094" # 外部客户端访问
environment:
- TZ=Asia/Shanghai
# KRaft 模式配置
- KAFKA_CFG_NODE_ID=0
- KAFKA_CFG_PROCESS_ROLES=controller,broker
- KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=0@kafka:9093
# 监听器配置
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093,EXTERNAL://:9094
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092,EXTERNAL://192.168.0.192:9094
- KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,EXTERNAL:PLAINTEXT
- KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
- KAFKA_CFG_INTER_BROKER_LISTENER_NAME=PLAINTEXT
# 单节点必须设为 1
- KAFKA_CFG_OFFSETS_TOPIC_REPLICATION_FACTOR=1
- KAFKA_CFG_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1
- KAFKA_CFG_TRANSACTION_STATE_LOG_MIN_ISR=1
- KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true
- ALLOW_PLAINTEXT_LISTENER=yes
volumes:
- ./data:/bitnami/kafka

然后执行

1
2
3
4
5
6
7
8
# 拉取
docker pull docker.1ms.run/bitnamilegacy/kafka:3.6.2
# 将拉取的镜像名字和yml文件中配置的对应上,打上标签。
docker tag docker.1ms.run/bitnamilegacy/kafka:3.6.2 bitnami/kafka:3.6.2
# 验证,能看到 bitnami/kafka 和 docker.1ms.run/bitnamilegacy/kafka 指向同一个镜像 ID
docker images | grep kafka
# 启动
docker compose up -d