5分钟快速上手:JDspyder京东抢购脚本终极实战指南
2026/7/22 0:48:05
视频看了几百小时还迷糊?关注我,几分钟让你秒懂!
在 Kafka 中,消费者消费消息后需要“提交偏移量”(offset commit),告诉 Kafka:“我已经处理到第 X 条消息了,下次从 X+1 开始给我发”。
Kafka 提供两种提交方式:
今天我们重点讲自动提交机制—— 它看似简单,但用不好会丢消息或重复消费!
spring: kafka: consumer: enable-auto-commit: true # 默认就是 true! auto-commit-interval: 5000 # 每 5 秒提交一次auto.commit.interval.ms(默认 5 秒)自动提交当前已拉取的最新 offset⚠️ 关键问题:自动提交的是“已拉取”的 offset,不是“已处理”的 offset!
你希望:只有邮件真正发送成功,才算消息消费成功。
@KafkaListener(topics = "user-register") public void handleRegister(User user) { // 1. 拉取消息(offset=101 被拉取) // 2. 自动提交线程将在 5 秒内提交 offset=101 try { emailService.sendWelcomeEmail(user); // 耗时 8 秒,且可能失败 } catch (Exception e) { // 邮件发送失败!但 offset 已经被自动提交了! log.error("邮件发送失败", e); } }| 时间 | 事件 |
|---|---|
| T=0s | 拉取 offset=101 的消息 |
| T=3s | 开始发送邮件(耗时 8 秒) |
| T=5s | 自动提交线程提交 offset=101(即使邮件还没发完!) |
| T=6s | 应用崩溃 / 重启 |
| T=7s | 重启后从 offset=102 开始消费 →offset=101 的消息永远丢失! |
🔥 这就是消息丢失的典型原因!
application.yml)spring: kafka: consumer: enable-auto-commit: false # 关闭自动提交! group-id: email-service@KafkaListener( topics = "user-register", groupId = "email-service" ) public void handleRegister(User user, Acknowledgment ack) { try { emailService.sendWelcomeEmail(user); // 业务成功 → 手动提交 offset ack.acknowledge(); } catch (Exception e) { // 失败时不提交,下次重启还会重新消费这条消息 log.error("处理失败,不提交 offset", e); // 可选:记录到死信队列,避免无限重试 } }@Configuration @EnableKafka public class KafkaConfig { @Bean public ConcurrentKafkaListenerContainerFactory<String, User> kafkaListenerContainerFactory( ConsumerFactory<String, User> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, User> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 设置为手动提交 factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); return factory; } }✅ 这样:只有
ack.acknowledge()被调用,offset 才会提交,确保“处理成功才提交”。
| 特性 | 自动提交(auto-commit) | 手动提交(manual commit) |
|---|---|---|
| 默认开启 | ✅ 是 | ❌ 否 |
| 提交时机 | 每隔 N 秒(与业务无关) | 由代码控制(业务成功后) |
| 消息可靠性 | 可能丢消息 | 不丢消息(至少一次) |
| 重复消费 | 不会(但可能丢) | 可能重复(需幂等) |
| 适用场景 | 日志收集、监控等允许丢失的场景 | 订单、支付、邮件等关键业务 |
ack.acknowledge()之后、实际业务完成前崩溃(极小概率),所以消费者逻辑要幂等(如用数据库唯一索引去重)。AckMode.BATCH,一批消息处理完再统一提交。kill -9模拟非优雅关闭,验证消息是否丢失。视频看了几百小时还迷糊?关注我,几分钟让你秒懂!