Spring Boot + Kafka 实战:生产者消费者完整示例
Spring Boot 集成 Spring Kafka 的完整示例,包含 HTTP 触发消息发送、KafkaTemplate 批量生产和 @KafkaListener 消费。
Spring Boot + Kafka 实战:生产与消费
Spring Kafka 集成
Spring Kafka 在原生 Kafka Client API 上做了封装,让生产者和消费者的使用体验大幅简化。
KafkaTemplate 生产者
@Service
public class OrderProducer {
@Autowired
private KafkaTemplate<String, OrderMessage> kafkaTemplate;
public void sendOrderCreated(Order order) {
OrderMessage msg = OrderMessage.builder()
.orderId(order.getId())
.userId(order.getUserId())
.amount(order.getTotalAmount())
.timestamp(System.currentTimeMillis())
.build();
// 发送消息
ListenableFuture<SendResult<String, OrderMessage>> future =
kafkaTemplate.send("order-events", msg.getOrderId(), msg);
// 异步处理结果
future.addCallback(
result -> log.info("订单消息发送成功: {}", msg.getOrderId()),
ex -> log.error("订单消息发送失败", ex)
);
}
}
@KafkaListener 消费者
@Component
public class OrderConsumer {
@KafkaListener(
topics = "order-events",
groupId = "order-processing-group",
containerFactory = "batchFactory"
)
public void consume(List<ConsumerRecord<String, OrderMessage>> records,
Acknowledgment ack) {
try {
for (ConsumerRecord<String, OrderMessage> record : records) {
processOrder(record.value());
}
// 处理完成后手动提交 offset
ack.acknowledge();
} catch (Exception e) {
log.error("消费订单消息失败", e);
// 不提交 offset,下次重新消费
}
}
}
序列化与反序列化
Kafka 传输的是字节数组,Java 对象需要序列化:
public class OrderMessageSerializer implements Serializer<OrderMessage> {
@Override
public byte[] serialize(String topic, OrderMessage data) {
return JsonUtils.toJson(data).getBytes(StandardCharsets.UTF_8);
}
}
public class OrderMessageDeserializer implements Deserializer<OrderMessage> {
@Override
public OrderMessage deserialize(String topic, byte[] data) {
return JsonUtils.fromJson(new String(data, StandardCharsets.UTF_8),
OrderMessage.class);
}
}
对应的生产者配置:
spring:
kafka:
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: com.example.config.OrderMessageSerializer
consumer:
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: com.example.config.OrderMessageDeserializer
properties:
spring.json.trusted.packages: "*"
幂等消费:唯一业务 ID
Kafka 是 at-least-once 语义,意味着消费者可能收到重复消息。消费端需要幂等:
@Service
public class OrderProcessor {
@Autowired
private RedisTemplate<String, String> redisTemplate;
public void processOrder(OrderMessage msg) {
// 先检查是否处理过
String processedKey = "processed:order:" + msg.getOrderId();
Boolean existed = redisTemplate.opsForValue().setIfAbsent(
processedKey, "1", Duration.ofHours(24));
if (Boolean.FALSE.equals(existed)) {
log.info("订单 {} 已处理过,跳过", msg.getOrderId());
return; // 幂等
}
// 核心业务逻辑
orderService.processOrder(msg);
}
}
异常处理与死信队列
@Bean
public ConcurrentKafkaListenerContainerFactory<String, OrderMessage> factory() {
ConcurrentKafkaListenerContainerFactory<String, OrderMessage> factory =
new ConcurrentKafkaListenerContainerFactory<>();
// 重试配置:重试 3 次,间隔 2 秒
factory.setRetryTemplate(new RetryTemplate() {{
setRetryOperations(new SimpleRetryPolicy(3));
setBackOffPolicy(new FixedBackOffPolicy() {{
setBackOffPeriod(2000);
}});
}});
// 重试耗尽后进入死信队列
factory.setErrorHandler((record, exception) -> {
kafkaTemplate.send("order-events-dlq", record.value());
log.error("消息处理失败,已转入死信队列", exception);
});
return factory;
}
消费者组与 offset
消费者组内每个消费者消费不同的 Partition。offset 提交方式:
- 自动提交(enable.auto.commit=true):定时提交,可能导致重复消费或丢消息
- 手动提交(enable.auto.commit=false):业务处理完成后手动 ack。推荐生产使用
spring:
kafka:
consumer:
enable-auto-commit: false
auto-offset-reset: earliest # 新消费者从最早的消息开始
版本兼容
Kafka 服务端版本和客户端版本不匹配可能导致连接失败。Spring Kafka 自动适配 Kafka 服务端版本,但如果用原生 Kafka Client API,需要确保客户端版本 ≤ 服务端版本。
消息顺序
Kafka 的消息顺序只在单 Partition 内保证。如果业务要求严格顺序(如订单状态流转),需要确保同一状态的订单发到同一个 Partition(用订单 ID 作为 Key)。