消息中间件是微服务架构解耦的核心组件,Kafka凭借高吞吐、低延迟、可持久化的特性成为分布式系统异步通信的首选方案。本文以Spring Boot集成Kafka为主线,覆盖Broker集群部署、生产者/消费者配置、精确一次语义、消息积压治理等生产环境关键问题。
Kafka集群部署与Topic规划
生产环境Kafka集群至少3个Broker节点,搭配ZooKeeper或KRaft模式(Kafka 3.3+内置共识协议无需ZooKeeper)。以KRaft模式部署为例,配置文件核心参数:
# server.properties (KRaft模式)
process.roles=broker,controller
node.id=1
controller.quorum.voters=1@broker1:9093,2@broker2:9093,3@broker3:9093
listeners=PLAINTEXT://:9092,CONTROLLER://:9093
advertised.listeners=PLAINTEXT://broker1: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
message.max.bytes=10485760
Topic规划遵循业务域隔离原则,命名规范为业务域.事件类型.版本,如order.created.v1、payment.result.v1。分区数根据消费组并发度设定,公式:分区数 = 目标TPS / 单分区TPS。副本数3,min.insync.replicas=2确保写入可靠性。创建Topic命令:
kafka-topics.sh --create \
--bootstrap-server broker1:9092 \
--topic order.created.v1 \
--partitions 12 \
--replication-factor 3 \
--config min.insync.replicas=2 \
--config retention.ms=604800000
Spring Boot生产者配置与消息发送
Spring Boot通过spring-kafka集成Kafka客户端。生产者配置需关注ack机制、重试策略、序列化和事务。acks=all确保消息写入所有ISR副本后才确认,配合retries和幂等生产者实现Exactly-Once语义:
# application.yml 生产者配置
spring:
kafka:
bootstrap-servers: broker1:9092,broker2:9092,broker3:9092
producer:
acks: all
retries: 3
batch-size: 16384
buffer-memory: 33554432
compression-type: lz4
properties:
enable.idempotence: true
max.in.flight.requests.per.connection: 5
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
// Java生产者服务
@Service
public class OrderEventProducer {
private final KafkaTemplate<String, OrderEvent> kafkaTemplate;
public OrderEventProducer(KafkaTemplate<String, OrderEvent> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void sendOrderCreated(OrderEvent event) {
// 使用orderId作为key保证同一订单事件有序
kafkaTemplate.send("order.created.v1", event.getOrderId(), event)
.addCallback(
result -> log.info("Message sent to partition {}",
result.getRecordMetadata().partition()),
ex -> log.error("Failed to send message", ex)
);
}
// 事务消息:跨Topic原子写入
@Transactional
public void processOrder(OrderEvent event) {
kafkaTemplate.executeInTransaction(template -> {
template.send("order.created.v1", event.getOrderId(), event);
template.send("inventory.deduct.v1", event.getProductId(),
new InventoryEvent(event.getProductId(), event.getQuantity()));
return true;
});
}
}
消费者配置与并发消费
消费者配置重点在于消费组管理、位移提交策略和并发控制。手动提交位移避免自动提交导致的消息丢失,消费失败时通过死信队列(DLT)隔离异常消息:
# application.yml 消费者配置
spring:
kafka:
consumer:
group-id: order-service
auto-offset-reset: earliest
enable-auto-commit: false
max-poll-records: 500
properties:
max.poll.interval.ms: 300000
session.timeout.ms: 30000
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "com.example.events"
listener:
ack-mode: manual_immediate
concurrency: 6
// Java消费者服务
@Component
public class OrderEventConsumer {
@KafkaListener(topics = "order.created.v1",
groupId = "order-service",
containerFactory = "kafkaListenerContainerFactory")
public void handleOrderCreated(
@Payload OrderEvent event,
@Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
@Header(KafkaHeaders.OFFSET) long offset,
Acknowledgment ack) {
try {
log.info("Processing order: {}, partition: {}, offset: {}",
event.getOrderId(), partition, offset);
orderService.processOrder(event);
ack.acknowledge();
} catch (Exception e) {
log.error("Failed to process order: {}", event.getOrderId(), e);
// 重试3次后发送到死信队列
if (event.getRetryCount() >= 3) {
kafkaTemplate.send("order.created.v1.DLT", event.getOrderId(), event);
ack.acknowledge();
} else {
event.incrementRetryCount();
throw e; // 触发重试
}
}
}
}
消息积压诊断与性能调优
消息积压是Kafka运维中最常见的问题。诊断流程:通过kafka-consumer-groups.sh查看消费组Lag,定位积压分区。常见原因及解决方案:
消费者处理速度不足时,增加concurrency配置提升消费线程数,但不超过分区数(超出部分的消费者空闲)。批量消费模式将max-poll-records调大至1000-2000,配合批量处理逻辑减少网络往返。消费逻辑中存在慢查询时,引入Redis缓存热点数据、优化SQL索引、异步化非核心逻辑。
# 消费组Lag监控
kafka-consumer-groups.sh \
--bootstrap-server broker1:9092 \
--describe \
--group order-service
# 输出示例:
# TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
# order.created.v1 0 15234 15300 66
# order.created.v1 1 14890 14890 0
# 扩容分区(注意:只能增加不能减少)
kafka-topics.sh --alter --bootstrap-server broker1:9092 \
--topic order.created.v1 --partitions 24
服务治理层面,通过Micrometer + Prometheus监控消费Lag、处理耗时、错误率三个核心指标。Lag持续增长触发自动扩容消费者实例(K8s HPA基于自定义指标)。消费者Rebalance导致消费暂停时,使用Cooperative Sticky Assignor替代默认的RangeAssignor减少Rebalance影响范围。Broker端优化包括增大socket.buffer.bytes提升网络吞吐、调整num.io.threads和num.network.threads匹配硬件资源、使用SSD存储日志段降低IO延迟。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/kafka-gao-tun-tu-xiao-xi-dui-lie-shi-zhan-springboot-ji/