消息队列 Kafka 实战:高吞吐场景下的生产者与消费者最佳实践
作者:洛水石
阅读时间:约 14 分钟
关键词:Kafka、消息队列、生产者、消费者、高吞吐、Exactly-Once、分区、副本
---
目录
- [引言:为什么选 Kafka?](#引言)
- [Kafka 核心概念速览](#kafka-核心概念速览)
- [生产者最佳实践](#生产者最佳实践)
- [消费者最佳实践](#消费者最佳实践)
- [Exactly-Once 语义实现](#exactly-once-语义实现)
- [高吞吐调优参数](#高吞吐调优参数)
- [监控与告警](#监控与告警)
- [总结](#总结)
---
引言:为什么选 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实战 —
更多推荐

所有评论(0)