Kafka 消息队列集群部署与运维

一、概述

Apache Kafka 是一个分布式流处理平台,广泛用于构建实时数据管道和流式应用。本文介绍 Kafka 集群的部署、配置及日常运维要点。

核心概念

概念 说明
Broker Kafka 服务器节点
Topic 消息分类主题
Partition 主题分区,实现并行处理
Replica 分区副本,保证高可用
Producer 消息生产者
Consumer 消息消费者
Consumer Group 消费者组,实现负载均衡
ZooKeeper 集群元数据管理(Kafka 2.8+ 可启用 KRaft 模式)

二、部署规划

2.1 硬件要求

组件 最低配置 推荐配置
Broker 4 核 8G 8 核 16G+
ZooKeeper 2 核 4G 4 核 8G
磁盘 SSD 100G NVMe SSD 500G+
网络 1Gbps 10Gbps

2.2 集群架构

1
2
3
4
5
6
7
8
9
10
11
                 ┌─────────────┐
│ ZooKeeper │
│ Cluster │
└──────┬──────┘

┌──────────────────┼──────────────────┐
│ │ │
┌────▼────┐ ┌────▼────┐ ┌────▼────┐
│ Broker 1│ │ Broker 2│ │ Broker 3│
│ :9092 │ │ :9092 │ │ :9092 │
└─────────┘ └─────────┘ └─────────┘

三、安装部署

3.1 下载与解压

1
2
3
4
5
6
# 下载 Kafka
wget https://downloads.apache.org/kafka/3.6.1/kafka_2.13-3.6.1.tgz

# 解压
tar -xzf kafka_2.13-3.6.1.tgz -C /opt/
ln -s /opt/kafka_2.13-3.6.1 /opt/kafka

3.2 配置 ZooKeeper

1
2
3
4
5
6
7
8
9
10
11
12
13
# 创建数据目录
mkdir -p /data/zookeeper

# 配置 zookeeper.properties
cat > /opt/kafka/config/zookeeper.properties << EOF
dataDir=/data/zookeeper
clientPort=2181
maxClientCnxns=0
admin.enableServer=false
tickTime=2000
initLimit=10
syncLimit=5
EOF

3.3 配置 Kafka Broker

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
30
31
32
33
34
35
# 创建数据目录
mkdir -p /data/kafka/logs

# 配置 server.properties(Broker 1 示例)
cat > /opt/kafka/config/server.properties << EOF
# 基础配置
broker.id=1
listeners=PLAINTEXT://:9092
advertised.listeners=PLAINTEXT://kafka1.example.com:9092

# ZooKeeper 连接
zookeeper.connect=zk1.example.com:2181,zk2.example.com:2181,zk3.example.com:2181

# 日志配置
log.dirs=/data/kafka/logs
num.partitions=3
default.replication.factor=3
min.insync.replicas=2

# 性能优化
num.network.threads=8
num.io.threads=16
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600

# 日志保留策略
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000

# 删除主题允许
delete.topic.enable=true
auto.create.topics.enable=false
EOF

3.4 创建 Systemd 服务

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
30
31
32
33
34
35
36
# ZooKeeper 服务
cat > /etc/systemd/system/zookeeper.service << EOF
[Unit]
Description=Apache ZooKeeper
After=network.target

[Service]
Type=forking
User=kafka
Group=kafka
ExecStart=/opt/kafka/bin/zookeeper-server-start.sh -daemon /opt/kafka/config/zookeeper.properties
ExecStop=/opt/kafka/bin/zookeeper-server-stop.sh
Restart=on-failure

[Install]
WantedBy=multi-user.target
EOF

# Kafka 服务
cat > /etc/systemd/system/kafka.service << EOF
[Unit]
Description=Apache Kafka
After=network.target zookeeper.service
Requires=zookeeper.service

[Service]
Type=forking
User=kafka
Group=kafka
ExecStart=/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties
ExecStop=/opt/kafka/bin/kafka-server-stop.sh
Restart=on-failure

[Install]
WantedBy=multi-user.target
EOF

3.5 启动服务

1
2
3
4
5
6
7
8
9
10
11
# 重载 systemd
systemctl daemon-reload

# 启动 ZooKeeper
systemctl enable zookeeper && systemctl start zookeeper

# 启动 Kafka
systemctl enable kafka && systemctl start kafka

# 检查状态
systemctl status zookeeper kafka

四、集群管理

4.1 创建 Topic

1
2
3
4
5
6
7
8
9
10
11
12
13
# 创建主题(3 分区,3 副本)
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka1.example.com:9092 \
--create \
--topic orders \
--partitions 3 \
--replication-factor 3

# 查看主题详情
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka1.example.com:9092 \
--describe \
--topic orders

4.2 生产与消费测试

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
# 生产者测试
/opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server kafka1.example.com:9092 \
--topic orders

# 消费者测试(从最新消息开始)
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka1.example.com:9092 \
--topic orders \
--from-beginning

# 消费者组管理
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server kafka1.example.com:9092 \
--list

/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server kafka1.example.com:9092 \
--describe \
--group my-consumer-group

4.3 分区扩容

1
2
3
4
5
6
# 增加分区数(只能增加不能减少)
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka1.example.com:9092 \
--alter \
--topic orders \
--partitions 6

五、监控与告警

5.1 JMX 监控配置

