📌 本文是「Kafka 3.9 教程」系列第 五 篇 · 系列目录
本文档基于 Apache Kafka 3.9 官方文档整理
目录
5.1 多数据中心部署
5.1.1 部署架构
推荐架构:每个数据中心部署本地 Kafka 集群
数据中心 A 数据中心 B
┌────────────────────┐ ┌────────────────────┐
│ Kafka 集群 A │ │ Kafka 集群 B │
│ ┌──────────────┐ │ │ ┌──────────────┐ │
│ │ Broker 1 │ │ │ │ Broker 1 │ │
│ │ Broker 2 │ │ │ │ Broker 2 │ │
│ │ Broker 3 │ │ │ │ Broker 3 │ │
│ └──────────────┘ │ │ └──────────────┘ │
└────────┬─────────┬─┘ └────────┬─────────┬─┘
│ │ │ │
│ │ │ │
应用实例 A 应用实例 B
│ │ │ │
└────┬────┘ └────┬────┘
│ │
└──────────────┬───────────────────────┘
│
┌────────┴────────┐
│ MirrorMaker │
│ 跨集群复制 │
└─────────────────┘
5.1.2 优势
- 独立运作:每个数据中心可独立运行
- 容灾能力:跨数据中心链路故障时,本地操作不受影响
- 低延迟:应用只与本地集群交互
- 集中管理:跨集群复制集中配置和管理
5.1.3 注意事项
不要在单一集群中跨越多个数据中心:
- 高复制延迟
- ZooKeeper/Kafka 可用性受影响
- 网络分区时不可用
5.1.4 网络配置
对于 WAN(广域网)环境:
# 增加 Socket 缓冲区大小
socket.send.buffer.bytes=65536
socket.receive.buffer.bytes=65536
socket.request.max.bytes=104857600
# 增加超时时间
request.timeout.ms=60000
5.2 跨集群复制(Geo-Replication)
5.2.1 MirrorMaker 2.0
MirrorMaker 2.0 是 Kafka 3.0+ 推荐的跨集群复制工具。
核心特性:
- 监控源集群主题变化
- 自动创建目标集群主题
- 复制配置同步
- 消费位置追踪
配置示例
# mm2-properties
# 源集群配置
clusters=source,destination
source.bootstrap.servers=source-cluster:9092
destination.bootstrap.servers=dest-cluster:9092
source->destination.enabled=true
source->destination.topics=.* # 复制所有主题
# ACL 复制
source->destination.emit.acls.enabled=false
# 消费者偏移量同步
source->destination.sync.group.offsets.enabled=true
source->destination.sync.group.offsets.interval.seconds=60
启动:
bin/connect-mirror-maker.sh mm2-properties
5.2.2 传统 MirrorMaker
# consumer.properties
bootstrap.servers=source-cluster:9092
group.id=mirror-maker
auto.offset.reset=earliest
# producer.properties
bootstrap.servers=dest-cluster:9092
acks=all
retries=3
bin/kafka-mirror-maker.sh \
--consumer.config consumer.properties \
--producer.config producer.properties \
--whitelist "topic1,topic2" \
--num.streams 4
5.2.3 复制策略
| 策略 | 说明 | 适用场景 |
|---|---|---|
| 主动-被动 | 只读源集群,写目标集群 | 灾备 |
| 主动-主动 | 双向复制 | 多活架构 |
| 聚合 | 多个集群聚合到一个 | 数据汇总 |
5.3 监控
5.3.1 监控指标分类
| 分类 | 指标示例 | 说明 |
|---|---|---|
| Broker | UnderReplicatedPartitions | 副本状态 |
| 主题 | TotalProduceRequests | 生产请求 |
| 消费者 | records-lag | 消费延迟 |
| ZooKeeper | ZooKeeperRequestLatencyMs | ZK 延迟 |
5.3.2 关键指标
Broker 级别指标:
| 指标 | 描述 | 告警阈值 |
|---|---|---|
| UnderReplicatedPartitions | 副本不足的分区数 | > 0 |
| ISRExpansionRate | ISR 扩展速率 | 持续 > 0 |
| ISRShrinkRate | ISR 收缩速率 | 持续 > 0 |
| LeaderCount | Leader 分区数 | 异常变化 |
| RequestHandlerAvgIdlePercent | 请求处理空闲率 | < 0.3 |
| PurgatorySize | 待处理请求数 | 持续增长 |
Producer 指标:
| 指标 | 描述 |
|---|---|
| request-rate | 请求速率 |
| request-latency-ms | 请求延迟 |
| batch-size-avg | 平均批次大小 |
| compression-rate-avg | 平均压缩率 |
| record-error-rate | 错误率 |
Consumer 指标:
| 指标 | 描述 |
|---|---|
| records-lag | 消费延迟(消息数) |
| records-lag-max | 最大消费延迟 |
| fetch-rate | 拉取速率 |
| fetch-latency-ms | 拉取延迟 |
5.3.3 监控工具
JMX 监控
# 启用 JMX
JMX_PORT=9999 bin/kafka-server-start.sh server.properties
Prometheus + Grafana
Prometheus 配置:
# prometheus.yml
scrape_configs:
- job_name: 'kafka'
static_configs:
- targets: ['localhost:9108']
metrics_path: /metrics
关键 JMX 指标:
| JMX MBean | 指标 |
|---|---|
| kafka.server:type=BrokerTopicMetrics,name=TotalProduceRequests | 生产请求数 |
| kafka.server:type=BrokerTopicMetrics,name=TotalFetchRequests | 消费请求数 |
| kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions | 副本不足 |
| kafka.consumer:type=consumer-fetch-manager-metrics,records-lag | 消费延迟 |
Kafka Manager
# 启动 Kafka Manager
bin/kafka-manager -Dhttp.port=9000
5.3.4 告警配置
# 告警规则示例 (Prometheus)
groups:
- name: kafka-alerts
rules:
- alert: KafkaUnderReplicatedPartitions
expr: kafka_server_replicamanager_underreplicatedpartitions > 0
for: 5m
labels:
severity: critical
annotations:
summary: "Kafka 副本不足"
- alert: KafkaConsumerLag
expr: kafka_consumer_records_lag_max > 10000
for: 5m
labels:
severity: warning
annotations:
summary: "Kafka 消费延迟过高"
5.4 KRaft 模式
5.4.1 什么是 KRaft?
KRaft 是 Kafka 2.8+ 引入的内置共识协议,用于替代 ZooKeeper。
优势:
- 简化部署:减少外部依赖
- 提高可扩展性:元数据管理更高效
- 简化运维:无需管理 ZooKeeper 集群
5.4.2 配置示例
# server.properties (KRaft 模式)
# 进程角色
process.roles=broker,controller
# 节点 ID
node.id=1
# Controller 投票者(所有 Controller 节点)
controller.quorum.voters=1@localhost:9093,2@broker2:9093,3@broker3:9093
# Controller 监听器
controller.listener.names=CONTROLLER
# Broker 监听器
listeners=PLAINTEXT://0.0.0.0:9092
# 对外公告地址
advertised.listeners=PLAINTEXT://localhost:9092
# 日志目录
log.dirs=/tmp/kafka-logs
# 集群 ID(首次启动自动生成,也可手动指定)
# cluster.id=xxx
5.4.3 集群启动顺序
- 格式化存储:
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties
- 启动所有节点:
# 所有节点启动
bin/kafka-server-start.sh config/kraft/server.properties
5.4.4 从 ZooKeeper 迁移到 KRaft
Kafka 3.0+ 支持从 ZooKeeper 模式滚动升级到 KRaft 模式。
迁移步骤:
- 升级到 Kafka 3.x(使用 ZooKeeper)
- 配置 KRaft 相关参数
- 执行迁移
- 重启 Broker(使用 KRaft)
5.5 分层存储(Tiered Storage)
5.5.1 概述
分层存储将数据分为热数据和冷数据:
- 热数据:存储在本地 SSD(高性能)
- 冷数据:存储在远程存储(如 S3,成本更低)
5.5.2 配置
# 启用分层存储
log.local.retention.bytes=10737418240 # 10GB 热数据
log.local.retention.ms=259200000 # 3天热数据
remote.log.storage.dir=/s3-bucket/kafka
remote.log.storage.system.enable=true
5.5.3 优势
| 优势 | 说明 |
|---|---|
| 降低存储成本 | 冷数据存 S3 更便宜 |
| 延长保留期 | 可保留更长时间的历史数据 |
| 提高性能 | 热数据只在 SSD 上 |
| 简化运维 | 自动分层管理 |
5.6 多租户
5.6.1 资源隔离
主题隔离:
- 按租户划分主题命名空间
- 使用前缀或后缀区分
# 租户 A 的主题
kafka-topics.sh --create --topic tenant-a-orders --bootstrap-server localhost:9092
# 租户 B 的主题
kafka-topics.sh --create --topic tenant-b-orders --bootstrap-server localhost:9092
5.6.2 Quota 管理
生产配额:
# 设置生产速率配额(字节/秒)
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
--add --allow-principal User:tenant-a \
--operation Write \
--topic tenant-a-*
kafka-configs.sh --alter \
--entity-type users \
--entity-name tenant-a \
--add-config producer_byte_rate=10485760 \
--bootstrap-server localhost:9092
消费配额:
# 设置消费速率配额(字节/秒)
kafka-configs.sh --alter \
--entity-type users \
--entity-name tenant-b \
--add-config consumer_byte_rate=20971520 \
--bootstrap-server localhost:9092
5.6.3 ACL 权限控制
# 生产权限
kafka-acls.sh --add \
--allow-principal User:tenant-a \
--operation Write \
--topic tenant-a-orders
# 消费权限
kafka-acls.sh --add \
--allow-principal User:tenant-a \
--operation Read \
--group tenant-a-group \
--topic tenant-a-orders
# 描述权限
kafka-acls.sh --list \
--topic tenant-a-orders \
--group tenant-a-group
附录:运维命令速查
主题管理
# 列出所有主题
kafka-topics.sh --list --bootstrap-server localhost:9092
# 创建主题
kafka-topics.sh --create --topic my-topic --bootstrap-server localhost:9092 \
--partitions 6 --replication-factor 3
# 查看主题详情
kafka-topics.sh --describe --topic my-topic --bootstrap-server localhost:9092
# 修改分区数
kafka-topics.sh --alter --topic my-topic --partitions 12 --bootstrap-server localhost:9092
# 删除主题
kafka-topics.sh --delete --topic my-topic --bootstrap-server localhost:9092
# 查看主题配置
kafka-configs.sh --describe --topic my-topic --bootstrap-server localhost:9092
消费者组管理
# 列出消费者组
kafka-consumer-groups.sh --list --bootstrap-server localhost:9092
# 查看消费者组详情
kafka-consumer-groups.sh --describe --group my-group --bootstrap-server localhost:9092
# 查看消费偏移量
kafka-consumer-groups.sh --describe --group my-group \
--topic my-topic --bootstrap-server localhost:9092
# 重置偏移量
kafka-consumer-groups.sh --reset-offsets \
--group my-group \
--topic my-topic \
--to-earliest \
--execute \
--bootstrap-server localhost:9092
Broker 管理
# 查看 Broker 列表
kafka-metadata-shell.sh --describe --entity-type brokers
# 查看 Broker 配置
kafka-configs.sh --describe --entity-type brokers --entity-name 0 --bootstrap-server localhost:9092
性能测试
# 生产者性能测试
kafka-producer-perf-test.sh \
--topic my-topic \
--num-records 1000000 \
--throughput 100000 \
--record-size 1000 \
--producer-props bootstrap.servers=localhost:9092
# 消费者性能测试
kafka-consumer-perf-test.sh \
--topic my-topic \
--messages 1000000 \
--broker-list localhost:9092
下一步学习
- 第六章:安全认证 - 安全配置详解
- 第七章:Kafka Connect - 数据集成框架
- 第八章:Kafka Streams - 流处理库