Kafka是分布式流处理平台,以高吞吐、低延迟和水平扩展能力成为微服务架构中消息中间件的首选。Kafka核心设计围绕分区(Partition)和副本(Replica)机制展开,理解分区策略、副本同步和消费者组协调机制,是构建可靠消息系统的关键。本文从架构设计到生产部署,详解Kafka核心机制与配置方法。
Kafka集群架构核心组件
Kafka集群包含三类节点角色:
Broker:Kafka服务节点,存储分区数据,处理生产者和消费者请求
Controller:集群中的一个Broker选举为Controller,负责分区Leader选举和副本重新分配
Coordinator:消费者组协调器,管理消费者组成员加入、退出和分区重平衡
数据模型层次:Topic(主题)→ Partition(分区)→ Replica(副本)→ Segment(日志段)。
分区机制与并行度设计
Topic划分为多个Partition,每个Partition是一个有序的、不可变的消息序列。Partition是Kafka并行处理的基本单元:生产者并行写入不同分区,消费者组内每个消费者独立消费不同分区。
分区数量规划原则:
1. 分区数 ≥ 消费者数,确保每个消费者至少分配一个分区
2. 单分区吞吐量约10MB/s(机械磁盘)或50MB/s(SSD),总吞吐需求 / 单分区吞吐 = 最小分区数
3. 分区数不宜过多:每个分区占用Broker内存(replica.fetch.max.bytes)和文件句柄,建议单Broker不超过4000个分区
创建Topic并指定分区数和副本因子:
# 创建Topic,6个分区,3副本
kafka-topics.sh --create \
--bootstrap-server kafka1:9092 \
--topic order-events \
--partitions 6 \
--replication-factor 3 \
--config retention.ms=604800000 \
--config segment.bytes=1073741824 \
--config cleanup.policy=delete
# 查看Topic详情
kafka-topics.sh --describe \
--bootstrap-server kafka1:9092 \
--topic order-events
# 增加分区数(只能增加不能减少)
kafka-topics.sh --alter \
--bootstrap-server kafka1:9092 \
--topic order-events \
--partitions 12
分区分配策略与生产者消息路由
生产者发送消息时需要指定目标分区,Kafka提供三种分区策略:
// Java生产者配置
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka1:9092,kafka2:9092,kafka3:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.RETRIES_CONFIG, 3);
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
props.put(ProducerConfig.LINGER_MS_CONFIG, 10);
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// 1. 指定分区号发送
ProducerRecord<String, String> record1 = new ProducerRecord<>(
"order-events", 0, "key", "value"
);
// 2. 指定Key发送,相同Key始终路由到同一分区(保证局部顺序)
ProducerRecord<String, String> record2 = new ProducerRecord<>(
"order-events", "user-123", "value"
);
// 3. 不指定Key和分区,使用默认分区器(轮询或Sticky分区器)
ProducerRecord<String, String> record3 = new ProducerRecord<>(
"order-events", "value"
);
producer.send(record2, (metadata, exception) -> {
if (exception == null) {
System.out.printf("Sent to partition %d, offset %d%n",
metadata.partition(), metadata.offset());
} else {
exception.printStackTrace();
}
});
acks参数控制消息持久化可靠性:
acks=0:生产者不等待Broker确认,最高吞吐但可能丢消息
acks=1:等待Leader副本写入确认,Leader宕机可能丢消息
acks=all(-1):等待所有ISR(In-Sync Replicas)副本写入确认,最高可靠性
幂等生产者配置(enable.idempotence=true)保证消息不重复,配合acks=all实现Exactly-Once语义。max.in.flight.requests.per.connection ≤ 5时保证顺序性。
副本同步机制与ISR管理
每个Partition有多个副本,其中一个为Leader,其余为Follower。所有读写请求由Leader处理,Follower从Leader拉取数据同步。
ISR(In-Sync Replicas)是当前与Leader保持同步的副本集合。Follower落后超过replica.lag.time.max.ms(默认30秒)时从ISR中移除。只有ISR中的副本才有资格被选为Leader。
# Broker配置 server.properties
# ISR最小同步副本数,acks=all时需要至少min.insync.replicas个副本在ISR中
min.insync.replicas=2
# Follower同步超时
replica.lag.time.max.ms=30000
replica.fetch.wait.max.ms=500
# 副本获取线程数
num.replica.fetchers=4
# Unclean Leader Election:是否允许非ISR副本成为Leader(数据丢失风险)
unclean.leader.election.enable=false
unclean.leader.election.enable=false表示当ISR中没有可用副本时,分区不可用而非选择数据不完整的副本,保证数据一致性。生产环境建议设为false。
消费者组与分区重平衡机制
消费者组(Consumer Group)实现发布订阅和点对点两种模式。同一组内每个分区只能被一个消费者消费,实现负载均衡。
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka1:9092,kafka2:9092,kafka3:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-processor-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500);
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300000);
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 10000);
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("order-events"));
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
try {
// 处理消息
processOrder(record.value());
} catch (Exception e) {
log.error("处理失败, offset={}", record.offset(), e);
// 处理失败时跳过或转入死信队列
}
}
// 手动同步提交offset
consumer.commitSync();
}
重平衡(Rebalance)触发条件:
1. 消费者加入或离开消费者组
2. 消费者心跳超时(session.timeout.ms内未发送心跳)
3. 消费者处理超时(max.poll.interval.ms内未调用poll)
4. Topic分区数变化
重平衡期间所有消费者停止消费,导致短暂不可用。通过Cooperative Rebalance(增量重平衡)减少影响:
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
"org.apache.kafka.clients.consumer.CooperativeStickyAssignor");
消息可靠投递与Exactly-Once语义
Kafka事务实现跨分区的原子写入,用于流处理中消费-处理-生产的Exactly-Once语义:
// 事务生产者配置
Properties props = new Properties();
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "tx-producer-1");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.ACKS_CONFIG, "all");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
// 消费者配置read_committed隔离级别
Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
// 事务处理
try {
producer.beginTransaction();
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
// 处理消息并发送结果到输出Topic
producer.send(new ProducerRecord<>("output-topic",
record.key(), transform(record.value())));
}
// 提交消费者offset(事务内)
producer.sendOffsetsToTransaction(
getCurrentOffsets(consumer), consumer.groupMetadata());
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
}
生产环境监控与运维
关键监控指标:
# 使用kafka-exporter + Prometheus + Grafana监控
# 核心JMX指标
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions
# 副本不同步的分区数,应始终为0
kafka.server:type=ReplicaManager,name=UnderMinIsrPartitionCount
# ISR低于min.insync.replicas的分区数,应始终为0
kafka.server:type=ControllerStats,name=ActiveControllerCount
# 活跃Controller数,应始终为1
kafka.network:type=RequestMetrics,name=TotalTimeMs,request=Produce
# 生产请求延迟P99,应低于50ms
常见运维操作:
# 优雅扩容:新增Broker后重新分配分区
kafka-reassign-partitions.sh --generate \
--bootstrap-server kafka1:9092 \
--topics-to-move-json-file topics.json \
--broker-list "4" > reassignment.json
kafka-reassign-partitions.sh --execute \
--bootstrap-server kafka1:9092 \
--reassignment-json-file reassignment.json
kafka-reassign-partitions.sh --verify \
--bootstrap-server kafka1:9092 \
--reassignment-json-file reassignment.json
# 消费者组管理
kafka-consumer-groups.sh --list --bootstrap-server kafka1:9092
kafka-consumer-groups.sh --describe --group order-processor-group --bootstrap-server kafka1:9092
kafka-consumer-groups.sh --reset-offsets --group order-processor-group \
--topic order-events --to-earliest --execute --bootstrap-server kafka1:9092
Kafka通过分区并行、副本冗余和消费者组协调机制,实现了高吞吐、高可用的消息处理能力。生产环境部署需重点配置acks策略、min.insync.replicas、幂等生产和消费者提交策略,在性能与可靠性间取得平衡。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/kafka-gao-tun-tu-xiao-xi-dui-lie-jia-gou-she-ji-yu-fen-qu/