返回文章归档
Kafka

ELFK-Kafka【消息队列】

日志消息队列系统架构

一、消息队列概述

1,什么是消息队列

· 传统业务系统的问题

1,用户访问效率问题

假设:我们自研的电商平台,双11有一个爆品被10w人同时点击购买;

  • 订单请求会根据网站开发的程序循序,逐次执行;
  • 每个业务程序执行都有时间消耗;
  • 等到所有的程序执行成功,才会已返回给用户下单成功的消息;
  • 导致问题:
问题 描述
【网站访问速度问题】 用户下单后,需要等待所有系统程序执行完成,才能得到下单结果反馈;
每个系统的执行速度根据当前服务器性能所决定,高并发时可能有较强延迟;
【服务器高负载问题】 在特殊消费日来临时,高并发需要虚增服务器来满足;
而这些新增的服务器,在非消费日又用不到(还必须准备着),导致企业资源浪费;

image-20251211160825251

2,服务器资源利用问题

假设:

  • 每天晚上只有三个小时是网站访问高峰期,需要12台服务器才能满足负载需求;
  • 而每天剩余的21个小时,都是低谷期,只需要2台服务器就能满足负载要求;
  • 那么也就是说,企业80%的资源,是限制的,浪费的;

image-20251211171856214

· 消息队列的作用

消息队列【Message Queue】:简称【MQ】

功能 解释说明
消息并发 异步执行消息,使得执行效率翻倍
削峰填谷 缓存消息,使得服务器资源可以最大化利用

1,消息并发

【消息并发】:

  • 就是将本来多系统的顺序执行逻辑,变为多系统并行执行;
  • 原本需要等待上一步完成才能执行下一步,使用并发消息策略可以多系统异步并行执行;

image-20251211164040563

2,削峰填谷

【削峰填谷】:

  • 就是将大量并发数据,以缓存的方式管理起来;
  • 然后有序地,释放执行,减缓后端系统服务器的压力;
  • 保证后端服务器在高并发时刻,可以有序地,安全的处理请求;而不用增加服务器资源;

image-20251211172920182

2,常见的MQ应用

常见MQ 吞吐量 延迟 消息顺序 适合场景 不适合场景 丢数据
Kafka 100w+/秒 2-100毫秒 分区内严格有序 文本日志
JSON数据
事假记录
监控指标
交易记录
用户行为
大文件:图片/视频)
复杂路由需求
极低延迟<=1毫秒的
几乎不会
RabbitMQ 1-5w/秒 0.1-1毫秒 无保证 微服务通信
任务队列
复杂路由消息
RPC调用
工作流编排
大数据量日志
海量消息重放
长期存储
适中
RocketMQ 10-50w/秒 1-5毫秒 严格有序 金融交易数据
电商订单流程
账单处理
需要事务的消息
只适合java 极低
AxtiveMQ(已过时) <=1w/秒 10-100毫秒 无保证 传统企业应用 极高

3,kafka数据与集群介绍

名词 解释说明
broker kafka集群的每个节点服务器,就是一个broker,broker列表:就是指kafka的集群所有节点
producer 【生产者】:就是往kafka集群【写入数据】的角色
consumer 【消费者】:就是找kafka集群【索要数据】的角色
consumer group 【消费者组】:就是找kafka集群【索要数据】的一组角色
topic 【主题】:就是【生产者】和【消费者】写入和索要数据的逻辑单元;【类似ES中的索引】
parition 【分区】:【topic主题】内部划分的区域;【类似ES中的分片】
offset 【数据消费记录】:【消费者】索要数据的位置,避免重复消费(重复读取);
replica 【副本】:分区内的备份副本;【类似ES中的副本分片】
注意:与ES不同,这里的副本数是指,主分区+副本备份的总和;
# 在kafka集群中,有两个维度的Leader和Follower角色
	- 节点级别:leader主节点,Follower从节点;
	- 数据级别:leader主数据,Follower副本数据;

image-20251213224718704

二、单点部署kafka

1,环境准备

· 机器准备

主机名 ip 配置要求
kafka01 10.0.0.111 2核4G
kafka02 10.0.0.112 2核4G
kafka03 10.0.0.113 2核4G

· 主机名hosts解析

cat >> /etc/hosts <<EOF
10.0.0.111  kafka01
10.0.0.112  kafka02
10.0.0.113  kafka03
EOF

2,下载安装包

· 下载kafka安装包

1,官网

https://kafka.apache.org/

image-20251211173924270

2,进入下载页面

# 点击右上角【DOWNLOAD KAFKA】进入下来页面

image-20251211174016247

3,选择版本下载

版本 上线时间 说明
【3.3.2】 2023年1月 kafka脱离zookeeper的最早稳定版本;
# 选择3.3.2版本的2进制安装包

image-20251211174722670

# 等待下载完成

image-20251211174938505

· 下载jdk11安装包

1,官网

# 进入官网:
https://www.oracle.com/java/technologies/downloads/archive/

image-20251212080320389

2,下载JDK11版本

# 选择X64社区版

image-20251212080445459

# 等待下载完成

image-20251212080605759

3,上传安装包

· 上传kafka安装包

[root@kafka01 ~ ]# rz -E
[root@kafka01 ~ ]# ls -l
-rw-r--r-- 1 root root 106619987 Dec 11 14:45 kafka_2.13-3.3.2.tgz

· 上传JDK安装包

[root@kafka01 ~ ]# rz -E
[root@kafka01 ~ ]# ls -l
-rw-r--r-- 1 root root 168251647 Dec 12 00:06 jdk-11.0.28_linux-x64_bin.tar.gz
-rw-r--r-- 1 root root 106619987 Dec 11 14:45 kafka_2.13-3.3.2.tgz

4,创建安装目录

[root@kafka01 ~ ]# mkdir /tools

5,安装JDK

· 解压JDK到安装目录

[root@kafka01 ~ ]# tar xf jdk-11.0.28_linux-x64_bin.tar.gz -C /tools/
[root@kafka01 ~ ]# ls -l /tools/
drwxr-xr-x 9 root root 4096 Dec 12 00:12 jdk-11.0.28

· 配置环境变量

[root@kafka01 ~ ]# vim /etc/profile
......
export JAVA_HOME="/tools/jdk-11.0.28"
export PATH="$PATH:$JAVA_HOME/bin:$JAVA_HOME/jre/bin"
export CLASSPATH="$CLASSPATH:$JAVA_HOME/lib/:$JAVA_HOME/jre/lib:$JAVA_HOME/lib/tools.jar"

# 变量生效
[root@kafka01 ~ ]# source /etc/profile

# 查看java
[root@kafka01 ~ ]# java -version
java version "11.0.28" 2025-07-15 LTS
Java(TM) SE Runtime Environment 18.9 (build 11.0.28+12-LTS-279)
Java HotSpot(TM) 64-Bit Server VM 18.9 (build 11.0.28+12-LTS-279, mixed mode)

6,单点部署kafka

· 解压安装包

[root@kafka01 ~ ]# tar xf kafka_2.13-3.3.2.tgz -C /tools/
[root@kafka01 ~ ]# ls -l /tools/kafka_2.13-3.3.2/
total 64
drwxr-xr-x 3 root root  4096 Dec 21  2022 bin
drwxr-xr-x 3 root root  4096 Dec 21  2022 config
drwxr-xr-x 2 root root  4096 Dec 11 23:51 libs
-rw-r--r-- 1 root root 14844 Dec 21  2022 LICENSE
drwxr-xr-x 2 root root  4096 Dec 21  2022 licenses
-rw-r--r-- 1 root root 28184 Dec 21  2022 NOTICE
drwxr-xr-x 2 root root  4096 Dec 21  2022 site-docs

· 创建数据和证书存储目录

[root@kafka01 ~ ]# mkdir  /tools/kafka_2.13-3.3.2/{data,certs}
[root@kafka01 ~ ]# ls -l /tools/kafka_2.13-3.3.2/
total 72
drwxr-xr-x 3 root root  4096 Dec 21  2022 bin
drwxr-xr-x 3 root root  4096 Dec 21  2022 config
drwxr-xr-x 2 root root  4096 Dec 12 04:47 certs
drwxr-xr-x 2 root root  4096 Dec 12 04:47 data
drwxr-xr-x 2 root root  4096 Dec 11 23:51 libs
-rw-r--r-- 1 root root 14844 Dec 21  2022 LICENSE
drwxr-xr-x 2 root root  4096 Dec 21  2022 licenses
-rw-r--r-- 1 root root 28184 Dec 21  2022 NOTICE
drwxr-xr-x 2 root root  4096 Dec 21  2022 site-docs

