作者:洛水石

阅读时间:约 14 分钟

关键词:Kafka、消息队列、生产者、消费者、高吞吐、Exactly-Once、分区、副本

---

目录

  1. [引言:为什么选 Kafka?](#引言)
  2. [Kafka 核心概念速览](#kafka-核心概念速览)
  3. [生产者最佳实践](#生产者最佳实践)
  4. [消费者最佳实践](#消费者最佳实践)
  5. [Exactly-Once 语义实现](#exactly-once-语义实现)
  6. [高吞吐调优参数](#高吞吐调优参数)
  7. [监控与告警](#监控与告警)
  8. [总结](#总结)

---

引言:为什么选 Kafka?

在微服务架构中,服务间通信方式主要有两种:同步调用(HTTP/RPC)和异步消息(MQ)。

当面临以下场景时,Kafka 是不二之选:

  • **日志采集**:日均百亿级日志,需要高吞吐写入
  • **事件驱动**:订单状态变更触发库存、物流、通知等多个服务
  • **流式处理**:实时数据分析、风控规则引擎
  • **削峰填谷**:秒杀场景下平滑流量峰值

Kafka 单机吞吐可达 **百万级 TPS**,这是 RabbitMQ、RocketMQ 难以企及的。

---

Kafka 核心概念速览

核心架构

Producer ──→  Kafka Cluster  ──→ Consumer

                │

    ┌───────────┼───────────┐

    │           │           │

  Broker 1   Broker 2   Broker 3

    │           │           │

  ┌─┴─┐      ┌─┴─┐      ┌─┴─┐

  P0  P1     P2  P3     P4  P5    ← Partition

  │   │      │   │      │   │

  L   F      L   F      L   F     ← Leader / Follower

关键术语

术语

说明

类比

Topic

消息主题,逻辑上的消息队列

数据库表

Partition

主题的分片,实现并行处理

表的分区

Offset

消息在分区中的唯一标识

自增 ID

Broker

Kafka 服务器节点

MySQL 实例

Consumer Group

消费者组,组内消费者共同消费

负载均衡集群

Replication

副本机制,保证高可用

主从复制

ISR

同步中的副本集合

健康从库

Topic

消息主题,逻辑上的消息队列

数据库表

Partition

主题的分片,实现并行处理

表的分区

Offset

消息在分区中的唯一标识

自增 ID

Broker

Kafka 服务器节点

MySQL 实例

Consumer Group

消费者组,组内消费者共同消费

负载均衡集群

Replication

副本机制,保证高可用

主从复制

ISR

同步中的副本集合

健康从库

Partition

主题的分片,实现并行处理

表的分区

Offset

消息在分区中的唯一标识

自增 ID

Broker

Kafka 服务器节点

MySQL 实例

Consumer Group

消费者组,组内消费者共同消费

负载均衡集群

Replication

副本机制,保证高可用

主从复制

ISR

同步中的副本集合

健康从库

Offset

消息在分区中的唯一标识

自增 ID

Broker

Kafka 服务器节点

MySQL 实例

Consumer Group

消费者组,组内消费者共同消费

负载均衡集群

Replication

副本机制,保证高可用

主从复制

ISR

同步中的副本集合

健康从库

Broker

Kafka 服务器节点

MySQL 实例

Consumer Group

消费者组,组内消费者共同消费

负载均衡集群

Replication

副本机制,保证高可用

主从复制

ISR

同步中的副本集合

健康从库

Consumer Group

消费者组,组内消费者共同消费

负载均衡集群

Replication

副本机制,保证高可用

主从复制

ISR

同步中的副本集合

健康从库

Replication

副本机制,保证高可用

主从复制

ISR

同步中的副本集合

健康从库

---

生产者最佳实践

1. 配置优化

Properties props = new Properties();

props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092");

props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");

props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

// ===== 高吞吐配置 =====

// 批量发送大小,默认 16KB,生产建议 32-128KB

props.put("batch.size", 32768);

// 批量发送等待时间,默认 0,建议 5-10ms

props.put("linger.ms", 10);

// 压缩算法,减少网络传输

props.put("compression.type", "lz4");  // lz4 / snappy / gzip

// 发送缓冲区大小

props.put("buffer.memory", 67108864);  // 64MB

// ===== 可靠性配置 =====

// acks=all 确保所有 ISR 副本都写入

props.put("acks", "all");

// 失败重试次数

props.put("retries", 3);

// 开启幂等性(Kafka 2.5+)

props.put("enable.idempotence", "true");

KafkaProducer<String, String> producer = new KafkaProducer<>(props);

2. 异步发送 + 回调

// 异步发送(推荐)

ProducerRecord<String, String> record = new ProducerRecord<>(

    "order-topic",

    order.getOrderId(),  // key,用于分区路由

    JSON.toJSONString(order)

);

producer.send(record, (metadata, exception) -> {

    if (exception != null) {

        // 记录失败日志,后续补偿

        log.error("Send failed, topic={}, partition={}, offset={}",

            metadata.topic(), metadata.partition(), metadata.offset(), exception);

        // 写入死信队列或数据库,定时重试

        deadLetterService.save(record, exception.getMessage());

    } else {

        log.info("Send success, topic={}, partition={}, offset={}",

            metadata.topic(), metadata.partition(), metadata.offset());

    }

});

3. 自定义分区策略

public class OrderPartitioner implements Partitioner {

    @Override

    public int partition(String topic, Object key, byte[] keyBytes,

                         Object value, byte[] valueBytes, Cluster cluster) {

        List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);

        int numPartitions = partitions.size();

        // 根据订单 ID 取模,保证同一订单消息进入同一分区(顺序性)

        String orderId = (String) key;

        int partition = Math.abs(orderId.hashCode()) % numPartitions;

        return partition;

    }

    @Override

    public void configure(Map<String, ?> configs) {}

    @Override

    public void close() {}

}

// 使用自定义分区器

props.put("partitioner.class", "com.example.kafka.OrderPartitioner");

4. 生产者拦截器(统一加 TraceId)

public class TraceInterceptor implements ProducerInterceptor<String, String> {

    @Override

    public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {

        String traceId = MDC.get("traceId");

        if (traceId != null) {

            // 在消息头中注入 traceId

            record.headers().add("traceId", traceId.getBytes(StandardCharsets.UTF_8));

        }

        return record;

    }

    @Override

    public void onAcknowledgement(RecordMetadata metadata, Exception exception) {}

    @Override

    public void close() {}

    @Override

    public void configure(Map<String, ?> configs) {}

}

// 注册拦截器

props.put("interceptor.classes", "com.example.kafka.TraceInterceptor");

---

消费者最佳实践

1. 消费者配置

Properties props = new Properties();

props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092");

props.put("group.id", "order-consumer-group");

props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");

props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");

// 手动提交偏移量(推荐,保证消息不丢失)

props.put("enable.auto.commit", "false");

// 消费起始位置:earliest / latest / none

props.put("auto.offset.reset", "earliest");

// 单次拉取最大数量

props.put("max.poll.records", 500);

// 两次 poll 之间的最大间隔,超时则触发再均衡

props.put("max.poll.interval.ms", 300000);  // 5 分钟

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);

consumer.subscribe(Arrays.asList("order-topic"));

2. 手动提交偏移量

while (true) {

    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));

    for (ConsumerRecord<String, String> record : records) {

        try {

            // 处理消息

            processMessage(record);

            // 逐条同步提交(最可靠,但性能较低)

            // consumer.commitSync();

        } catch (Exception e) {

            log.error("Process failed, offset={}, error={}", record.offset(), e.getMessage());

            // 不提交偏移量,下次重新消费

            break;

        }

    }

    // 批量异步提交(推荐)

    consumer.commitAsync((offsets, exception) -> {

        if (exception != null) {

            log.error("Commit failed: {}", exception.getMessage());

        }

    });

}

3. 多线程消费模型

public class MultiThreadConsumer {

    private final KafkaConsumer<String, String> consumer;

    private final ExecutorService executor;

    public void start() {

        consumer.subscribe(Arrays.asList("order-topic"));

        while (true) {

            ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));

            // 按分区拆分任务

            for (TopicPartition partition : records.partitions()) {

                List<ConsumerRecord<String, String>> partitionRecords = records.records(partition);

                executor.submit(() -> {

                    for (ConsumerRecord<String, String> record : partitionRecords) {

                        processMessage(record);

                    }

                    // 处理完该分区所有消息后提交偏移量

                    synchronized (consumer) {

                        consumer.commitSync();

                    }

                });

            }

        }

    }

}

4. 消息处理幂等设计

@Service

public class OrderConsumer {

    @Autowired

    private RedisTemplate<String, String> redisTemplate;

    @KafkaListener(topics = "order-topic", groupId = "order-group")

    public void consume(ConsumerRecord<String, String> record) {

        String messageId = getMessageId(record);

        // 1. 幂等校验:Redis setnx

        Boolean isNew = redisTemplate.opsForValue()

            .setIfAbsent("msg:" + messageId, "1", Duration.ofHours(24));

        if (Boolean.FALSE.equals(isNew)) {

            log.warn("Duplicate message ignored, id={}", messageId);

            return;

        }

        // 2. 业务处理

        processOrder(record.value());

    }

    private String getMessageId(ConsumerRecord<String, String> record) {

        // 从 header 中获取 traceId 或生成唯一 ID

        Header header = record.headers().lastHeader("traceId");

        if (header != null) {

            return new String(header.value(), StandardCharsets.UTF_8);

        }

        return record.topic() + "-" + record.partition() + "-" + record.offset();

    }

}

---

Exactly-Once 语义实现

Kafka 从 0.11 版本开始支持 Exactly-Once 语义,核心依赖三个机制:

1. 生产者幂等性

props.put("enable.idempotence", "true");  // 自动开启

props.put("acks", "all");

props.put("retries", Integer.MAX_VALUE);

props.put("max.in.flight.requests.per.connection", 5);  // Kafka 2.5+ 支持 5

2. 事务生产者

Properties props = new Properties();

props.put("transactional.id", "order-producer-1");  // 唯一事务 ID

KafkaProducer<String, String> producer = new KafkaProducer<>(props);

producer.initTransactions();

try {

    producer.beginTransaction();

    // 发送多条消息

    producer.send(new ProducerRecord<>("order-topic", order1));

    producer.send(new ProducerRecord<>("payment-topic", payment1));

    // 发送偏移量到消费者组(consume-transform-produce 模式)

    producer.sendOffsetsToTransaction(

        consumer.position(consumer.assignment()),

        consumer.groupMetadata()

    );

    producer.commitTransaction();

} catch (Exception e) {

    producer.abortTransaction();

    throw e;

}

3. 消费者隔离级别

// 只读取已提交的事务消息

props.put("isolation.level", "read_committed");

// 读取所有消息(包括未提交的)

props.put("isolation.level", "read_uncommitted");  // 默认

---

高吞吐调优参数

Broker 端调优

server.properties

网络线程数

num.network.threads=8

IO 线程数

num.io.threads=16

发送缓冲区

socket.send.buffer.bytes=102400

接收缓冲区

socket.receive.buffer.bytes=102400

日志刷盘策略(权衡吞吐与可靠性)

log.flush.interval.messages=10000

log.flush.interval.ms=1000

副本同步策略

replica.fetch.max.bytes=1048576

replica.socket.timeout.ms=30000

Topic 级别调优

创建高吞吐 Topic

kafka-topics.sh --create \

  --topic high-throughput-topic \

  --partitions 12 \

  --replication-factor 3 \

  --config min.insync.replicas=2 \

  --config compression.type=lz4 \

  --config segment.bytes=536870912

JVM 调优

export KAFKA_HEAP_OPTS="-Xmx6g -Xms6g"

export KAFKA_JVM_PERFORMANCE_OPTS="-server -XX:+UseG1GC -XX:MaxGCPauseMillis=20"

---

监控与告警

核心监控指标

指标

说明

告警阈值

UnderReplicatedPartitions

副本不同步的分区数

> 0

OfflinePartitions

离线分区数

> 0

ActiveControllerCount

活跃 Controller 数

!= 1

ConsumerLag

消费者延迟

> 10000

RequestLatency

请求延迟

P99 > 100ms

ISRShrinkExpandRate

ISR 收缩/扩展速率

频繁变动需关注

UnderReplicatedPartitions

副本不同步的分区数

> 0

OfflinePartitions

离线分区数

> 0

ActiveControllerCount

活跃 Controller 数

!= 1

ConsumerLag

消费者延迟

> 10000

RequestLatency

请求延迟

P99 > 100ms

ISRShrinkExpandRate

ISR 收缩/扩展速率

频繁变动需关注

OfflinePartitions

离线分区数

> 0

ActiveControllerCount

活跃 Controller 数

!= 1

ConsumerLag

消费者延迟

> 10000

RequestLatency

请求延迟

P99 > 100ms

ISRShrinkExpandRate

ISR 收缩/扩展速率

频繁变动需关注

ActiveControllerCount

活跃 Controller 数

!= 1

ConsumerLag

消费者延迟

> 10000

RequestLatency

请求延迟

P99 > 100ms

ISRShrinkExpandRate

ISR 收缩/扩展速率

频繁变动需关注

ConsumerLag

消费者延迟

> 10000

RequestLatency

请求延迟

P99 > 100ms

ISRShrinkExpandRate

ISR 收缩/扩展速率

频繁变动需关注

RequestLatency

请求延迟

P99 > 100ms

ISRShrinkExpandRate

ISR 收缩/扩展速率

频繁变动需关注

Prometheus + Grafana 监控

kafka-metrics.yaml

  • name: kafka_server_replica_manager_under_replicated_partitions

  help: Number of under-replicated partitions

  type: GAUGE

  • name: kafka_consumer_group_lag

  help: Consumer group lag by topic partition

  type: GAUGE

  labels:

    - group

    - topic

    - partition

消费者延迟告警规则

alert-rules.yaml

groups:

  - name: kafka-alerts

    rules:

      - alert: KafkaConsumerLagHigh

        expr: kafka_consumer_group_lag > 10000

        for: 5m

        labels:

          severity: warning

        annotations:

          summary: "Kafka consumer lag is high"

          description: "Consumer group {{ $labels.group }} lag is {{ $value }}"

---

总结

Kafka 使用 checklist

生产者:

  • [ ] 开启幂等性 `enable.idempotence=true`
  • [ ] 使用异步发送 + 回调处理失败
  • [ ] 合理设置 `batch.size` 和 `linger.ms`
  • [ ] 启用压缩 `compression.type=lz4`
  • [ ] 自定义分区器保证消息顺序

消费者:

  • [ ] 关闭自动提交,使用手动提交
  • [ ] 实现消息幂等(Redis setnx 或数据库唯一索引)
  • [ ] 控制 `max.poll.records` 避免处理超时
  • [ ] 使用多线程按分区并行消费

运维:

  • [ ] 监控 UnderReplicatedPartitions 和 ConsumerLag
  • [ ] 配置 `min.insync.replicas >= 2`
  • [ ] 定期做分区再均衡(reassignment)
  • [ ] JVM 使用 G1GC,避免 Full GC 停顿

一句话总结

**Kafka 的高吞吐来自于批处理、压缩和分区并行;高可靠来自于副本、ISR 和 acks=all;Exactly-Once 来自于幂等 + 事务 + read_committed。**

---

更多硬核技术文章每周更新。

配图1:Kafka核心架构与组件关系

配图2:生产者与消费者最佳实践

配图3:消费者延迟监控与告警

— 作者:洛水石 | 消息队列 | Kafka实战 —

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