Kafka是分布式消息中间件中吞吐量最高的系统之一,单集群可支撑百万级TPS。在后端开发中,Kafka用于异步解耦、削峰填谷和事件驱动架构。高并发设计场景下,合理的分区分配策略和消费者组配置直接决定系统的消息处理能力和容错水平。
Kafka集群架构与Broker节点规划
Kafka集群由多个Broker节点组成,每个Broker存储部分分区数据。为保证高可用,每个分区有多个副本分布在不同Broker上,其中Leader副本负责读写,Follower副本同步数据。服务治理中,Broker数量建议为奇数以支持多数派选举,分区数根据消费者并发度规划。
# server.properties - Broker核心配置
broker.id=1
listeners=PLAINTEXT://0.0.0.0:9092
advertised.listeners=PLAINTEXT://kafka-1:9092
log.dirs=/data/kafka-logs
num.partitions=6
default.replication.factor=3
min.insync.replicas=2
log.retention.hours=168
log.segment.bytes=1073741824
log.cleanup.policy=delete
zookeeper.connect=zk-1:2181,zk-2:2181,zk-3:2181
auto.create.topics.enable=false
unclean.leader.election.enable=false
num.network.threads=8
num.io.threads=16
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
queued.max.requests=500
Topic创建与分区副本分配策略
Topic是Kafka的消息分类单位,分区是并行处理的基本单位。分区数决定了最大消费者并发数,副本数决定数据冗余级别。业务中台建设中,核心业务的Topic通常配置3副本、6到12分区,保证单分区故障不影响整体可用性。分区分配策略决定副本在Broker间的分布,影响容错能力。
# 创建Topic:指定分区数和副本数
kafka-topics.sh --create \
--bootstrap-server kafka-1:9092 \
--topic order-events \
--partitions 12 \
--replication-factor 3 \
--config min.insync.replicas=2 \
--config retention.ms=604800000 \
--config segment.bytes=1073741824 \
--config cleanup.policy=delete
# 查看Topic分区分布
kafka-topics.sh --describe \
--bootstrap-server kafka-1:9092 \
--topic order-events
# 手动指定副本分布
# rack-awareness:确保副本跨机架分布
kafka-topics.sh --alter \
--bootstrap-server kafka-1:9092 \
--topic order-events \
--replica-assignment 1:2:3,2:3:1,3:1:2,1:3:2,2:1:3,3:2:1
# 增加分区(只能增加不能减少)
kafka-topics.sh --alter \
--bootstrap-server kafka-1:9092 \
--topic order-events \
--partitions 24
# Java客户端创建Topic
Properties props = new Properties();
props.put("bootstrap.servers", "kafka-1:9092");
AdminClient admin = AdminClient.create(props);
NewTopic topic = new NewTopic("payment-events", 12, (short) 3);
Map<String, String> configs = new HashMap<>();
configs.put("retention.ms", "604800000");
configs.put("min.insync.replicas", "2");
topic.configs(configs);
admin.createTopics(Collections.singleton(topic));
admin.close();
消费者组负载均衡与分区再平衡机制
消费者组是Kafka实现消息并行处理的核心机制。同一组内的消费者平均分配分区,每个分区只被组内一个消费者消费。消费者加入或离开时触发再平衡(rebalance),分区重新分配。分布式事务场景中,再平衡期间的消息处理中断是常见的可用性风险,需要通过合理的配置和消费者设计来缓解。
// Java消费者配置与分区分配策略
Properties props = new Properties();
props.put("bootstrap.servers", "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put("group.id", "order-consumer-group");
props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", StringDeserializer.class.getName());
// 分区分配策略
// RangeAssignor(默认):按分区范围分配
// RoundRobinAssignor:轮询分配
// StickyAssignor:粘性分配,减少再平衡时的分区迁移
// CooperativeStickyAssignor:协作式粘性分配,增量再平衡
props.put("partition.assignment.strategy",
CooperativeStickyAssignor.class.getName());
// 消费偏移量管理
props.put("enable.auto.commit", "false");
props.put("auto.offset.reset", "earliest");
props.put("max.poll.records", "500");
props.put("max.poll.interval.ms", "300000");
props.put("session.timeout.ms", "45000");
props.put("heartbeat.interval.ms", "15000");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("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={}, error={}", record.offset(), e.getMessage());
handleFailure(record, e);
}
}
// 手动提交偏移量
if (!records.isEmpty()) {
consumer.commitSync();
}
}
生产者ACK机制与消息可靠性保障
Kafka生产者的acks参数决定消息确认级别。acks=0不等待确认,吞吐量最高但可能丢消息;acks=1等待Leader确认,Leader故障时可能丢消息;acks=all等待所有ISR副本确认,配合min.insync.replicas实现最高可靠性。消息中间件在高并发场景下需要在可靠性和吞吐量之间权衡。
// Java生产者配置
Properties props = new Properties();
props.put("bootstrap.servers", "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put("key.serializer", StringSerializer.class.getName());
props.put("value.serializer", StringSerializer.class.getName());
// ACK机制
props.put("acks", "all");
props.put("min.insync.replicas", "2");
// 重试配置
props.put("retries", 10);
props.put("retry.backoff.ms", "100");
props.put("max.in.flight.requests.per.connection", "5");
props.put("enable.idempotence", "true"); // 幂等生产者
props.put("transactional.id", "order-tx-1"); // 事务ID
// 批处理配置
props.put("batch.size", "65536");
props.put("linger.ms", "10");
props.put("buffer.memory", "67108864");
props.put("compression.type", "lz4");
// 超时配置
props.put("delivery.timeout.ms", "120000");
props.put("request.timeout.ms", "30000");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// 事务消息发送
producer.initTransactions();
try {
producer.beginTransaction();
// 发送多条消息
ProducerRecord<String, String> record1 = new ProducerRecord<>(
"order-events", "order-123", orderJson);
ProducerRecord<String, String> record2 = new ProducerRecord<>(
"payment-events", "payment-123", paymentJson);
producer.send(record1);
producer.send(record2);
// 提交事务
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
log.error("事务失败", e);
}
// 异步发送回调
producer.send(record, (metadata, exception) -> {
if (exception != null) {
log.error("发送失败: {}", exception.getMessage());
// 重试逻辑
} else {
log.debug("发送成功: topic={}, partition={}, offset={}",
metadata.topic(), metadata.partition(), metadata.offset());
}
});
消费者端幂等处理与消息积压应对方案
消息积压是Kafka运维中的常见问题。原因包括消费者处理速度不足、分区数过少或消费者实例数不足。API接口规范要求消费者具备幂等处理能力,避免重复消费导致数据不一致。应对积压的方案包括增加分区和消费者实例、批量处理、异步消费等。
// 幂等消费者实现
public class IdempotentConsumer {
private final KafkaConsumer<String, String> consumer;
private final RedisTemplate<String, String> redis;
private static final String PROCESSED_KEY = "kafka:processed:";
private static final int PROCESSED_TTL = 86400; // 24小时
public void consume() {
consumer.subscribe(Collections.singletonList("order-events"));
while (running) {
ConsumerRecords<String, String> records = consumer.poll(
Duration.ofMillis(1000));
List<ConsumerRecord<String, String>> batch = new ArrayList<>();
for (ConsumerRecord<String, String> record : records) {
String messageId = record.topic() + ":" +
record.partition() + ":" + record.offset();
// Redis SETNX 实现幂等判断
Boolean isNew = redis.opsForValue()
.setIfAbsent(PROCESSED_KEY + messageId, "1",
PROCESSED_TTL, TimeUnit.SECONDS);
if (Boolean.TRUE.equals(isNew)) {
batch.add(record);
}
// 已处理过的消息跳过
}
if (!batch.isEmpty()) {
// 批量处理
processBatch(batch);
consumer.commitSync();
}
}
}
// 批量处理提升吞吐量
private void processBatch(List<ConsumerRecord<String, String>> batch) {
List<OrderEvent> events = batch.stream()
.map(r -> JSON.parseObject(r.value(), OrderEvent.class))
.collect(Collectors.toList());
// 批量写入数据库
orderService.batchSave(events);
}
}
Kafka在高并发设计和微服务架构中的角色不仅是消息传递,更是系统解耦和弹性扩展的核心组件。通过合理的分区规划、消费者组管理和生产者配置,可以在保证消息可靠性的前提下实现高吞吐量。服务治理中,Kafka的监控指标(如消费者延迟、分区不平衡、ISR缩减)需要纳入告警体系,及时发现和处理潜在问题。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/kafka-xiao-xi-dui-lie-gao-ke-yong-bu-shu-shi-zhan-fen-qu/