· 生成证书

1,生成CA根证书

# 进入证书目录
[root@kafka01 ~ ]# cd /tools/kafka_2.13-3.3.2/certs

# 生成ca的私钥和公钥
[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# openssl req -x509 -newkey rsa:4096 -sha256 -days 3650 -nodes \
  -subj "/C=CN/ST=Beijing/L=Beijing/O=bakwite/CN=kafka Root CA" \
  -keyout kafka-ca.key -out kafka-ca.crt

2,生成服务端证书

# 生成服务端私钥与请求文件
[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# openssl req -newkey rsa:4096 -nodes -sha256 \
  -subj "/C=CN/ST=Beijing/L=Beijing/O=bakwite/CN=kafka-cluster" \
  -addext "subjectAltName=DNS:kafka01,DNS:kafka02,DNS:kafka03,IP:10.0.0.111,IP:10.0.0.112,IP:10.0.0.113,IP:127.0.0.1" \
  -keyout kafka-server.key -out kafka-server.csr
  
# 生成服务端公钥(集群通讯证书)
[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# openssl x509 -req -in kafka-server.csr -CA kafka-ca.crt -CAkey kafka-ca.key -CAcreateserial \
  -days 365 -out kafka-server.crt \
  -extfile <(printf "subjectAltName=DNS:kafka01,DNS:kafka02,DNS:kafka03,IP:10.0.0.111,IP:10.0.0.112,IP:10.0.0.113,IP:127.0.0.1")

3,服务端证书合并

[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# cat kafka-server.crt kafka-server.key > kafka-server.pem

查看所有证书

[root@kafka01 /tools/kafka_2.13-3.3.2/certs ]# ls -l
-rw-r--r-- 1 kafka kafka 2009 Dec 13 10:24 kafka-ca.crt
-rw------- 1 kafka kafka 3272 Dec 13 10:24 kafka-ca.key
-rw-r--r-- 1 kafka kafka 2069 Dec 13 10:25 kafka-server.crt
-rw-r--r-- 1 kafka kafka 1793 Dec 13 10:25 kafka-server.csr
-rw------- 1 kafka kafka 3272 Dec 13 10:25 kafka-server.key
-rw-r--r-- 1 kafka kafka 5341 Dec 13 10:25 kafka-server.pem

· 编辑配置文件

[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
# ===== 基础配置 =====
# 节点角色
# 控制器节点 (Controller):仅参与集群管理,不处理客户端的生产消费请求。
# 代理节点 (Broker):仅负责存储数据和响应客户端请求,不参与集群管理。
process.roles=broker,controller 
# 节点id(集群唯一)
node.id=1                          
# 总裁选举列表
# # 集群部署
# controller.quorum.voters=1@kafka01:9093,2@kafka02:9093,3@kafka03:9093
# # 单点部署
controller.quorum.voters=1@kafka01:9093

# ===== 网络配置 =====
# 【SASL_SSL://IP:9092】:客户端连接入口,(SASL_SSL表示证书加用户名密码)
listener.security.protocol.map=CONTROLLER:PLAINTEXT,SASL_SSL:SASL_SSL
listeners=SASL_SSL://10.0.0.111:9092,CONTROLLER://10.0.0.111:9093
# 对外公布的地址:客户端实际连接的地址
advertised.listeners=SASL_SSL://10.0.0.111:9092
# broker间通信使用的监听器名称
inter.broker.listener.name=SASL_SSL

# ===== 存储配置 =====
# 数据持久化存储目录,可以写多个,以【,】分割
log.dirs=/tools/kafka_2.13-3.3.2/data
# 默认topic数据保存的时间(毫秒)                                      
log.retention.ms=604800000
# 默认topic数据保存的大小限制(byte)
log.segment.bytes=1073741824
# 消费者offset记录保存时间(单位:分钟[7天])
offsets.retention.minutes=10080

# ===== 集群配置 =====
# 禁止自动创建topic
auto.create.topics.enable=false

# 创建topic时的默认分区数,分区数决定并行度,通常按业务需求设置
# 单点部署设置为1,集群部署设置为【集群数】
num.partitions=1                            
# 默认副本因子:每个分区的副本数,通常设为3,保证高可用
# 单点部署设置为1,集群部署设置为【集群数】
default.replication.factor=1                
# 最小同步副本数,当生产者来确认时,需要这么多副本确认
# 单点部署设置为1,集群部署设置为【集群数-1】
min.insync.replicas=1                       
# 【__consumer_offsets】主题的副本数,存储消费者偏移量,很重要;
# 单点部署设置为1,集群部署设置为【集群数】
offsets.topic.replication.factor=1          
# 事务状态日志的副本数;(如果使用Kafka事务才起作用)
# 单点部署设置为1,集群部署设置为【集群数】
transaction.state.log.replication.factor=1  
# 事务状态日志的最小同步副本数
# 单点部署设置为1,集群部署设置为【集群数-1】 
transaction.state.log.min.isr=1             

# ===== 控制器配置 =====
# Controller监听器名称,必须与listeners中的CONTROLLER对应
controller.listener.names=CONTROLLER

# ===== 性能优化 =====
# 线程数设置(CPU核心数的2倍)
num.network.threads=4
# IO线程数:处理磁盘读写(磁盘数*8)
num.io.threads=8
# socket发送缓冲区大小:1MB
socket.send.buffer.bytes=1024000
# socket接收缓冲区大小:1MB
socket.receive.buffer.bytes=1024000  
# 最大请求大小:100MB
socket.request.max.bytes=104857600   

# ===== 副本配置 =====
# 主分区宕机时,不允许未同步完成同步的副本分区成为leader分区
unclean.leader.election.enable=false
# 是否自动平衡leader分布
auto.leader.rebalance.enable=true
# 副本滞后时间阈值:30秒(超过则认为副本失效)
replica.lag.time.max.ms=30000        

# ===== 安全认证配置 =====
# SSL配置
# 指定PEM格式
ssl.keystore.type=PEM                
ssl.keystore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
# 关闭客户端双向SSL认证(开启双向认证:required)
ssl.client.auth=none  
# 明确设为空,表示不验证主机名
ssl.endpoint.identification.algorithm=  
# 在配置文件中添加信任库配置[自己链接自己时需要信任证书]
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-ca.crt

# SASL配置(用户密码认证)启用的SASL机制:
# # 【PLAIN】(用户名密码明文传输)【SCRAM-SHA-256】(更安全)
sasl.enabled.mechanisms=PLAIN
# broker间通信使用的SASL机制
sasl.mechanism.inter.broker.protocol=PLAIN

· 创建kafka用户名密码文件

写法:[user_用户名="密码"]

说明:

  • 【KafkaServer用户】:是broker节点接收链接请求时的认证用户;
  • 【username/password】:表示节点之间通讯验证的用户;
  • 【user_用户名="密码"】:表示客户端的链接验证用户;
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/kafka_users.conf
KafkaServer {
    org.apache.kafka.common.security.plain.PlainLoginModule required
    username="broker"
    password="123456"
    user_admin="123456"
    user_filebeat="123456"
    user_logstash="123456"
    user_broker="123456";
};

· 编辑启动脚本

将用户文件加入到启动文件中,使其生效;

[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/bin/kafka-server-start.sh
export KAFKA_OPTS="-Djava.security.auth.login.config=/tools/kafka_2.13-3.3.2/config/kafka_users.conf"
......

· 生成集群ID格式化存储目录

注意:单点部署也需要格式化存储目录

1,生成集群ID

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-storage.sh random-uuid

zP0Jy27XTyKo090WVkFqUQ

2,格式化存储目录

必须格式化,否则存储会有问题

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-storage.sh format \
	-t zP0Jy27XTyKo090WVkFqUQ \
	-c /tools/kafka_2.13-3.3.2/config/server.properties

Formatting /tools/kafka_2.13-3.3.2/data with metadata.version 3.3-IV3.

· 前台启动kafka测试(可以不前台启动)

1,前台启动

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-server-start.sh /tools/kafka_2.13-3.3.2/config/server.properties 

2,查看启动端口

[root@kafka01 ~ ]# ss -tnulp          
tcp   LISTEN 0    50    [::ffff:10.0.0.111]:9092  *:*  users:(("java",pid=124457,fd=144))     
tcp   LISTEN 0    50    [::ffff:10.0.0.111]:9093  *:*  users:(("java",pid=124457,fd=124))

· 配置system启动

先【ctrl + c】关闭kafka前台运行

1,创建kafka用户

[root@kafka01 ~ ]# useradd -M -s /sbin/nologin kafka

2,授权kafka安装目录

[root@kafka01 ~ ]# chown -R kafka.kafka /tools/kafka_2.13-3.3.2/

3,编辑system文件

[root@kafka01 ~ ]# vim /lib/systemd/system/kafka.service
[Unit]
Description=Apache Kafka Server
After=network.target

[Service]
Type=simple
User=kafka
Group=kafka
Environment="JAVA_HOME=/tools/jdk-11.0.28"
Environment="KAFKA_HEAP_OPTS=-Xmx1g -Xms1g"
ExecStart=/tools/kafka_2.13-3.3.2/bin/kafka-server-start.sh /tools/kafka_2.13-3.3.2/config/server.properties
ExecStop=/tools/kafka_2.13-3.3.2/bin/kafka-server-stop.sh
Restart=on-failure
RestartSec=10
LimitNOFILE=65536

[Install]
WantedBy=multi-user.target

4,system启动kafka

[root@kafka01 ~ ]# systemctl daemon-reload 
[root@kafka01 ~ ]# systemctl enable --now kafka.service

三、集群部署kafka

1,分发安装目录文件

· 关闭单点kafka

[root@kafka01 ~ ]# systemctl stop kafka.service

· 分发安装目录给其他节点

[root@kafka01 ~ ]# scp -r /tools 10.0.0.112:/
[root@kafka01 ~ ]# scp -r /tools 10.0.0.113:/

· 分发全局变量文件

[root@kafka01 ~ ]# scp /etc/profile 10.0.0.112:/etc/  
[root@kafka01 ~ ]# scp /etc/profile 10.0.0.113:/etc/

# 拷贝后执行生效
[root@kafka02 ~ ]# source /etc/profile
[root@kafka03 ~ ]# source /etc/profile

· 分发system启动文件

[root@kafka01 ~ ]# scp /lib/systemd/system/kafka.service 10.0.0.112:/lib/systemd/system/
[root@kafka01 ~ ]# scp /lib/systemd/system/kafka.service 10.0.0.113:/lib/systemd/system/

2,其他节点创建kafka用户

[root@kafka02 ~ ]# useradd -M -s /sbin/nologin kafka
[root@kafka03 ~ ]# useradd -M -s /sbin/nologin kafka

3,其他节点授权kafka目录

[root@kafka02 ~ ]# chown -R kafka.kafka /tools/kafka_2.13-3.3.2/
[root@kafka03 ~ ]# chown -R kafka.kafka /tools/kafka_2.13-3.3.2/

4,编辑三个节点配置文件

· kafka01节点

[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
# ===== 基础配置 =====
# 节点角色
process.roles=broker,controller 
# 节点id(集群唯一)
node.id=1                          
# 总裁选举列表
# # 集群部署
controller.quorum.voters=1@kafka01:9093,2@kafka02:9093,3@kafka03:9093
# # 单点部署
#controller.quorum.voters=1@kafka01:9093

# ===== 网络配置 =====
# 【SASL_SSL://IP:9092】:客户端连接入口,(SASL_SSL表示证书加用户名密码)
listener.security.protocol.map=CONTROLLER:PLAINTEXT,SASL_SSL:SASL_SSL
listeners=SASL_SSL://10.0.0.111:9092,CONTROLLER://10.0.0.111:9093
# 对外公布的地址:客户端实际连接的地址
advertised.listeners=SASL_SSL://10.0.0.111:9092
# broker间通信使用的监听器名称
inter.broker.listener.name=SASL_SSL

# ===== 存储配置 =====
# 数据持久化存储目录,可以写多个,以【,】分割
log.dirs=/tools/kafka_2.13-3.3.2/data
# 默认topic数据保存的时间(毫秒)
log.retention.ms=604800000
# 默认topic数据保存的最大大小(byte);
log.segment.bytes=1073741824


# ===== 集群配置 =====
# 关闭自动创建topic
auto.create.topics.enable=false

# 创建topic时的默认分区数,分区数决定并行度,通常按业务需求设置
# 单点部署设置为1,集群部署设置为【集群数】
num.partitions=3
# 默认副本因子:每个topic分区的副本数,通常设为3,保证高可用
# 单点部署设置为1,集群部署设置为【集群数】
default.replication.factor=3
# 最小同步副本数,当生产者写入数据时,需要等待起码多少个副本写入完成,才算做写入数据成功?
# 单点部署设置为1,集群部署设置为【集群数-1】
min.insync.replicas=2
# 【__consumer_offsets】主题的副本数,存储消费者偏移量,很重要;
# 单点部署设置为1,集群部署设置为【集群数】
offsets.topic.replication.factor=3
# 事务状态日志的副本数;(如果使用Kafka事务才起作用)
# 单点部署设置为1,集群部署设置为【集群数】
transaction.state.log.replication.factor=3
# 事务状态日志的最小同步副本数
# 单点部署设置为1,集群部署设置为【集群数-1】 
transaction.state.log.min.isr=2
# 消费者offset记录保存时间(分钟)
offsets.retention.minutes=10080


# ===== 控制器配置 =====
# Controller监听器名称,必须与listeners中的CONTROLLER对应
controller.listener.names=CONTROLLER

# ===== 性能优化 =====
# 线程数设置(CPU核心数的2倍)
num.network.threads=4
# IO线程数:处理磁盘读写(磁盘数*8)
num.io.threads=8
# socket发送缓冲区大小:1MB
socket.send.buffer.bytes=1024000
# socket接收缓冲区大小:1MB
socket.receive.buffer.bytes=1024000  
# 最大请求大小:100MB
socket.request.max.bytes=104857600   

# ===== 副本配置 =====
# 主分区宕机时,不允许未同步完成同步的副本分区成为leader分区
unclean.leader.election.enable=false
# 是否自动平衡leader分布
auto.leader.rebalance.enable=true
# 副本滞后时间阈值:30秒(超过则认为副本失效)
replica.lag.time.max.ms=30000        

# ===== 安全认证配置 =====
# SSL配置
# 指定PEM格式
ssl.keystore.type=PEM                
ssl.keystore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
# 关闭客户端双向SSL认证
ssl.client.auth=none  
# 明确设为空,表示不验证主机名
ssl.endpoint.identification.algorithm=  
# 在配置文件中添加信任库配置[自己链接自己时需要信任证书]
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem

# SASL配置(用户密码认证)
# 启用的SASL机制:
# # 【PLAIN】(用户名密码明文传输)【SCRAM-SHA-256】(更安全)
sasl.enabled.mechanisms=PLAIN
# broker间通信使用的SASL机制
sasl.mechanism.inter.broker.protocol=PLAIN

· kafka02节点

[root@kafka02 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
process.roles=broker,controller 
# 节点id(集群唯一):这里每个节点都要不一样
node.id=2                         
# 总裁选举列表
# # 集群部署
controller.quorum.voters=1@kafka01:9093,2@kafka02:9093,3@kafka03:9093
# # 单点部署
#controller.quorum.voters=1@kafka01:9093

# 这里ip地址改为节点的IP
# 【SASL_SSL://IP:9092】:客户端连接入口,(SASL_SSL表示证书加用户名密码)
listener.security.protocol.map=CONTROLLER:PLAINTEXT,SASL_SSL:SASL_SSL
listeners=SASL_SSL://10.0.0.112:9092,CONTROLLER://10.0.0.112:9093
# 对外公布的地址:客户端实际连接的地址
advertised.listeners=SASL_SSL://10.0.0.112:9092
# broker间通信使用的监听器名称
inter.broker.listener.name=SASL_SSL

# 以下:【所有节点都一样】
log.dirs=/tools/kafka_2.13-3.3.2/data
# 默认topic数据保存的时间(毫秒)
log.retention.ms=604800000
log.segment.bytes=1073741824
auto.create.topics.enable=false
num.partitions=3
default.replication.factor=3
min.insync.replicas=2
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
offsets.retention.minutes=10080
controller.listener.names=CONTROLLER
num.network.threads=4
num.io.threads=8
socket.send.buffer.bytes=1024000
socket.receive.buffer.bytes=1024000  
socket.request.max.bytes=104857600   
unclean.leader.election.enable=false
auto.leader.rebalance.enable=true
replica.lag.time.max.ms=30000        
ssl.keystore.type=PEM                
ssl.keystore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
ssl.client.auth=none  
ssl.endpoint.identification.algorithm=  
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
sasl.enabled.mechanisms=PLAIN
sasl.mechanism.inter.broker.protocol=PLAIN

· kafka03节点

[root@kafka03 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
process.roles=broker,controller 
# 节点id(集群唯一):这里每个节点都要不一样
node.id=3                        
# 总裁选举列表
# # 集群部署
controller.quorum.voters=1@kafka01:9093,2@kafka02:9093,3@kafka03:9093
# # 单点部署
#controller.quorum.voters=1@kafka01:9093

# 这里ip地址改为节点的IP
# 【SASL_SSL://IP:9092】:客户端连接入口,(SASL_SSL表示证书加用户名密码)
listener.security.protocol.map=CONTROLLER:PLAINTEXT,SASL_SSL:SASL_SSL
listeners=SASL_SSL://10.0.0.113:9092,CONTROLLER://10.0.0.113:9093
# 对外公布的地址:客户端实际连接的地址
advertised.listeners=SASL_SSL://10.0.0.113:9092
# broker间通信使用的监听器名称
inter.broker.listener.name=SASL_SSL

# 以下:【所有节点都一样】
log.dirs=/tools/kafka_2.13-3.3.2/data
# 默认topic数据保存的时间(毫秒)
log.retention.ms=604800000
log.segment.bytes=1073741824
auto.create.topics.enable=false
num.partitions=3
default.replication.factor=3
min.insync.replicas=2
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
offsets.retention.minutes=10080
controller.listener.names=CONTROLLER
num.network.threads=4
num.io.threads=8
socket.send.buffer.bytes=1024000
socket.receive.buffer.bytes=1024000  
socket.request.max.bytes=104857600   
unclean.leader.election.enable=false
auto.leader.rebalance.enable=true
replica.lag.time.max.ms=30000        
ssl.keystore.type=PEM                
ssl.keystore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
ssl.client.auth=none  
ssl.endpoint.identification.algorithm=  
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem
sasl.enabled.mechanisms=PLAIN
sasl.mechanism.inter.broker.protocol=PLAIN

5,清理历史数据

之前单点部署时,产生的集群元数据;

[root@kafka01 ~ ]# rm -rf /tools/kafka_2.13-3.3.2/data/*
[root@kafka02 ~ ]# rm -rf /tools/kafka_2.13-3.3.2/data/*
[root@kafka03 ~ ]# rm -rf /tools/kafka_2.13-3.3.2/data/*

6,格式化目录

· kafka01生成集群ID

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-storage.sh random-uuid
u7LyLt7yRj-EaS8g2yh32g

· 所有节点初始化数据目录

/tools/kafka_2.13-3.3.2/bin/kafka-storage.sh format \
	-t u7LyLt7yRj-EaS8g2yh32g \
	-c /tools/kafka_2.13-3.3.2/config/server.properties

7,启动集群

· 所有节点system启动

systemctl daemon-reload 
systemctl enable kafka.service
systemctl start kafka.service

· 节点查看

[root@kafka01 ~ ]# ss -tnulp
......
tcp     LISTEN    0         50           [::ffff:10.0.0.111]:9092                   *:*       
tcp     LISTEN    0         50           [::ffff:10.0.0.111]:9093                   *:* 

8,创建客户端用户(文件)

[root@kafka01 ~ ]# vim login-kafka.conf 
security.protocol=SASL_SSL
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required \
    username="filebeat" \
    password="123456";
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem

9,查看验证集群

· 查看集群元数据

显示字段 解释说明
ClusterId: 集群的ID(部署时,格式化目录生成的那个个)
LeaderId: 主节点是哪个ID
LeaderEpoch: 集群的leader更换了多少次了?
HighWatermark: 【高水位线】:已经提交的元数据的日志偏移量
MaxFollowerLag: 【最大落后偏移量】:同步数据落后了多少偏移量
MaxFollowerLagTimeMs: 【最大落后时间】:同步数据落后了多长时间
CurrentVoters: 当前参与选举的broker(集群节点)列表
CurrentObservers: 当前不参与选举的(观察者)broker列表
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-metadata-quorum.sh \
    --bootstrap-server 10.0.0.111:9092 \
    --command-config login-kafka.conf \
    describe --status
ClusterId:              cuR3KFFBR8S9OlmZAduISw
LeaderId:               1
LeaderEpoch:            371
HighWatermark:          1506
MaxFollowerLag:         0
MaxFollowerLagTimeMs:   0
CurrentVoters:          [1,2,3]
CurrentObservers:       []

· 查看集群节点详细信息

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-broker-api-versions.sh \
    --bootstrap-server 10.0.0.111:9092 \
    --command-config login-kafka.conf 
10.0.0.111:9092 (id: 1 rack: null) -> (
	Produce(0): 0 to 9 [usable: 9],
	Fetch(1): 0 to 13 [usable: 13],
	ListOffsets(2): 0 to 7 [usable: 7],
	Metadata(3): 0 to 12 [usable: 12],
	LeaderAndIsr(4): UNSUPPORTED,
	StopReplica(5): UNSUPPORTED,
	UpdateMetadata(6): UNSUPPORTED,
	ControlledShutdown(7): UNSUPPORTED,
	OffsetCommit(8): 0 to 8 [usable: 8],
	OffsetFetch(9): 0 to 8 [usable: 8],
	FindCoordinator(10): 0 to 4 [usable: 4],
	JoinGroup(11): 0 to 9 [usable: 9],
	Heartbeat(12): 0 to 4 [usable: 4],
	LeaveGroup(13): 0 to 5 [usable: 5],
	SyncGroup(14): 0 to 5 [usable: 5],
	DescribeGroups(15): 0 to 5 [usable: 5],
	ListGroups(16): 0 to 4 [usable: 4],
	SaslHandshake(17): 0 to 1 [usable: 1],
	ApiVersions(18): 0 to 3 [usable: 3],
	CreateTopics(19): 0 to 7 [usable: 7],
	DeleteTopics(20): 0 to 6 [usable: 6],
	DeleteRecords(21): 0 to 2 [usable: 2],
	InitProducerId(22): 0 to 4 [usable: 4],
	OffsetForLeaderEpoch(23): 0 to 4 [usable: 4],
	AddPartitionsToTxn(24): 0 to 3 [usable: 3],
	AddOffsetsToTxn(25): 0 to 3 [usable: 3],
	EndTxn(26): 0 to 3 [usable: 3],
	WriteTxnMarkers(27): 0 to 1 [usable: 1],
	TxnOffsetCommit(28): 0 to 3 [usable: 3],
	DescribeAcls(29): 0 to 3 [usable: 3],
	CreateAcls(30): 0 to 3 [usable: 3],
	DeleteAcls(31): 0 to 3 [usable: 3],
	DescribeConfigs(32): 0 to 4 [usable: 4],
	AlterConfigs(33): 0 to 2 [usable: 2],
	AlterReplicaLogDirs(34): 0 to 2 [usable: 2],
	DescribeLogDirs(35): 0 to 4 [usable: 4],
	SaslAuthenticate(36): 0 to 2 [usable: 2],
	CreatePartitions(37): 0 to 3 [usable: 3],
	CreateDelegationToken(38): UNSUPPORTED,
	RenewDelegationToken(39): UNSUPPORTED,
	ExpireDelegationToken(40): UNSUPPORTED,
	DescribeDelegationToken(41): UNSUPPORTED,
	DeleteGroups(42): 0 to 2 [usable: 2],
	ElectLeaders(43): 0 to 2 [usable: 2],
	IncrementalAlterConfigs(44): 0 to 1 [usable: 1],
	AlterPartitionReassignments(45): 0 [usable: 0],
	ListPartitionReassignments(46): 0 [usable: 0],
	OffsetDelete(47): 0 [usable: 0],
	DescribeClientQuotas(48): 0 to 1 [usable: 1],
	AlterClientQuotas(49): 0 to 1 [usable: 1],
	DescribeUserScramCredentials(50): UNSUPPORTED,
	AlterUserScramCredentials(51): UNSUPPORTED,
	DescribeQuorum(55): 0 to 1 [usable: 1],
	AlterPartition(56): UNSUPPORTED,
	UpdateFeatures(57): 0 to 1 [usable: 1],
	DescribeCluster(60): 0 [usable: 0],
	DescribeProducers(61): 0 [usable: 0],
	UnregisterBroker(64): 0 [usable: 0],
	DescribeTransactions(65): 0 [usable: 0],
	ListTransactions(66): 0 [usable: 0],
	AllocateProducerIds(67): UNSUPPORTED
)
10.0.0.112:9092 (id: 2 rack: null) -> (
	Produce(0): 0 to 9 [usable: 9],
	Fetch(1): 0 to 13 [usable: 13],
	ListOffsets(2): 0 to 7 [usable: 7],
	Metadata(3): 0 to 12 [usable: 12],
	LeaderAndIsr(4): UNSUPPORTED,
	StopReplica(5): UNSUPPORTED,
	UpdateMetadata(6): UNSUPPORTED,
	ControlledShutdown(7): UNSUPPORTED,
	OffsetCommit(8): 0 to 8 [usable: 8],
	OffsetFetch(9): 0 to 8 [usable: 8],
	FindCoordinator(10): 0 to 4 [usable: 4],
	JoinGroup(11): 0 to 9 [usable: 9],
	Heartbeat(12): 0 to 4 [usable: 4],
	LeaveGroup(13): 0 to 5 [usable: 5],
	SyncGroup(14): 0 to 5 [usable: 5],
	DescribeGroups(15): 0 to 5 [usable: 5],
	ListGroups(16): 0 to 4 [usable: 4],
	SaslHandshake(17): 0 to 1 [usable: 1],
	ApiVersions(18): 0 to 3 [usable: 3],
	CreateTopics(19): 0 to 7 [usable: 7],
	DeleteTopics(20): 0 to 6 [usable: 6],
	DeleteRecords(21): 0 to 2 [usable: 2],
	InitProducerId(22): 0 to 4 [usable: 4],
	OffsetForLeaderEpoch(23): 0 to 4 [usable: 4],
	AddPartitionsToTxn(24): 0 to 3 [usable: 3],
	AddOffsetsToTxn(25): 0 to 3 [usable: 3],
	EndTxn(26): 0 to 3 [usable: 3],
	WriteTxnMarkers(27): 0 to 1 [usable: 1],
	TxnOffsetCommit(28): 0 to 3 [usable: 3],
	DescribeAcls(29): 0 to 3 [usable: 3],
	CreateAcls(30): 0 to 3 [usable: 3],
	DeleteAcls(31): 0 to 3 [usable: 3],
	DescribeConfigs(32): 0 to 4 [usable: 4],
	AlterConfigs(33): 0 to 2 [usable: 2],
	AlterReplicaLogDirs(34): 0 to 2 [usable: 2],
	DescribeLogDirs(35): 0 to 4 [usable: 4],
	SaslAuthenticate(36): 0 to 2 [usable: 2],
	CreatePartitions(37): 0 to 3 [usable: 3],
	CreateDelegationToken(38): UNSUPPORTED,
	RenewDelegationToken(39): UNSUPPORTED,
	ExpireDelegationToken(40): UNSUPPORTED,
	DescribeDelegationToken(41): UNSUPPORTED,
	DeleteGroups(42): 0 to 2 [usable: 2],
	ElectLeaders(43): 0 to 2 [usable: 2],
	IncrementalAlterConfigs(44): 0 to 1 [usable: 1],
	AlterPartitionReassignments(45): 0 [usable: 0],
	ListPartitionReassignments(46): 0 [usable: 0],
	OffsetDelete(47): 0 [usable: 0],
	DescribeClientQuotas(48): 0 to 1 [usable: 1],
	AlterClientQuotas(49): 0 to 1 [usable: 1],
	DescribeUserScramCredentials(50): UNSUPPORTED,
	AlterUserScramCredentials(51): UNSUPPORTED,
	DescribeQuorum(55): 0 to 1 [usable: 1],
	AlterPartition(56): UNSUPPORTED,
	UpdateFeatures(57): 0 to 1 [usable: 1],
	DescribeCluster(60): 0 [usable: 0],
	DescribeProducers(61): 0 [usable: 0],
	UnregisterBroker(64): 0 [usable: 0],
	DescribeTransactions(65): 0 [usable: 0],
	ListTransactions(66): 0 [usable: 0],
	AllocateProducerIds(67): UNSUPPORTED
)
10.0.0.113:9092 (id: 3 rack: null) -> (
	Produce(0): 0 to 9 [usable: 9],
	Fetch(1): 0 to 13 [usable: 13],
	ListOffsets(2): 0 to 7 [usable: 7],
	Metadata(3): 0 to 12 [usable: 12],
	LeaderAndIsr(4): UNSUPPORTED,
	StopReplica(5): UNSUPPORTED,
	UpdateMetadata(6): UNSUPPORTED,
	ControlledShutdown(7): UNSUPPORTED,
	OffsetCommit(8): 0 to 8 [usable: 8],
	OffsetFetch(9): 0 to 8 [usable: 8],
	FindCoordinator(10): 0 to 4 [usable: 4],
	JoinGroup(11): 0 to 9 [usable: 9],
	Heartbeat(12): 0 to 4 [usable: 4],
	LeaveGroup(13): 0 to 5 [usable: 5],
	SyncGroup(14): 0 to 5 [usable: 5],
	DescribeGroups(15): 0 to 5 [usable: 5],
	ListGroups(16): 0 to 4 [usable: 4],
	SaslHandshake(17): 0 to 1 [usable: 1],
	ApiVersions(18): 0 to 3 [usable: 3],
	CreateTopics(19): 0 to 7 [usable: 7],
	DeleteTopics(20): 0 to 6 [usable: 6],
	DeleteRecords(21): 0 to 2 [usable: 2],
	InitProducerId(22): 0 to 4 [usable: 4],
	OffsetForLeaderEpoch(23): 0 to 4 [usable: 4],
	AddPartitionsToTxn(24): 0 to 3 [usable: 3],
	AddOffsetsToTxn(25): 0 to 3 [usable: 3],
	EndTxn(26): 0 to 3 [usable: 3],
	WriteTxnMarkers(27): 0 to 1 [usable: 1],
	TxnOffsetCommit(28): 0 to 3 [usable: 3],
	DescribeAcls(29): 0 to 3 [usable: 3],
	CreateAcls(30): 0 to 3 [usable: 3],
	DeleteAcls(31): 0 to 3 [usable: 3],
	DescribeConfigs(32): 0 to 4 [usable: 4],
	AlterConfigs(33): 0 to 2 [usable: 2],
	AlterReplicaLogDirs(34): 0 to 2 [usable: 2],
	DescribeLogDirs(35): 0 to 4 [usable: 4],
	SaslAuthenticate(36): 0 to 2 [usable: 2],
	CreatePartitions(37): 0 to 3 [usable: 3],
	CreateDelegationToken(38): UNSUPPORTED,
	RenewDelegationToken(39): UNSUPPORTED,
	ExpireDelegationToken(40): UNSUPPORTED,
	DescribeDelegationToken(41): UNSUPPORTED,
	DeleteGroups(42): 0 to 2 [usable: 2],
	ElectLeaders(43): 0 to 2 [usable: 2],
	IncrementalAlterConfigs(44): 0 to 1 [usable: 1],
	AlterPartitionReassignments(45): 0 [usable: 0],
	ListPartitionReassignments(46): 0 [usable: 0],
	OffsetDelete(47): 0 [usable: 0],
	DescribeClientQuotas(48): 0 to 1 [usable: 1],
	AlterClientQuotas(49): 0 to 1 [usable: 1],
	DescribeUserScramCredentials(50): UNSUPPORTED,
	AlterUserScramCredentials(51): UNSUPPORTED,
	DescribeQuorum(55): 0 to 1 [usable: 1],
	AlterPartition(56): UNSUPPORTED,
	UpdateFeatures(57): 0 to 1 [usable: 1],
	DescribeCluster(60): 0 [usable: 0],
	DescribeProducers(61): 0 [usable: 0],
	UnregisterBroker(64): 0 [usable: 0],
	DescribeTransactions(65): 0 [usable: 0],
	ListTransactions(66): 0 [usable: 0],
	AllocateProducerIds(67): UNSUPPORTED
)

· 查看集群健康状态

显示字段 含义说明
NodeId 集群的ID
LogEndOffset 元数据日志的结束偏移量(现在有2039条日志数据)
Lag 副本落后于Leader的偏移量(表示数据完全同步)
LastFetchTimestamp 最后一次拉取数据的时间戳
LastCaughtUpTimestamp 最后一次追上Leader的时间戳
Status 节点的角色
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-metadata-quorum.sh \
    --bootstrap-server 10.0.0.111:9092 \
    --command-config login-kafka.conf \
    describe --replication
NodeId	LogEndOffset	Lag	LastFetchTimestamp	LastCaughtUpTimestamp	Status  	
1     	2039        	0  	1765622013103     	1765622013103        	Leader  	
2     	2039        	0  	1765622012819     	1765622012819        	Follower	
3     	2039        	0  	1765622012819     	1765622012819        	Follower

四、kafka数据流管理

1,topic管理

先查看topic列表:当前什么都没有

[root@kafka01 ~ ]# cd /tools/kafka_2.13-3.3.2/

[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
    --bootstrap-server 10.0.0.111:9092 \
    --command-config login-kafka.conf \
    --list

· 创建topic

参数 说明
--create 表示对topic的动作为:创建
--bootstrap-server 指定创建的节点
--command-config 指定使用的用户身份文件
--partitions 3 指定分区数(不指定会使用配置文件中的默认参数值)
--replication-factor 3 指定副本数(不指定会使用配置文件中的默认参数值)
注意:kafka的分区副本与ES的概念不同,是本分数加上原数据的总数为副本数
[root@kafka01 ~ ]# cd /tools/kafka_2.13-3.3.2/              #login-kafka.conf必须在当前目录

[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
  --create \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf \
  --topic bakwite-topic01 \
  --partitions 3 \
  --replication-factor 3

Created topic bakwite-topic01

创建后,查看topic列表

[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
    --bootstrap-server 10.0.0.111:9092 \
    --command-config login-kafka.conf \
    --list

bakwite-topic01

· 查看topic

1,查看topic列表

参数 说明
--list 表是对topic的操作为:显示列表
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
    --bootstrap-server 10.0.0.111:9092 \
    --command-config login-kafka.conf \
    --list
    
bakwite-topic01

2,查看topic详情

参数 说明
--describe 表示要查看topic的详情
--topic 指定要查看的topic分区;如果不写,会查看所有
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
  --describe \
  --topic bakwite-topic01 \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf 

Topic: bakwite-topic01	TopicId: FGAY7jg1SPqLBXXMVpkuvw	PartitionCount: 3	ReplicationFactor: 3	Configs: min.insync.replicas=2,unclean.leader.election.enable=false
	Topic: bakwite-topic01	Partition: 0	Leader: 1	Replicas: 1,2,3	Isr: 1,2,3
	Topic: bakwite-topic01	Partition: 1	Leader: 2	Replicas: 2,3,1	Isr: 2,3,1
	Topic: bakwite-topic01	Partition: 2	Leader: 3	Replicas: 3,1,2	Isr: 3,1,2

显示结果说明:

显示结果 二级字段 解释说明
Topic: bakwite-topic01 表示topic名称
TopicId: FGAY7...... 表示topic的ID
PartitionCount: 3 表示该topic的分区数
ReplicationFactor: 3 表示该topic的副本数
Configs: min.insync.replicas=2, 最小同步副本数
生产者写入数据时,必须两个副本写入完成,才算成功
unclean.leader.election.enable=false 当leader宕机时,是否允许从非同步副本中选举新leader
leader 当前分区的leader所在broker节点;
Isr 当前处于同步状态的副本;因该是1.2.3都有才算完整;
说明,备份没有问题;

3,各节点查看目录

[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0  #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:43 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2

[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:43 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1  #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:36 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2

[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:36 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2  # 主

· 修改topic

参数 说明
--alter 修改topic主题

注意:

  • 只能修改分区数,副本数一经创建,无法修改;
  • 修改分区,也只能增加,不能减少;
# 将分区数从:3,修改为:5;
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
  --alter \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf \
  --topic bakwite-topic01 \
  --partitions 5

修改后查看topic详情

[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
  --describe \
  --topic bakwite-topic01 \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf
Topic: bakwite-topic01	TopicId: PN4Ew3jXRWmloK-sRUA80Q	PartitionCount: 5	ReplicationFactor: 3	Configs: min.insync.replicas=2,unclean.leader.election.enable=false
	Topic: bakwite-topic01	Partition: 0	Leader: 2	Replicas: 2,3,1	Isr: 2,3,1
	Topic: bakwite-topic01	Partition: 1	Leader: 3	Replicas: 3,1,2	Isr: 3,1,2
	Topic: bakwite-topic01	Partition: 2	Leader: 1	Replicas: 1,2,3	Isr: 1,2,3
	Topic: bakwite-topic01	Partition: 3	Leader: 2	Replicas: 2,3,1	Isr: 2,3,1
	Topic: bakwite-topic01	Partition: 4	Leader: 3	Replicas: 3,1,2	Isr: 3,1,2

查看各个节点目录验证

[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:43 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2   #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-3
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-4

[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0   #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:43 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:36 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-3   #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-4

[root@kafka01 /tools/kafka_2.13-3.3.2]# ls -ld /tools/kafka_2.13-3.3.2/data/bakwite-topic01-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:33 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-1   #主
drwxr-xr-x 2 kafka kafka 4096 Dec 14 05:36 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-2
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-3
drwxr-xr-x 2 kafka kafka 4096 Dec 14 11:23 /tools/kafka_2.13-3.3.2/data/bakwite-topic01-4   #主

· 删除topic

参数 说明
--delete 表示对topic进行删除操作
[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
  --delete \
  --topic bakwite-topic01 \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf 

2,生产者写入数据

· 编辑生产者配置文件

[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/producer.properties
# 链接kafaka的入口
bootstrap.servers=kafka01:9092,kafka02:9092,kafka03:9092
# 消息压缩类型【none】表示不压缩
# 可以设置为:gzip,snappy、lz4、zstd,但是会增加CPU开销
# 推荐【snappy】或【lz4】
compression.type=snappy

# 消息从发送到收到确认的总超时时间(2分钟)
delivery.timeout.ms=120000
# 单次请求(发送或获取元数据)的超时时间(30秒)
request.timeout.ms=30000
# 启用幂等生产者;
enable.idempotence=true
# # 开启幂等性,默认包含下面两项参数;
# # 消息确认机制,保证数据不丢失的关键配置
# # all表示所有ISR副本确认,才算写入成功,否则写入失败;
# acks=all
# # 当确定发送失败后的重试次数【下面这个值表示无限重试】
# retries=2147483647

# 批次和缓冲区
# 最多等待5ms打包
linger.ms=5
# 16KB批次大小
batch.size=16384   
# 32MB发送缓冲区
buffer.memory=33554432          


# 单个连接上允许的未确认请求最大数量
max.in.flight.requests.per.connection=5

# 指定key和value的序列化器
key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer

# 链接集群用户
security.protocol=SASL_SSL
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="123456";
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem

· 创建topic

[root@kafka01 ~ ]# cd /tools/kafka_2.13-3.3.2/   

[root@kafka01 /tools/kafka_2.13-3.3.2]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
  --create \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf \
  --topic bakwite-topic01 \
  --partitions 3 \
  --replication-factor 3

Created topic bakwite-topic01.

· 生产者写入数据

1,交互式写入数据

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
  --topic bakwite-topic01
>hello,world!
>bakwite!
>go on kfka!
>

2,非交互式写入数据

# 单行数据写入
[root@kafka01 ~ ]# echo "hello,bakwite!" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
  --topic bakwite-topic01

# 多行写入
[root@kafka01 ~ ]# cat <<EOF | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
  --topic bakwite-topic01
hello,wa
kafka new message
gogo!
EOF

# 从文件中读取写入
[root@kafka01 ~ ]# cat 1.txt
hei,kafka

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --producer.config login-kafka.conf \
  --topic bakwite-topic01 < 1.txt

3,消费者读取数据

· 编辑消费者配置文件

1,编辑消费者配置文件

[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/consumer.properties 
# 连接
bootstrap.servers=10.0.0.111:9092,10.0.0.112:9092,10.0.0.113:9092
# 消费者组ID
group.id=wa-group

# 可靠性
# 是否自动提交offset(消费位移)
enable.auto.commit=true
# 当没有初始offset或offset失效时,从哪里开始消费
# latest:表示最新的消息开始消费
# earliest:表示最早的消息开始消费(重头开始)
# none:没有offset时抛出异常
auto.offset.reset=earliest

# 性能
# 消费者从broker拉取数据的最小字节数
# 设置大了可以减少请求,但是会增加延迟;
fetch.min.bytes=1
# 等待fetch.min.bytes数据的最长时间(毫秒)
fetch.max.wait.ms=500
# 每个partition一次拉取的最大字节数(默认1MB)
max.partition.fetch.bytes=1048576
# 消费者一次从kafka拉取消息的最大数量;
max.poll.records=1000

# 容错
# 消费者会话超时时间(30秒)
session.timeout.ms=30000
# 消费者带一次拉取完毕后,超过5分钟,kafka会认为这个消费者挂了,触发冲平衡,重新分配消费者;
max.poll.interval.ms=300000
# 消费者定期发送心跳证明自己存活
heartbeat.interval.ms=3000

# 反序列化
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

# 链接集群用户
security.protocol=SASL_SSL
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="123456";
ssl.truststore.type=PEM
ssl.truststore.location=/tools/kafka_2.13-3.3.2/certs/kafka-server.pem

2,分发其他节点

[root@kafka01 ~ ]# scp /tools/kafka_2.13-3.3.2/config/consumer.properties 10.0.0.112:/tools/kafka_2.13-3.3.2/config/consumer.properties
[root@kafka01 ~ ]# scp /tools/kafka_2.13-3.3.2/config/consumer.properties 10.0.0.113:/tools/kafka_2.13-3.3.2/config/consumer.properties

· 读取所有数据

参数 说明
--from-beginning 从头开始读取数据(不写这个参数,默认值读新的写入数据)
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
  --topic bakwite-topic01 \
  --from-beginning 

hello,world!
bakwite!
go on kfka!
hello,bakwite!
hello,wa
kafka new message
gogo!
hei,kafka

· 读取前5条数据

参数 说明
--max-messages 5 读取多少条数据
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
  --topic bakwite-topic01 \
  --from-beginning \
  --max-messages 5
hello,world!
bakwite!
go on kfka!
hello,bakwite!
hello,wa
Processed a total of 5 messages

· 消费者组读取数据

1,生产者写入数据原理

【生产者】写入数据:

  • 只写入topic主题的主分区中
  • 写入是否成功,取决于topic的配置;
    • 副本分区至少成功复制完多少个,算作写入成功?配置文件参数:【min.insync.replicas=2】;
    • 也就是说,加上主分区,共计3个副本数据,成功写入2个,才算作数据写入完成;

image-20251214153514092

2,消费者读取数据原理

【消费者】读取数据:

  • 消费者读取数据,只读取【leader主分区副本数据】

    • kafka为了保证数据的强一致性,就是leader分区与follower分区数据相同;
    • 要求读取写入都必须是leader的分区;
    • 因为,在读取时,副本follower很可能还没复制完,导致数据不一致,读取不全;
  • 消费者读数据,会在kafka的data目录下创建读取记录的目录;

    • 记录着读取数据的offset偏移量;
    • 这些offset偏移量目录,可以让集群监控到各个节点的消费及情况;
    • 即便集群宕机重启,消费者根据这些offset记录也能知道从哪里继续消费读数据;
    • 避免数据重复读取;
[root@kafka01 ~ ]# ls -ld /tools/kafka_2.13-3.3.2/data/__consumer_offsets-*
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-0
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-1
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-10
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-11
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-12
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-13
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-14
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-15
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-16
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-17
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-18
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-19
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-2
drwxr-xr-x 2 kafka kafka 4096 Dec 14 06:42 /tools/kafka_2.13-3.3.2/data/__consumer_offsets-20
......

image-20251214155131136

3,kafka消费者读数据瓶颈

消费者读数据瓶颈问题 说明
broker单点故障 当消费者在一个broker节点正在读数据,突然这台机器挂了,消费者也就挂了,整个消费停止;
高并发数据积压 当topic有大量数据积压时,一个消费者消费性能有限,读取缓慢;

4,消费者组

为了解决【消费者单点故障】与【高并发数据积压】的问题:

  • kafak采用了比较特别的手段:【消费者组】;
  • 就是创建消费者组:启动多个消费者,绑定这个消费者组;
    • kafka会将topic内的所有分区,自动分配给组内消费者成员进行消费;
    • 如果有新的消费者加入消费者组,会自动触发重平衡,再次平均分配分区;
    • 组内【消费者】与【分区】是【1对多】的关系;
      • 一个分区最多有一个消费者;
      • 而一个消费者最多可以消费多个分区;
    • 当触发重平衡后,组内成员会自动接管上一任消费者的offset,保正数据不会重复消费;
  • 也就是说:
    • 当我启动三个消费者进程分别在三台broker主机上;
    • 就可以解决单点故障问题,和高并发数据积压问题;

image-20251214195400762

将用户文件分发给其他节点

[root@kafka01 ~ ]# scp login-kafka.conf 10.0.0.112:~/
[root@kafka01 ~ ]# scp login-kafka.conf 10.0.0.113:~/

创建topic

# 创建topic
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
  --create \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf \
  --topic bakwite-topic01 \
  --partitions 3 \
  --replication-factor 3
Created topic bakwite-topic01

三个节点分别启动消费者,并【指定相同消费者组】

  • 注意:消费者组会自动创建
  • 当消费者组未被使用时,超过7天,自动删除,无需人为干预;
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
  --topic bakwite-topic01 

[root@kafka02 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
  --topic bakwite-topic01 

[root@kafka03 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties \
  --topic bakwite-topic01 

查看消费者组列表

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-consumer-groups.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf \
  --all-groups \
  --list 

wa-group

查看消费者组详细状态

image-20251215103112775

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-consumer-groups.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf \
  --group wa-group \
  --describe

生产者,多次写入数据

  • 总结:
  • 生产者每次写入数据,组内只有一个消费者会读取,不会重复消费;
  • 即便一台机器挂了,也不影响消费数据;
[root@kafka01 ~ ]# echo "这是一条测试消息" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
  --topic bakwite-topic01
  
[root@kafka01 ~ ]# echo "这是一条测试消息" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
  --topic bakwite-topic01

[root@kafka01 ~ ]# echo "这是一条测试消息" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
  --topic bakwite-topic01
  
[root@kafka01 ~ ]# echo "这是一条测试消息" | /tools/kafka_2.13-3.3.2/bin/kafka-console-producer.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --producer.config /tools/kafka_2.13-3.3.2/config/producer.properties \
  --topic bakwite-topic01

image-20251214201652281

4,kafka数据相关

· 什么情况下数据会丢

当生产者向leader分区写入数据后:

  • 副本分区开始同步leader主分区的数据;
  • 当一个broker节点尚未同步完成,leader节点突然宕机,同时kafka将这个主节点选举为leader分区;
  • 此时,所有其他节点开始与这个broker节点同步数据(因为它是leader了),导致数据不完整;

image-20251215104836155

【ISR】:主题分区同步情况;

  • 当我们查看主题详情,看到最后最后一列叫做【ISR】
  • 【In-Sync Replicas】 = 同步副本
  • ISR这一列,就可以看到,当前副本分区的同步情况;
# 查看主题topic详情
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-topics.sh \
  --describe \
  --topic bakwite-topic01 \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf 
Topic: bakwite-topic01	TopicId: ZAz2sFLxTnaXr6GtRtErGw	PartitionCount: 3	ReplicationFactor: 3	Configs: min.insync.replicas=2,segment.bytes=1073741824,retention.ms=604800000,unclean.leader.election.enable=false
	Topic: bakwite-topic01	Partition: 0	Leader: 3   Replicas: 3,1,2	Isr: 3,2,1
	Topic: bakwite-topic01	Partition: 1	Leader: 1   Replicas: 1,2,3	Isr: 3,2,1
	Topic: bakwite-topic01	Partition: 2	Leader: 2   Replicas: 2,3,1	Isr: 3,2,1

· kafka如何保证数据不丢失

1,生产者配置文件

目标让生产者写入数据时:

  • 启动幂等性:【enable.idempotence=true】

  • 必须所有副本确认同步完成后,才算写入成功;

  • 否则,判定失败,生产者重新写入数据;

[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/producer.properties
......
# 消息从发送到收到确认的总超时时间(2分钟)
delivery.timeout.ms=120000
# 单次请求(发送或获取元数据)的超时时间(30秒)
request.timeout.ms=30000
# 启用幂等生产者;
enable.idempotence=true
# # 开启幂等性,默认包含下面两项参数;
# # 消息确认机制,保证数据不丢失的关键配置
# # all表示所有ISR副本确认,才算写入成功,否则写入失败;
# acks=all
# # 当确定发送失败后的重试次数【下面这个值表示无限重试】
# retries=2147483647
......

2,kafka配置文件

目标让kafka保证主分区宕机时:

  • 必须同步完成的副本分区才可以被选举为主分区;
    • 只有同步完成的副本分区,才有资格成为leader分区;
    • 【unclean.leader.election.enable=false】
  • 保证有同步完的备选分区
    • 生产者写入时,必须起码有2个副本分区同步完成,才算做写入完成;
    • 【min.insync.replicas=2】
[root@kafka01 ~ ]# vim /tools/kafka_2.13-3.3.2/config/server.properties
...
# 必须两个副本同步完成,才算写入成功,否则,写入失败(回复生产者写入失败);
min.insync.replicas=2
# topic副本设置为3;保证数据安全;
default.replication.factor=3
# 主分区宕机时,不允许未同步完成同步的副本分区成为leader分区
unclean.leader.election.enable=false
# 是否自动平衡leader分布(消费者重平衡)
auto.leader.rebalance.enable=true
# 副本滞后时间阈值:30秒(超过则认为副本同步失效)
replica.lag.time.max.ms=30000 

· kafka的kraft选举机制

1,Contrller与Broker角色

Controller节点核心功能:

  • 主题与分区管理:负责增删改查“主题topic”、分区扩容工作;
  • 分区Leader选举:当broker宕机时,Controller会为受到影响的分区快速选举出新的Leader副本,确保服务可用;
  • 集群节点扩缩容:Controller能够监听Broker的加入和退出,并执行响应的“善后”工作;
  • 元数据管理:Controller维护着最全的集群元数据,并负责将这些信息同步给集群内所有的Broker;
kafka节点类型 名称 作用
broker 集群数据存储节点 kafka的集群节点
Controller 控制节点 kafka集群的管理功能节点

2,Zookeeper与Kraft集群模式区别

Zookeeper模式下集群模式:

  • Broker在Zookeeper集群上创建临时节点来竞争Controller;
  • 由Zookeeper的ZAB协议保证只有一个可以选举成功;
  • 元数据(ISR、分区、主题)存储位置:
    • 在Zookeeper中,Contrller需要频繁读取Zookeeper获取元数据;
    • 再同步给其他Broker节点;

Kraft模式下集群模式:

  • Controller不再是1个节点,而是一组(3、5个);
  • 由Raft协议投片选举出Contrller,其他参与选举的节点作为备用,
  • 主节点宕机,备用节点快速接管;
  • 元数据(ISR、分区、主题)存储位置:
    • 存储在kafka的一个元数据主题中:__cluster_metadata
    • 由raft协议复制到所有Contrller节点;
    • 本质上就是一个内置的、高可用的元数据日志(如MySQL的binlog),不再需要Zookeeper;
集群模式 对应
Kraft kafka3.3.2之后版本
Zookeeper kafka3.3.2之前版本,kafka4.0移除Zookeeper;

5,kfka环境要求

场景 broker kraft-Controller 副本 副本确认 吞吐量 CPU/内存 存储容量
小型 3 3(与broker合并) 3 2 100M/s 4核/8G <1T
中小型 5-10 3(独立) 3 2 100-500M/s 8核/64G 1-20T
大型 10-20 3/5(独立) 3 2 500-2000M/s 32核/64G 20TB-1PB
超大 2-50 5(独立) 3 2 >2000M/s 32核/64G PB级别

五、ELFK架构实战

1,filebeat配置写入kafka

· 编辑filebeat配置文件

[root@kibana ~ ]# vim /tools/filebeat-8.18.8-linux-x86_64/filebeat.yml
filebeat.inputs:
- type: filestream
  id: bakwite_01
  enabled: true
  paths:
    - /test/1.txt

# 输出到kafka
output.kafka:
  hosts:
    - 10.0.0.111:9092
    - 10.0.0.112:9092
    - 10.0.0.113:9092
  # 写入到topic主题
  topic: bakwite-topic01
  #【幂等性】
  # 保证副本同步完成才算写入成功,否则重新写入;
  required_acks: -1
  # 只要失败就无限次重新写入,保证数据不丢失
  max_retries: -1
  # 写入超过30s算失败,重写
  timeout: 30s
  # 压缩方式
  compression: gzip
  # 轮询kafka节点写入
  partition.round_robin:
    # true表示:只写入存活的broker节点;
    reachable_only: true

  # ===== SSL =====
  ssl.enabled: true
  ssl.certificate_authorities:
    - /tools/filebeat-8.18.8-linux-x86_64/certs/kafka-ca.crt

  # ===== SASL =====
  sasl.mechanism: PLAIN
  username: filebeat
  password: "123456"

· 复制kafka的服务证书

[root@kibana ~ ]# scp 10.0.0.111:/tools/kafka_2.13-3.3.2/certs/kafka-ca.crt /tools/filebeat-8.18.8-linux-x86_64/certs/

· 重启filebeat

[root@kibana ~ ]# systemctl restart filebeat.service 

· kafka启动消费者

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-console-consumer.sh  \
	--bootstrap-server 10.0.0.111:9092   \
	--consumer.config /tools/kafka_2.13-3.3.2/config/consumer.properties   \
	--topic bakwite-topic01

· 输入文件写入内容

[root@kibana ~ ]# echo "hello filebeat" >> /test/1.txt
[root@kibana ~ ]# echo "hello filebeat" >> /test/1.txt
[root@kibana ~ ]# echo "hello filebeat" >> /test/1.txt

· 查看消费者终端

{
	"@timestamp":"2025-12-15T05:25:32.977Z",
	"@metadata":{"beat":"filebeat","type":"_doc","version":"8.18.8"},
	"host":{"name":"kibana"},
	"agent":{
		"ephemeral_id":"7c76994b-bc88-4f30-9e11-322dc4bfd151",
		"id":"2f97e00a-7e4c-45d4-8bd1-a65742eaf4f8",
		"name":"kibana",
		"type":"filebeat",
		"version":"8.18.8"
	},
	"log":{
		"file":{
			"path":"/test/1.txt",
			"device_id":"64768",
			"inode":"4587523"
		},
		"offset":1128
	},
	"message":"hello filebeat",
	"input":{"type":"filestream"},
	"ecs":{"version":"8.0.0"}
}

image-20251215133018569

2,logstash读取kafka数据

· 生成JVM信任证书

1,复制kafka的ca证书

[root@kibana ~ ]# scp 10.0.0.111:/tools/kafka_2.13-3.3.2/certs/kafka-ca.crt /tools/logstash-8.18.8/certs/

2,使用JAVA证书工具生成证书

[root@kibana ~ ]# /tools/jdk1.8.0_461/bin/keytool -import   \
	-alias kafka-ca  \
	-file /tools/logstash-8.18.8/certs/kafka-ca.crt \
    -keystore /tools/logstash-8.18.8/certs/kafka.truststore.jks \
    -storepass 123456
.......
Trust this certificate? [no]:  y   # 这里需要确认【y】


#查看生成的证书
[root@kibana ~ ]# ls -l /tools/logstash-8.18.8/certs/
......
-rw-r--r-- 1 root root 1506 Dec 15 05:46 kafka.truststore.jks

· 编辑logstash配置文件

[root@kibana ~ ]# vim /tools/logstash-8.18.8/conf.d/test.conf 
input {
  kafka {
    bootstrap_servers => "10.0.0.111:9092,10.0.0.112:9092,10.0.0.113:9092"

    topics => ["bakwite-topic01"]
    group_id => "xjzw-gp"

    # ===== 安全 =====
    security_protocol => "SASL_SSL"
    sasl_mechanism => "PLAIN"
    # 用户名密码
    sasl_jaas_config => "org.apache.kafka.common.security.plain.PlainLoginModule required username='logstash' password='123456';"
    # JVM专有证书
    ssl_truststore_location => "/tools/logstash-8.18.8/certs/kafka.truststore.jks"
    ssl_truststore_password => "123456"

    # 从最新的offset开始消费
    # earliest 表示从头开始消费
    # latest 表示只读最新的,之前的不读;
    auto_offset_reset => "earliest"
    # 自动提交消费offset信息
    enable_auto_commit => true

    # 消费者线程(线程越多速度越快,跟cpu核心有关)
    consumer_threads => 3
    # 最慢3秒拉取一次消息
    poll_timeout_ms => 3000
    # 一次最多拉取500条消息
    max_poll_records => 500
  }
}

filter {
  mutate {
    remove_field => ["log","host","event","@version","host"]
  }
}

output {
  elasticsearch {
    hosts => ["https://10.0.0.101:10200", "https://10.0.0.102:10200", "https://10.0.0.103:10200"]
    index => "kafka-%{+YYYY.MM.dd}"
    ssl => true
    ssl_certificate_verification => true
    cacert => "/tools/logstash-8.18.8/certs/ca.crt"
    user => "elastic"
    password => "123456"
  }
}

· 重启logstash

[root@kibana ~ ]# systemctl restart logstash.service

· 写入数据测试

1,filebeat输入文件写入内容

[root@kibana ~ ]# echo "hello bakwite" >> /test/1.txt

2,查看kibana的索引

image-20251215143941298

3,查看消费者组验证

[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-consumer-groups.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf \
  --all-groups \
  --list 
  
[root@kafka01 ~ ]# /tools/kafka_2.13-3.3.2/bin/kafka-consumer-groups.sh \
  --bootstrap-server 10.0.0.111:9092 \
  --command-config login-kafka.conf \
  --group xjzw-gp \
  --describe

END

Discussion

评论

加载中
正在检查登录状态…