3.大数据架构技术 下——搭建kafka
大数据架构部署技术 下——搭建kafka(KRaft模式)集群
文章目录
一、kafka简述
Kafka是由LinkedIn开发的一个分布式流处理平台,作为一种消息中间件,主要用于消费数据、解耦、冗余、异步通信、流量削峰、缓冲等
Kafka支持多种消息格式,如JSON、XML、Avro等,具有高吞吐量、可持久化、可水平扩展、高容错性、支持数据压缩、支持批量和流式处理等特性
在大数据架构中Zookeeper作为集群的分布式协调服务为kafka和Hadoop的高可用集群提供多Master节点选主、节点状态监控、动态配置变更、注册与发现等服务。所以在kafka2.8版本以前与Hadoop都依赖于Zookeeper,从2.8版本开始kafka引入KRaft模式一种自研的共识协议,作为预览/早期访问,但依然建议使用外部的Zookeeper集群,一直到3.3版本KRaft才被标记为生产可用,最终在3.5版本ZooKeeper模式被标记为deprecated而KRaft模式成为默认,现在kafka4.0是完全移除ZooKeeper代码的版本,4.x全线默认且仅支持KRaft模式
新组织模式下Hadoop生态仍然需要Zookeeper为其服务,但kafka已经通过KRaft实现元数据的自治,数据集成关系不变的情况下使 Kafka 运维更轻量、故障隔离更好、分区上限更高
kafka核心概念
Kafka主要由Producer、Broker、Consumer三部分组成
kafka角色组成
topic:是Kafka中消息的逻辑分类,类似于数据库中的"表"。生产者将消息发送到指定的topic,消费者从指定的topic读取消息。一个Kafka集群可以有多个topic,每个topic之间相互隔离
partition:Kafka中的消息分区,每个topic可以包含多个partition,每个partition可以存储多个消息,这是Kafka实现高吞吐和水平扩展的核心机制。消息被写入topic后会被分配到不同的partition中,每个partition是一个有序的、不可变的消息序列,消息在partition内部按写入顺序分配一个递增的偏移量(Offset)
分区的核心作用在于并行处理和水平扩展;多个partition可以分布在不同的broker上,消费者组内的多个消费者可以同时消费不同partition实现并行处理;当单个broker的读写能力不足时,增加partition数量并将它们分散到更多broker上即可提升吞吐量
partition是topic的物理实现单元,topic只是逻辑概念,真正的消息读写发生在partition上
默认参数中分区数量设置为3,分区数量决定了topic最大的消费并行度,一个消费者组内最多只能有与分区数相同数量的消费者同时工作,超出的消费者将处于空闲状态
replica(副本)与replication factor(副本因子):为了防止某个broker宕机导致数据丢失,Kafka为每个partition创建多个replica,分散存储在不同的broker上。默认每个partition有3个副本,
ISR(同步副本集):并非所有副本都有资格参与Leader选举。只有与Leader保持同步的副本才属于ISR。如果某个Follower的数据落后太多或网络异常导致心跳超时,它会被移出ISR
broker:Kafka集群中的节点会在Zookeeper注册并保持相关的元数据更新,负责存储消息、处理客户端请求、转发消息等
producer:生产者,负责将消息发送到Kafka集群
consumer:消费者,负责从Kafka集群中消费消息
consumer group(消费者组):消费者组是Kafka消费端的核心概念。多个消费者可以组成一个消费者组,共同消费一个topic的数据,实现负载均衡。一个partition在同一消费者组内只能被一个消费者消费,保证消息不会被重复处理,反之一个消费者可以消费多个partition,同时消费者组内消费者的数量不能超过partition数量,多余超出的消费者会处于空闲状态。不同消费者组之间互不影响,同一消息可以被不同消费者组各自消费一次。因此分区数决定了单个消费者组的最大消费并行度
Offset(偏移量):类似链表的指针,每个消息在partition内都有一个唯一且递增的编号,称为Offset(偏移量)。消费者通过Offset来记录自己消费到了哪个位置(committed offset,消费者提交的消费位点)。如果消费者宕机后重启,会从上次提交的 Offset 继续消费,避免消息丢失或重复,并且服务端不会记录消费者的消费进度,由消费者自行管理 Offset
KRaft新模式下的角色变动
在KRaft模式下角色发生变动,其中producer、consumer、topic、partition等角色与老模式保持一致
partition:Leader选举由controller节点负责,不再依赖ZooKeeper
broker:在KRaft模式下的注册由Zookeeper被替换为controller
controller:负责管理集群中的所有broker,包括broker的注册、注销、选举等操作,并通过KRaft协议实现集群的共识保证稳定性和一致性
controller 专注元数据,broker 专注消息读写
二、kafka集群节点环境部署
系统规划
| 主机名 | 硬件性能 | 组件服务 |
|---|---|---|
| 192.168.8.20 kafka-ansible | 1C2G 20G | ansible |
| 192.168.8.21 kafka-node1 | 4C4G 20G | Kafka |
| 192.168.8.22 kafka-node2 | 4C4G 20G | Kafka |
| 192.168.8.23 kafka-node3 | 4C4G 20G | Kafka |
案例部署三节点的kafka集群,同样使用上篇中hadoop集群的node1-3节点就不更改主机名了,如果是自行生成新节点可以直接参照上篇中的环境准备工作,重新部署环境、注意主机名的变动,之后再下载kafka,也可以加入ansible方便后续操控
三、部署kafka
安装包版本选择
案例使用KRaft模式的kafka做部署,可以从官网下载kafka4.1.2版本(目前较为稳定的版本),案例部署通过wget获取安装包
https://kafka.apache.org/downloads
kafka配置说明
kafka集群的配置文件只需要修改server.properties,均在/usr/local/kafka/config目录下,并且官方有提供单节点的快速启动文档
https://kafka.apache.org/quickstart/
配置文件说明
config目录下有五个配置文件分别对应各个角色
server.properties:kafka 主配置文件,包含节点标识、角色、监听器、存储等核心配置
producer.properties:对应producer角色
consumer.properties:对应consumer角色
broker.properties:对应broker角色
controller.properties:对应controller角色
配置参数说明
node.id:节点id,相当于节点的"身份证号",集群内不论什么节点的id都不相同
process.roles:节点角色broker、controller,标识节点在集群中扮演的角色,多角色用逗号隔开。小集群各节点既管元数据又处理消息,大集群隔离角色
controller.quorum.voters:定义仲裁组(controller节点)的成员列表,即哪些节点参与Raft共识投票,并且所有节点(包括纯broker节点)都需要配置都需要知道谁能投票,格式为node.id@host:port,多个节点用逗号隔开
注意:controller.quorum.voters和controller.quorum.bootstrap.servers两个参数的区别在于前者解决谁来投票后者解决群众去哪里找"党",前者是必须配置的,后者是可选配置,只有纯brocker节点才需要知道去哪里找controller节点并且无法提供Raft选举所需的完整成员列表;在扩缩容时如果是纯broker节点靠后者就可以自动发现注册集群所以不需要重启集群,但如果是controller节点,所有broker节点成员都需要知道这个新增的节点必须重启集群服务
listeners:监听器,broker监听地址来自客户端的链接,controller监听Raft协议的链接;不绑定具体主机名直接监听所有地址防止hosts解析失败导致服务无法启动
advertised.listeners:节点对外暴露的广播地址
log.dirs:存储目录
num.partitions:设置partition分区数量
default.replication.factor:每个partition分区的副本数量
min.insync.replicas:最少同步副本数
default.replication.factor:副本总数
server.properties配置
修改配置文件/usr/local/kafka/config/server.properties,修改以下配置
# 节点id,集群中每个节点需要唯一,案例中的node1-3节点分别配置为0、1、2
node.id=0
# 关键节点角色,同一个节点既管元数据又处理消息
process.roles=broker,controller
# controller投票节点列表
controller.quorum.voters=0@node1:9093,1@node2:9093,2@node3:9093
# 可选,提供冗余发现能力
# controller.quorum.bootstrap.servers=node1:9093,node2:9093,node3:9093
# 监听器配置,broker监听地址来自客户端的链接,controller监听Raft协议的链接
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
# 对外暴露的广播地址,注意同步其他节点的主机名
advertised.listeners=PLAINTEXT://node1:9092
# 存储目录
log.dirs=/var/kafka/data
自动化部署kafka集群
更改配置文件
可以参考上述配置内容手动修改,案例复用Hadoop集群中的ansible环境编写playbook剧本,使用master节点操控node1-3节点修改配置文件
cd /ansible
# 更新inventory,新增kafka组
echo "
[master]
master ansible_host=192.168.8.20
[kafka]
node1 ansible_host=192.168.8.21
node2 ansible_host=192.168.8.22
node3 ansible_host=192.168.8.23" > /ansible/inventory
server.properties配置
准备server.properties配置文件的模板文件
# 创建server.properties.j2模板文件
cat >> /ansible/template/server.properties.j2 << "EOF"
##############################################
# Kafka KRaft 集群配置
# 节点: {{ inventory_hostname }} (node.id={{ kafka_node_ids[inventory_hostname] }})
##############################################
#------------- 节点标识与角色 ---------------
node.id={{ kafka_node_ids[inventory_hostname] }}
process.roles=broker,controller
#------------- Controller Quorum -----------
controller.quorum.voters=0@node1:9093,1@node2:9093,2@node3:9093
controller.listener.names=CONTROLLER
#------------- 监听器配置 ------------------
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
advertised.listeners=PLAINTEXT://{{ inventory_hostname }}:9092
#------------- 数据存储 --------------------
log.dirs={{ kafka_data_dir }}
#------------- Topic 默认值 ----------------
num.partitions=3
default.replication.factor=3
min.insync.replicas=2
#------------- 日志保留 --------------------
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000
#------------- 网络 ------------------------
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
#------------- Group Coordinator ----------
group.initial.rebalance.delay.ms=3
EOF
service服务文件
准备service服务文件的模板文件
# 创建kafka.service.j2模板文件,用于systemd管理kafka服务
cat >> /ansible/template/kafka.service.j2 << "EOF"
[Unit]
Description=Apache Kafka Broker (KRaft)
Documentation=https://kafka.apache.org/documentation/
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User={{ kafka_user }}
Group={{ kafka_group }}
Environment="JAVA_HOME={{ java_home }}"
Environment="KAFKA_HOME={{ kafka_home }}"
Environment="KAFKA_LOG4J_OPTS=-Dlog4j.configuration=file:{{ kafka_home }}/config/log4j2.properties"
Environment="KAFKA_HEAP_OPTS=-Xms1g -Xmx1g"
ExecStart={{ kafka_home }}/bin/kafka-server-start.sh {{ kafka_home }}/config/server.properties
ExecStop={{ kafka_home }}/bin/kafka-server-stop.sh
Restart=on-failure
RestartSec=10
LimitNOFILE=65536
[Install]
WantedBy=multi-user.target
EOF
编写playbook剧本
# 编写set_kafka_conf.yml主剧本
cat >> /ansible/set_kafka_conf.yml << "EOF"
---
##############################################
# Kafka KRaft 集群自动化部署剧本
##############################################
- name: 部署 Kafka KRaft 集群
hosts: kafka
become: yes
gather_facts: yes
vars:
#------------- 版本与路径 ---------------
kafka_version: "4.1.2"
scala_version: "2.13"
kafka_install_dir: "/usr/local"
kafka_home: "/usr/local/kafka"
kafka_data_dir: "/var/kafka/data"
kafka_user: "root"
kafka_group: "root"
java_home: "/usr/lib/jvm/java-17-openjdk" # 此处按实际java版本更改
#------------- 节点 ID 映射 ------------
kafka_node_ids:
node1: 0
node2: 1
node3: 2
tasks:
# ======================================
# 环境检查
# ======================================
- name: "验证 Java 版本"
shell: "java -version"
register: java_check
changed_when: false
- name: "显示 Java 版本"
debug:
msg: "{{ inventory_hostname }}: {{ java_check.stderr_lines[0] }}"
- name: "验证主机名解析"
command: "getent hosts {{ item }}"
loop: "{{ kafka_node_ids.keys() | list }}"
changed_when: false
- name: "检查端口占用"
shell: "ss -tlnp | grep -E ':9092|:9093' || true"
register: port_check
changed_when: false
- name: "端口冲突警告"
debug:
msg: "警告: {{ inventory_hostname }} 上 9092 或 9093 端口已被占用: {{ port_check.stdout }}"
when: port_check.stdout | length > 0
# ======================================
# 安装 Kafka
# ======================================
- name: "检查 Kafka 是否已安装"
stat:
path: "{{ kafka_home }}/bin/kafka-server-start.sh"
register: kafka_installed
- name: "下载 Kafka 二进制包"
get_url:
url: "https://downloads.apache.org/kafka/{{ kafka_version }}/kafka_{{ scala_version }}-{{ kafka_version }}.tgz"
dest: "/tmp/kafka_{{ scala_version }}-{{ kafka_version }}.tgz"
mode: '0644'
timeout: 300
when: not kafka_installed.stat.exists
tags: [install]
- name: "解压 Kafka"
unarchive:
src: "/tmp/kafka_{{ scala_version }}-{{ kafka_version }}.tgz"
dest: "{{ kafka_install_dir }}"
remote_src: yes
when: not kafka_installed.stat.exists
tags: [install]
- name: "创建软链接"
file:
src: "{{ kafka_install_dir }}/kafka_{{ scala_version }}-{{ kafka_version }}"
dest: "{{ kafka_home }}"
state: link
force: yes
tags: [install]
- name: "清理安装包"
file:
path: "/tmp/kafka_{{ scala_version }}-{{ kafka_version }}.tgz"
state: absent
when: not kafka_installed.stat.exists
# ======================================
# 生成集群 ID(只在 node1 执行一次)
# ======================================
- name: "生成 Cluster UUID"
shell: "{{ kafka_home }}/bin/kafka-storage.sh random-uuid"
register: uuid_result
delegate_to: node1
run_once: true
changed_when: false
- name: "存储 Cluster UUID 到所有节点"
set_fact:
cluster_uuid: "{{ hostvars['node1']['uuid_result'].stdout }}"
- name: "显示 Cluster UUID"
debug:
msg: "集群 UUID: {{ cluster_uuid }}"
run_once: true
# ======================================
# 配置
# ======================================
- name: "创建数据目录"
file:
path: "{{ kafka_data_dir }}"
state: directory
owner: "{{ kafka_user }}"
group: "{{ kafka_group }}"
mode: '0755'
- name: "部署 server.properties"
template:
src: "/ansible/template/server.properties.j2"
dest: "{{ kafka_home }}/config/server.properties"
owner: "{{ kafka_user }}"
group: "{{ kafka_group }}"
mode: '0644'
backup: yes
notify: restart kafka
- name: "部署 systemd 服务文件"
template:
src: "/ansible/template/kafka.service.j2"
dest: "/etc/systemd/system/kafka.service"
mode: '0644'
- name: "重载 systemd"
systemd:
daemon_reload: yes
# ======================================
# 格式化存储(仅首次)
# ======================================
- name: "检查是否已格式化"
stat:
path: "{{ kafka_data_dir }}/meta.properties"
register: formatted
- name: "执行存储格式化"
shell: >
{{ kafka_home }}/bin/kafka-storage.sh format
-t {{ cluster_uuid }}
-c {{ kafka_home }}/config/server.properties
when: not formatted.stat.exists
# ======================================
# 启动服务
# ======================================
- name: "启动并设置开机自启"
systemd:
name: kafka
state: started
enabled: yes
- name: "等待 Kafka 启动"
wait_for:
port: 9092
host: "{{ inventory_hostname }}"
timeout: 60
state: started
- name: "等待 Controller 就绪"
wait_for:
port: 9093
host: "{{ inventory_hostname }}"
timeout: 30
state: started
# ======================================
# 通知处理器
# ======================================
handlers:
- name: restart kafka
systemd:
name: kafka
state: restarted
EOF

