为学习记录,顺便记录下遇到的问题


一、Kafka是什么?

Apache Kafka 是一个分布式流处理平台,主要用于构建实时数据管道和流应用。它具有高吞吐量、低延迟、可扩展性和持久性等特点,广泛应用于日志收集、消息系统、事件溯源等场景。

1.Kafka的主要模块

  1. Producer(生产者)

功能:将数据发布到 Kafka 的 Topic 中。

详情:

生产者负责创建消息并发送到 Kafka Broker。

支持同步和异步发送模式。

可以指定分区(Partition)和键(Key)来控制消息的路由。

  1. Consumer(消费者)

功能:从 Kafka 的 Topic 中读取数据。

详情:

消费者以消费者组(Consumer Group)的形式工作,组内的消费者共同消费一个 Topic 的消息。

每个分区只能被消费者组中的一个消费者消费。

支持手动或自动提交偏移量(Offset)。

  1. Broker(代理)

功能:Kafka 集群中的单个节点,负责存储和转发消息。

详情:

每个 Broker 是一个独立的 Kafka 服务器。

Broker 接收生产者的消息并持久化到磁盘。

处理消费者的拉取请求,返回消息。

  1. Topic(主题)

功能:消息的逻辑分类,生产者将消息发送到 Topic,消费者从 Topic 读取消息。

详情:

Topic 可以分为多个分区(Partition),每个分区是一个有序的消息队列。

分区允许 Topic 水平扩展,提高并发性能。

  1. Partition(分区)

功能:Topic 的物理分片,每个分区是一个有序、不可变的消息序列。

详情:

分区在多个 Broker 上分布,实现负载均衡和高可用性。

每个分区可以有多个副本(Replica),其中一个为 Leader,其他为 Follower。

  1. Replica(副本)

功能:分区的备份,用于提高数据的可靠性和可用性。

详情:

Leader 副本处理读写请求,Follower 副本从 Leader 同步数据。

如果 Leader 失效,Kafka 会从 Follower 中选举新的 Leader。

  1. Zookeeper

功能:Kafka 依赖 Zookeeper 进行元数据管理和协调。

详情:

存储 Broker、Topic 和分区的元数据。

管理 Broker 和消费者的状态。

在 Kafka 2.8.0 及以上版本中,Kafka 引入了 KRaft 模式,逐步减少对 Zookeeper 的依赖。

二、Kafka集群部署

系统及软件版本:
注意:kafka_2.13-3.8.1 中 2.13 为 scala 版本

软件版本
UbuntuUbuntu 22.04.5 LTS
Kafkakafka_2.13-3.8.1
Zookeeper3.8.4
OPENJDK17.0.14

1.zookeeper集群搭建

# 安装openjdk
sudo apt-get install openjdk-17-jdk
# 下载kafka和zookeeper
wget https://downloads.apache.org/zookeeper/stable/apache-zookeeper-3.8.4-bin.tar.gz
wget https://downloads.apache.org/kafka/3.8.1/kafka_2.13-3.8.1.tgz
# 解压zookeeper,修改配置文件并启动
tar -zxvf apache-zookeeper-3.8.4-bin.tar.gz
# 创建数据存放目录
mkdir -p /data/zookeeper
# 主要改动的配置为
grep -vE "^#|^$" /usr/local/apache-zookeeper-3.8.4-bin/conf/zoo.cfg
# zookeeper默认时间单位,2000毫秒,2秒
tickTime=2000
# 超时时间 20s
initLimit=10
# 心跳检测时间 
syncLimit=5
# 数据存放目录
dataDir=/data/zookeeper
# 监听端口
clientPort=2181
# 2888集群通信端口,3888集群选举端口
server.1=192.168.1.5:2888:3888
server.2=192.168.1.6:2888:3888
server.3=192.168.1.7:2888:3888
# 打开使用四字命令的限制,*为可以使用所有四字命令,不建议这样做
4lw.commands.whitelist=*

zookeeper集群启动后,主节点为leader,其他节点为follower
在每台机器的 dataDir 目录下创建 myid 文件,并写入对应的服务器 ID:

# 机器1
echo "1" > /data/zookeeper/myid
# 机器2
echo "2" > /data/zookeeper/myid
# 机器3
echo "2" > /data/zookeeper/myid

启动服务:

# 这里可以自己去配置systemd或者开机自启
/usr/local/apache-zookeeper-3.8.4-bin/bin/zkServer.sh start

查看节点状态:

# 安装nc命令
apt-get install netcat
# 查看节点状态
echo stat | nc 192.168.1.5 2181
# 显示信息:
Zookeeper version: 3.8.4-9316c2a7a97e1666d8f4593f34dd6fc36ecc436c, built on 2024-02-12 22:16 UTC
Clients:
 /192.168.1.5:59396[0](queued=0,recved=1,sent=0)

Latency min/avg/max: 0/0.0/0
Received: 4
Sent: 5
Connections: 1
Outstanding: 0
Zxid: 0x300000004
Mode: follower
Node count: 5

至此,zookeeper集群搭建成功

2.kafka集群搭建

kafka分区策略:

轮询: 均匀分布到所有分区。
key哈希: 根据消息的 Key 计算哈希值,将相同 Key 的消息发送到同一个分区。
自定义分区器: 实现自定义的分区逻辑。

kafka 根据分区策略分配到分区后,会写入到 leader 副本中,然后再同步到其他副本,最后确认发送消息成功

# 修改配置文件。vim config/server.properties
broker.id=1
listeners=PLAINTEXT://192.168.1.101:9092
advertised.listeners=PLAINTEXT://192.168.1.101:9092
log.dirs=/tmp/kafka-logs-1
zookeeper.connect=192.168.1.101:2181,192.168.1.102:2181,192.168.1.103:2181
  1. broker.id : 集群kafka没台机器的唯一标识
  2. listeners 和 advertised.listeners : kafka监听的地址,如果advertised.listeners不设置,则会采用listeners的地址,advertised.listeners 外部客户端监听的地址
  3. log.dirs : 为kafka的数据存放目录
    zookeeper.connect: 连接zookeeper的集群,存放kafka中tipic和其他资源的元数据

其他配置

# 保留 7 天的日志
log.retention.hours=168  
# 每个分区最大保留 1GB 数据
log.retention.bytes=1073741824  
# 调整网络缓冲区的大小
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
# 启用消息压缩
message.compress=snappy

启动服务

bin/kafka-server-start.sh config/server.properties

验证
在任意一台机器上运行以下命令,查看 Kafka 集群状态:

bin/kafka-topics.sh --describe --zookeeper 192.168.1.101:2181

总结

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