Kafka消息队列高可用部署实战:分区副本机制与生产者消费者配置

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/

(0)
小编小编
上一篇 9小时前
下一篇 9小时前

相关推荐