1
2
3
4
5
6
# 在 kafka-server-start.sh 中添加 JMX 参数
export JMX_PORT=9999
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote \
-Dcom.sun.management.jmxremote.authenticate=false \
-Dcom.sun.management.jmxremote.ssl=false \
-Djava.rmi.server.hostname=<broker-ip>"

5.2 关键监控指标

指标 说明 告警阈值
UnderReplicatedPartitions 未同步副本数 > 0
OfflinePartitionsCount 离线分区数 > 0
RequestHandlerAvgIdlePercent 请求处理空闲率 < 30%
NetworkProcessorAvgIdlePercent 网络处理空闲率 < 30%
LogFlushRateAndTimeMs 日志刷盘延迟 > 1000ms

5.3 Prometheus Exporter 配置

1
2
3
4
5
6
# kafka_exporter 配置
scrape_configs:
- job_name: 'kafka'
static_configs:
- targets: ['kafka1.example.com:9308']
metrics_path: /metrics

六、日常运维

6.1 日志清理

1
2
# 手动触发日志清理
/opt/kafka/bin/kafka-run-class.sh kafka.admin.LogCleaner

6.2 数据迁移

1
2
3
4
5
6
7
8
9
10
11
12
13
14
# 重新分配分区副本
cat > /tmp/reassignment.json << EOF
{
"version": 1,
"partitions": [
{"topic": "orders", "partition": 0, "replicas": [1, 2, 3]}
]
}
EOF

/opt/kafka/bin/kafka-reassign-partitions.sh \
--bootstrap-server kafka1.example.com:9092 \
--reassignment-json-file /tmp/reassignment.json \
--execute

6.3 集群扩容

1
2
3
4
# 1. 部署新 Broker
# 2. 配置 broker.id 和 advertised.listeners
# 3. 启动服务
# 4. 重新分配分区以平衡负载

七、故障排查

7.1 Broker 无法启动

1
2
3
4
5
6
7
8
9
10
11
# 检查 ZooKeeper 连接
/opt/kafka/bin/zookeeper-shell.sh zk1.example.com:2181 ls /brokers/ids

# 检查日志
tail -f /data/kafka/logs/server.log

# 常见问题:
# 1. broker.id 冲突
# 2. 端口被占用
# 3. 数据目录权限问题
# 4. ZooKeeper 连接失败

7.2 消息积压

1
2
3
4
5
6
7
8
9
10
# 查看消费者延迟
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server kafka1.example.com:9092 \
--describe \
--group my-group

# 解决方案:
# 1. 增加消费者实例
# 2. 增加分区数
# 3. 优化消费者处理逻辑

7.3 ISR 收缩

1
2
3
4
5
6
7
8
9
10
# 查看 ISR 状态
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka1.example.com:9092 \
--describe \
--topic orders

# 解决方案:
# 1. 检查网络连通性
# 2. 检查磁盘 IO 性能
# 3. 调整 replica.lag.time.max.ms

八、安全配置

8.1 SASL 认证

1
2
3
4
5
6
7
8
# server.properties
listeners=SASL_PLAINTEXT://:9092
security.inter.broker.protocol=SASL_PLAINTEXT
sasl.mechanism.inter.broker.protocol=PLAIN
sasl.enabled.mechanisms=PLAIN

authorizer.class.name=kafka.security.authorizer.AclAuthorizer
allow.everyone.if.no.acl.found=false

8.2 ACL 权限控制

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# 创建用户
/opt/kafka/bin/kafka-configs.sh \
--bootstrap-server kafka1.example.com:9092 \
--alter \
--add-config 'SCRAM-SHA-256=[password=secret123]' \
--entity-name appuser \
--entity-type users

# 设置 ACL
/opt/kafka/bin/kafka-acls.sh \
--bootstrap-server kafka1.example.com:9092 \
--add \
--allow-principal User:appuser \
--operation Read \
--operation Write \
--topic orders

九、性能调优

9.1 Producer 优化

1
2
3
4
5
6
7
# 生产者配置
batch.size=65536
linger.ms=5
buffer.memory=33554432
compression.type=snappy
acks=all
retries=3

9.2 Consumer 优化

1
2
3
4
5
6
# 消费者配置
fetch.min.bytes=1048576
fetch.max.wait.ms=500
max.partition.fetch.bytes=10485760
session.timeout.ms=30000
heartbeat.interval.ms=10000

9.3 Broker 优化

1
2
3
4
5
6
7
8
9
10
11
# JVM 参数
-Xms6g -Xmx6g
-XX:MetaspaceSize=96m -XX:MaxMetaspaceSize=256m
-XX:+UseG1GC
-XX:MaxGCPauseMillis=20
-XX:InitiatingHeapOccupancyPercent=35

# 系统参数
vm.swappiness=1
vm.dirty_ratio=80
vm.dirty_background_ratio=5

十、总结

Kafka 作为高吞吐量的分布式消息系统,在生产环境中需要重点关注:

  1. 高可用:合理设置副本数和 min.insync.replicas
  2. 性能:根据业务场景调整分区数和批处理参数
  3. 监控:建立完善的监控告警体系
  4. 安全:启用认证和权限控制
  5. 容量规划:定期评估磁盘和网络资源

通过合理的部署和运维策略,Kafka 能够稳定支撑大规模实时数据处理场景。