Kafka是分布式流处理平台,在高吞吐消息场景下广泛应用。单机Kafka可支撑百万级TPS,但生产环境需要多节点集群保障可用性和数据可靠性。Kafka的高可用架构依赖分区副本机制——每个Partition有多个副本分布在不同Broker上,Leader副本处理读写请求,Follower副本同步数据,Leader故障时Follower自动接管。
集群部署与分区副本配置
三节点Kafka集群(KRaft模式,无需ZooKeeper)配置示例。每个Broker的server.properties:
# Broker 1
broker.id=1
listeners=PLAINTEXT://10.0.1.11:9092
controller.quorum.voters=1@10.0.1.11:9093,2@10.0.1.12:9093,3@10.0.1.13:9093
process.roles=broker,controller
node.id=1
log.dirs=/data/kafka/logs
num.partitions=6
default.replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=false
# 三个Broker的broker.id和node.id分别为1、2、3
# listeners中IP分别为10.0.1.11、12、13
default.replication.factor=3设置新创建Topic默认3副本。min.insync.replicas=2要求至少2个副本确认写入才算成功,配合acks=all实现强持久性。unclean.leader.election.enable=false禁止非同步副本成为Leader,避免数据丢失。
创建Topic时指定分区数和副本数:
bin/kafka-topics.sh --bootstrap-server 10.0.1.11:9092 \
--create --topic order-events \
--partitions 12 \
--replication-factor 3
# 查看Topic分区分布
bin/kafka-topics.sh --bootstrap-server 10.0.1.11:9092 \
--describe --topic order-events
# 输出示例
Topic: order-events Partitions: 12 ReplicationFactor: 3
Partition 0: Leader=1 Replicas: 1,2,3 Isr: 1,2,3
Partition 1: Leader=2 Replicas: 2,3,1 Isr: 2,3,1
Partition 2: Leader=3 Replicas: 3,1,2 Isr: 3,1,2
...
分区数决定了并行度,一般设置为消费者数量的整数倍。Isr(In-Sync Replicas)列出已同步的副本集合,Isr数量小于min.insync.replicas时写入会失败。
生产者确认机制与重试策略
Java客户端生产者配置:
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "10.0.1.11:9092,10.0.1.12:9092,10.0.1.13: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.ENABLE_IDEMPOTENCE_CONFIG, "true");
props.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE);
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "5");
// 批量与压缩
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 32768);
props.put(ProducerConfig.LINGER_MS_CONFIG, 10);
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
// 超时
props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 120000);
props.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 30000);
props.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 100);
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
acks=all要求所有Isr副本确认写入,持久性最强但延迟最高。enable.idempotence=true启用幂等生产者,配合max.in.flight.requests.per.connection <= 5保证同分区消息顺序性,避免重试导致乱序。
batch.size控制批次大小(字节),linger.ms控制批次等待时间(毫秒)。linger.ms=10表示最多等待10ms积攒消息后发送,以少量延迟换取更高吞吐。compression.type=lz4压缩消息,减少网络传输量和磁盘占用。
发送消息并处理回调:
ProducerRecord<String, String> record = new ProducerRecord<>(
"order-events",
orderId, // key,保证同key消息到同一分区
jsonPayload
);
producer.send(record, (metadata, exception) -> {
if (exception != null) {
log.error("发送失败: orderId={}, error={}", orderId, exception.getMessage());
// 写入死信队列或本地补偿
deadLetterQueue.send(record);
} else {
log.debug("发送成功: partition={}, offset={}",
metadata.partition(), metadata.offset());
}
});
消费者组与位移管理
消费者配置与手动位移提交:
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "10.0.1.11:9092,10.0.1.12:9092,10.0.1.13:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-processor");
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);
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());
// 记录失败位移,不中断消费
failedOffsets.add(new TopicPartition(record.topic(), record.partition()), record.offset());
}
}
// 批量处理完成后提交位移
consumer.commitSync();
}
// 精确位移提交
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
offsets.put(new TopicPartition("order-events", 0), new OffsetAndMetadata(lastOffset + 1));
consumer.commitSync(offsets);
enable.auto.commit=false关闭自动提交,改为处理完成后手动提交。max.poll.interval.ms=300000限制两次poll最大间隔5分钟,处理超时消费者会被踢出组触发Rebalance。max.poll.records=500限制单次poll返回的最大消息数,控制批量大小。
手动提交位移时,提交的位移值是下一条待消费消息的位移,即lastOffset + 1。提交失败位移之前的位移会导致已处理但未提交的消息在Rebalance后重复消费,因此消费逻辑需要实现幂等性。
精确一次语义实现
Kafka 0.11+通过事务API实现Exactly-Once语义。典型场景:消费Topic-A处理后写入Topic-B,位移提交和消息生产在同一事务中:
// 生产者开启事务
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-tx-1");
KafkaProducer<String, String> txProducer = new KafkaProducer<>(props);
txProducer.initTransactions();
// 消费者读取已提交消息
props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
txProducer.beginTransaction();
try {
for (ConsumerRecord<String, String> record : records) {
String result = processOrder(record.value());
txProducer.send(new ProducerRecord<>("order-results", record.key(), result));
}
// 提交消费位移到事务中
txProducer.sendOffsetsToTransaction(
getOffsetsToCommit(records), consumer.groupMetadata()
);
txProducer.commitTransaction();
} catch (Exception e) {
txProducer.abortTransaction();
log.error("事务回滚: {}", e.getMessage());
}
}
transactional.id必须全局唯一且持久化,Kafka通过它识别跨重启的生产者实例,中止未完成的事务。isolation.level=read_committed使消费者只读取已提交的消息,看不到事务中正在写入的消息。
事务的性能开销主要来自额外的broker交互(InitProducerId、BeginTxn、CommitTxn),单条消息的事务吞吐约为非事务的70%。批量事务(一次事务处理多条消息)可摊薄开销,适合低频高可靠性场景。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/kafka-xiao-xi-dui-lie-gao-ke-yong-bu-shu-shi-zhan-fen-qu-fu/