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
| 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
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
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
| 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
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
| systemctl daemon-reload
systemctl enable zookeeper && systemctl start zookeeper
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
| /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
| 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
| 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 集群扩容
七、故障排查
7.1 Broker 无法启动
1 2 3 4 5 6 7 8 9 10 11
| /opt/kafka/bin/zookeeper-shell.sh zk1.example.com:2181 ls /brokers/ids
tail -f /data/kafka/logs/server.log
|
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
|
7.3 ISR 收缩
1 2 3 4 5 6 7 8 9 10
| /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka1.example.com:9092 \ --describe \ --topic orders
|
八、安全配置
8.1 SASL 认证
1 2 3 4 5 6 7 8
| 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
/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
| -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 作为高吞吐量的分布式消息系统,在生产环境中需要重点关注:
- 高可用:合理设置副本数和 min.insync.replicas
- 性能:根据业务场景调整分区数和批处理参数
- 监控:建立完善的监控告警体系
- 安全:启用认证和权限控制
- 容量规划:定期评估磁盘和网络资源
通过合理的部署和运维策略,Kafka 能够稳定支撑大规模实时数据处理场景。