Kafka是分布式消息中间件的核心选型,在高并发设计、微服务架构事件驱动通信、日志收集等场景中广泛应用。Kafka集群部署涉及Broker配置、Topic分区策略、消费者组负载均衡三个关键环节。合理的分区数和消费者实例配比直接影响吞吐量和消息有序性,是后端开发中消息中间件调优的重点。
Kafka集群架构与Broker配置
Kafka集群由多个Broker节点组成,每个Broker存储部分分区副本。Topic的分区分布在不同Broker上,通过副本机制实现高可用。Leader副本处理读写请求,Follower副本同步数据,Leader宕机时Follower自动选举接管。
# docker-compose.yml - 3节点Kafka集群(KRaft模式,无需ZooKeeper)
version: "3.8"
services:
kafka-1:
image: bitnami/kafka:3.7.0
container_name: kafka-1
ports:
- "9092:9092"
environment:
KAFKA_CFG_NODE_ID: 1
KAFKA_CFG_PROCESS_ROLES: "controller,broker"
KAFKA_CFG_LISTENERS: "PLAINTEXT://:9092,CONTROLLER://:9093"
KAFKA_CFG_ADVERTISED_LISTENERS: "PLAINTEXT://kafka-1:9092"
KAFKA_CFG_CONTROLLER_LISTENER_NAMES: "CONTROLLER"
KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: "1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093"
KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: "CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT"
ALLOW_PLAINTEXT_LISTENER: "yes"
volumes:
- kafka1_data:/bitnami/kafka
kafka-2:
image: bitnami/kafka:3.7.0
container_name: kafka-2
ports:
- "9094:9092"
environment:
KAFKA_CFG_NODE_ID: 2
KAFKA_CFG_PROCESS_ROLES: "controller,broker"
KAFKA_CFG_LISTENERS: "PLAINTEXT://:9092,CONTROLLER://:9093"
KAFKA_CFG_ADVERTISED_LISTENERS: "PLAINTEXT://kafka-2:9092"
KAFKA_CFG_CONTROLLER_LISTENER_NAMES: "CONTROLLER"
KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: "1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093"
KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: "CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT"
ALLOW_PLAINTEXT_LISTENER: "yes"
volumes:
- kafka2_data:/bitnami/kafka
kafka-3:
image: bitnami/kafka:3.7.0
container_name: kafka-3
ports:
- "9095:9092"
environment:
KAFKA_CFG_NODE_ID: 3
KAFKA_CFG_PROCESS_ROLES: "controller,broker"
KAFKA_CFG_LISTENERS: "PLAINTEXT://:9092,CONTROLLER://:9093"
KAFKA_CFG_ADVERTISED_LISTENERS: "PLAINTEXT://kafka-3:9092"
KAFKA_CFG_CONTROLLER_LISTENER_NAMES: "CONTROLLER"
KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: "1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093"
KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: "CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT"
ALLOW_PLAINTEXT_LISTENER: "yes"
volumes:
- kafka3_data:/bitnami/kafka
volumes:
kafka1_data:
kafka2_data:
kafka3_data:
Topic创建与分区策略配置
分区数决定并行度,每个分区在同一消费者组内只能被一个消费者实例消费。分区数应大于等于消费者实例数,否则有消费者空闲:
# 创建Topic:12个分区,3副本
kafka-topics.sh --create \
--bootstrap-server kafka-1:9092 \
--topic order-events \
--partitions 12 \
--replication-factor 3
# 查看Topic详情
kafka-topics.sh --describe \
--bootstrap-server kafka-1:9092 \
--topic order-events
# 修改分区数(只能增加,不能减少)
kafka-topics.sh --alter \
--bootstrap-server kafka-1:9092 \
--topic order-events \
--partitions 24
Java生产者客户端配置与发送策略
import org.apache.kafka.clients.producer.*;
import java.util.Properties;
public class OrderEventProducer {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// acks=all:等待所有副本确认,保证不丢消息
props.put("acks", "all");
// 启用幂等生产者,防止网络重试导致重复
props.put("enable.idempotence", "true");
// 批量发送:缓冲区大小和等待时间
props.put("batch.size", 16384);
props.put("linger.ms", 10);
// 重试配置
props.put("retries", 3);
props.put("max.in.flight.requests.per.connection", 5);
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
for (int i = 0; i < 1000; i++) {
String orderId = "order-" + i;
String eventData = String.format("{"id":"%s","amount":%.2f}", orderId, i * 10.5);
// 使用orderId作为key,保证同一订单的事件进入同一分区
ProducerRecord<String, String> record =
new ProducerRecord<>("order-events", orderId, eventData);
producer.send(record, (metadata, e) -> {
if (e != null) {
System.err.println("发送失败: " + e.getMessage());
} else {
System.out.printf("发送成功: partition=%d, offset=%d%n",
metadata.partition(), metadata.offset());
}
});
}
producer.flush();
producer.close();
}
}
消费者组负载均衡与消息消费
消费者组分区分配策略
Kafka支持三种分区分配策略,通过partition.assignment.strategy配置:
import org.apache.kafka.clients.consumer.*;
import java.time.Duration;
import java.util.*;
public class OrderEventConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put("group.id", "order-processing-group");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
// 从最早未消费的位置开始
props.put("auto.offset.reset", "earliest");
// 手动提交offset
props.put("enable.auto.commit", "false");
// 分区分配策略:CooperativeStickyAssignor(平滑再均衡)
props.put("partition.assignment.strategy",
"org.apache.kafka.clients.consumer.CooperativeStickyAssignor");
// 单次poll最大记录数
props.put("max.poll.records", 500);
// 消费者处理消息的最大时间
props.put("max.poll.interval.ms", 300000);
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("order-events"));
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
// 处理消息
processOrder(record.value());
}
// 处理完成后手动提交offset
if (!records.isEmpty()) {
consumer.commitSync();
}
}
} finally {
consumer.close();
}
}
private static void processOrder(String eventData) {
// 业务处理逻辑
System.out.println("处理订单: " + eventData);
}
}
消费者扩缩容时的再均衡行为
当消费者组中实例数变化时,Kafka触发再均衡(rebalance),重新分配分区。CooperativeStickyAssignor策略下,再均衡只迁移需要变更的分区,不影响其他消费者继续消费。默认的RangeAssignor策略会导致全量再均衡,所有消费者短暂停止消费。
# 消费者组状态查看
kafka-consumer-groups.sh --describe \
--bootstrap-server kafka-1:9092 \
--group order-processing-group
# 输出示例(12分区,3消费者,每人4分区):
# GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
# order-processing-group order-events 0 15000 15200 200
# order-processing-group order-events 1 14800 15000 200
# ...
消息有序性与Exactly-Once语义保障
分区消息有序性
Kafka保证单个分区内消息有序,跨分区不保证顺序。需要全局有序的场景只能使用单分区,但会丧失并行能力。实践中按业务Key路由到同一分区即可保证同一业务实体的消息有序。
Exactly-Once事务配置
// 生产者端配置事务
props.put("transactional.id", "order-tx-1");
producer.initTransactions();
try {
producer.beginTransaction();
producer.send(new ProducerRecord<>("order-events", key, value));
producer.send(new ProducerRecord<>("inventory-events", key, inventoryUpdate));
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
}
消费者端配合事务需要将读取消息和写入结果作为原子操作:
// 消费-处理-写入另一个Topic的Exactly-Once
props.put("isolation.level", "read_committed"); // 只读取已提交事务的消息
Kafka集群运维与监控指标
关键监控指标包括:UnderReplicatedPartitions(未同步分区数,应为0)、OfflinePartitions(离线分区数,应为0)、ActiveControllerCount(活跃控制器数,必须为1)、MessagesInPerSec(入站消息速率)、BytesInPerSec/BytesOutPerSec(吞吐量)。
分区数规划建议:单个分区吞吐量约10MB/s,集群目标吞吐量除以单分区吞吐量得到最小分区数。消费者实例数不超过分区数,超过部分空闲。Topic保留时间根据业务需求配置,日志类Topic可设7天,业务事件类可设永久或更长。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/kafka-xiao-xi-zhong-jian-jian-ji-qun-bu-shu-yu-xiao-fei-zhe/