启动集群服务
# 执行playbook剧本
ansible-playbook set_kafka_conf.yml
# 验证kafka集群
ansible kafka -m shell -a "jps"
# 创建一个topic
/usr/local/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --partitions 1 --replication-factor 1 --topic msgtest
# 在节点上使用生产者
/usr/local/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic msgtest
# 在另一节点上使用消费者
/usr/local/kafka/bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic msgtest
# 查看topic列表
/usr/local/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
# 查看topic详情(分区分布、Leader 所在节点、ISR 列表)
/usr/local/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--describe --topic msgtest
# 查看消费者组
/usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--list
# 查看消费者组的消费进度
/usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group <group-id>

执行测试数据消费

关闭集群服务
# 关闭服务并停止开机自启
ansible kafka -m shell -a "systemctl disable kafka --now"
四、手动部署(替代方案)
下载安装包
从官网下载kafka4.1.2版本tar包并解压到/usr/local目录下
# 解压到/usr/local目录下
tar -zxf kafka_2.13-4.1.2.tgz -C /usr/local/ --transform='s/kafka_2.13-4.1.2/kafka/'
更改配置文件
在node1-3节点各修改一遍配置文件
# 节点id,集群中每个节点需要唯一,案例中的node1-3节点分别配置为0、1、2
node.id=0
# 关键节点角色,同一个节点既管元数据又处理消息
process.roles=broker,controller
# controller投票节点列表
controller.quorum.voters=0@node1:9093,1@node2:9093,2@node3:9093
# 可选,提供冗余发现能力
# controller.quorum.bootstrap.servers=node1:9093,node2:9093,node3:9093
# 监听器配置,broker监听地址来自客户端的链接,controller监听Raft协议的链接
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
# 对外暴露的广播地址,注意同步其他节点的主机名
advertised.listeners=PLAINTEXT://node1:9092
# 存储目录
log.dirs=/var/kafka/data
生成集群UUID
首先生成集群的UUID,是KRaft模式引入的设计目的是标识集群的唯一ID保存到环境变量中,防止不同集群之间的broker节点互相干扰
# 生成集群UUID,进入kafka的根目录下执行
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
# 验证集群ID
echo $KAFKA_CLUSTER_ID
格式化节点存储目录
使用生成的UUID格式化存储目录
## 在node1-3所有节点都执行一遍
# 创建存储目录
mkdir -p /var/kafka/data
# 格式化节点存储目录
/usr/local/kafka/bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c /usr/local/kafka/config/server.properties
注意:官方文档中的命令包含standalone参数表示单机节点,集群中需要删除参数
启动集群服务
在所有节点执行命令
## 在node1-3所有节点都执行一遍
/usr/local/kafka/bin/kafka-server-start.sh -daemon /usr/local/kafka/config/server.properties
更多推荐
所有评论(0)