尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Spring Boot 3 多实例定时任务防重复执行:基于 Kafka 的分布式锁与 JVM GC 优化实践

Spring Boot 3 多实例定时任务防重复执行:基于 Kafka 的分布式锁与 JVM GC 优化实践 Spring Boot 3 多实例定时任务防重复执行基于 Kafka 的分布式锁与 JVM GC 优化实践问题背景在电商、金融等业务场景中订单对账、供应商结算、用户积分发放等定时任务是核心流程这类任务重复执行会直接导致资损、重复触达用户等严重问题。传统基于Scheduled的定时任务在单实例部署时运行正常但一旦因K8s扩容、节点故障转移部署多实例同一任务会被多个实例同时执行引发生产事故。现有主流解决方案各有痛点基于数据库的ShedLock需要额外维护锁表高并发场景下DB压力大时锁获取延迟高甚至失败Redis分布式锁需要单独维护Redis集群存在锁续期、脑裂风险且无执行过程追溯能力ZooKeeper锁运维成本高不适合中小团队。若企业已部署Kafka集群能否基于Kafka实现低运维成本的分布式锁同时任务执行期间若发生JVM长时间GC停顿可能导致锁心跳发送延迟引发锁误释放又该如何规避本文围绕这两个问题提出一套基于Spring Boot 3、Kafka、JVM优化的完整方案。方案设计本方案三个技术栈的分工明确无强行拼接 1.Spring Boot 3作为调度层提供Scheduled定时任务入口整合分布式锁逻辑是任务调度的核心载体 2.Kafka作为分布式锁实现层利用其消费者组分区独占、消息持久化、高可用的特性实现锁的持有、心跳维持、过期判定同时全流程锁操作有消息日志可追溯 3.JVM内存与GC优化作为保障层通过选择低停顿GC算法、优化堆内存配置避免任务执行期间的长停顿STW导致锁心跳延迟引发锁误释放。整体执行流程每个实例启动后定时任务触发时先向Kafka锁主题发送抢锁请求仅拿到锁的实例执行任务执行期间定时发送心跳维持锁有效性任务完成后释放锁锁过期时间设置为任务最大执行时长的1.5倍避免任务超时导致锁提前释放。关键原理Kafka分布式锁核心原理锁主题设计创建专属锁主题scheduled-task-lock分区数固定为1保证全局仅有一个消费者能消费到该分区消息避免多实例同时持锁副本数设为3保证高可用避免Broker宕机导致锁消息丢失抢锁逻辑任务触发时生产者向锁主题发送带任务标识taskId执行时间戳的抢锁消息配置acksall保证消息持久化所有实例的消费者监听该主题仅能消费到抢锁消息的实例获得锁执行权锁维持与过期持锁实例每隔锁过期时间的1/3发送一次心跳消息带当前时间戳其他实例通过比对最近一次心跳的时间戳判定锁是否过期若超过锁过期时间未收到心跳则认为锁失效可重新抢锁释放锁任务执行完成后持锁实例发送释放锁消息或直接停止发送心跳等待锁自然过期。JVM优化核心原理任务执行期间若发生Full GC会导致JVM STWStop-The-World若STW时间超过心跳发送间隔持锁实例无法及时发送心跳其他实例会误认为锁失效抢锁执行导致任务重复。因此选择JDK 17正式可用的ZGC算法其停顿时间与堆大小无关最大停顿不超过1ms完全避免长停顿导致的心跳延迟问题。同时固定堆大小、关闭自适应调整避免堆动态变化触发Full GC。完整示例环境要求JDK 17/21ZGC生产可用最低版本Spring Boot 3.2.xKafka 3.0已部署集群spring-kafka 3.0.x与Spring Boot 3版本匹配1. 依赖引入dependencies !-- Spring Boot 3 核心依赖 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter/artifactId version3.2.5/version /dependency !-- Spring Kafka 依赖 -- dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId version3.0.15/version /dependency !-- 可选GC监控依赖 -- dependency groupIdio.micrometer/groupId artifactIdmicrometer-registry-prometheus/artifactId /dependency /dependencies2. 配置类实现2.1 应用配置application.ymlspring: kafka: bootstrap-servers: your-kafka-cluster:9092 producer: acks: all # 保证锁消息不丢 retries: 3 # 发送失败重试3次 linger-ms: 5 # 减少消息发送延迟 batch-size: 16384 consumer: group-id: scheduled-task-lock-group # 锁消费者组所有实例共用 auto-offset-reset: latest # 只消费最新消息 enable-auto-commit: false # 手动提交偏移量 scheduled: task: lock-topic: scheduled-task-lock # 锁主题名称 lock-expire-seconds: 1800 # 锁过期时间30分钟根据任务最大执行时长调整 heartbeat-interval-seconds: 600 # 心跳间隔10分钟小于锁过期时间的1/32.2 Kafka分布式锁实现import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.time.Instant; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicLong; Component public class KafkaDistributedLock { private final KafkaTemplateString, String kafkaTemplate; private final String lockTopic; private final long lockExpireMs; private final long heartbeatIntervalMs; // 存储持锁任务ID对应的最后心跳时间 private final MapString, Long heartbeatMap new ConcurrentHashMap(); // 存储持锁任务ID对应的实例ID private final MapString, String lockHolderMap new ConcurrentHashMap(); // 当前实例ID可配置为IP端口保证唯一 private final String instanceId instance- System.currentTimeMillis(); public KafkaDistributedLock(KafkaTemplateString, String kafkaTemplate, org.springframework.core.env.Environment env) { this.kafkaTemplate kafkaTemplate; this.lockTopic env.getProperty(scheduled.task.lock-topic, scheduled-task-lock); this.lockExpireMs Long.parseLong(env.getProperty(scheduled.task.lock-expire-seconds, 1800)) * 1000; this.heartbeatIntervalMs Long.parseLong(env.getProperty(scheduled.task.heartbeat-interval-seconds, 600)) * 1000; } /** * 抢锁方法 * param taskId 任务唯一标识 * return 是否抢锁成功 */ public boolean tryLock(String taskId) { String lockMessage String.format(LOCK|%s|%s|%d, taskId, instanceId, Instant.now().toEpochMilli()); // 发送抢锁消息key为taskId保证同一任务消息发到同一分区 kafkaTemplate.send(lockTopic, taskId, lockMessage); // 等待100ms判断是否拿到锁可根据网络情况调整 try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } // 判断当前实例是否是持锁者 return instanceId.equals(lockHolderMap.get(taskId)); } /** * 释放锁 * param taskId 任务唯一标识 */ public void releaseLock(String taskId) { String releaseMessage String.format(RELEASE|%s|%s|%d, taskId, instanceId, Instant.now().toEpochMilli()); kafkaTemplate.send(lockTopic, taskId, releaseMessage); lockHolderMap.remove(taskId); heartbeatMap.remove(taskId); } /** * 监听锁主题消息处理抢锁、心跳、释放逻辑 */ KafkaListener(topics ${scheduled.task.lock-topic}, groupId scheduled-task-lock-group) public void listenLockMessage(String message) { String[] parts message.split(\\|); if (parts.length 4) return; String type parts[0]; String taskId parts[1]; String senderInstance parts[2]; long timestamp Long.parseLong(parts[3]); switch (type) { case LOCK: // 抢锁消息如果当前没有持锁者或者持锁者心跳过期则更新持锁者 Long lastHeartbeat heartbeatMap.get(taskId); if (lastHeartbeat null || (Instant.now().toEpochMilli() - lastHeartbeat) lockExpireMs) { lockHolderMap.put(taskId, senderInstance); heartbeatMap.put(taskId, timestamp); } break; case HEARTBEAT: // 心跳消息更新对应任务的心跳时间 if (senderInstance.equals(lockHolderMap.get(taskId))) { heartbeatMap.put(taskId, timestamp); } break; case RELEASE: // 释放锁消息清除持锁信息 if (senderInstance.equals(lockHolderMap.get(taskId))) { lockHolderMap.remove(taskId); heartbeatMap.remove(taskId); } break; default: break; } } /** * 定时清理过期锁避免持锁实例宕机后锁永久持有 */ Scheduled(fixedRate 60000) // 每分钟检查一次 public void cleanExpiredLock() { long now Instant.now().toEpochMilli(); heartbeatMap.forEach((taskId, lastHeartbeat) - { if (now - lastHeartbeat lockExpireMs) { lockHolderMap.remove(taskId); heartbeatMap.remove(taskId); } }); } /** * 定时发送心跳维持锁有效性 */ Scheduled(fixedRate ${scheduled.task.heartbeat-interval-seconds:600}000) public void sendHeartbeat() { long now Instant.now().toEpochMilli(); lockHolderMap.forEach((taskId, instance) - { if (instance.equals(this.instanceId)) { String heartbeatMessage String.format(HEARTBEAT|%s|%s|%d, taskId, instanceId, now); kafkaTemplate.send(lockTopic, taskId, heartbeatMessage); } }); } }2.3 定时任务改造import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import jakarta.annotation.Resource; Component public class OrderReconcileTask { Resource private KafkaDistributedLock distributedLock; // 每日凌晨1点执行对账任务taskId为任务唯一标识 Scheduled(cron 0 0 1 * * ?) public void reconcileOrder() { String taskId order-reconcile-daily; // 尝试获取分布式锁 if (!distributedLock.tryLock(taskId)) { System.out.println(其他实例已持有锁当前实例跳过执行); return; } try { System.out.println(当前实例拿到锁开始执行对账任务); // 业务逻辑拉取订单、对账、生成账单 Thread.sleep(1000 * 60 * 10); // 模拟任务执行10分钟 System.out.println(对账任务执行完成); } catch (Exception e) { // 任务异常也要记录日志锁会自动过期 e.printStackTrace(); } finally { // 释放锁 distributedLock.releaseLock(taskId); } } }3. 锁主题创建脚本执行以下命令创建Kafka锁主题分区数必须为1kafka-topics.sh --create --bootstrap-server your-kafka-cluster:9092 \ --topic scheduled-task-lock \ --partitions 1 \ --replication-factor 34. JVM启动参数配置-Xms4g -Xmx4g \ -XX:UseZGC -XX:UnlockExperimentalVMOptions \ -XX:DisableExplicitGC \ -Xlog:gc*:file./logs/gc.log:time,tid,tags:filecount10,filesize100M参数说明 --Xms/-Xmx固定堆大小为4G不超过32G以保证ZGC指针压缩生效可根据任务内存需求调整 --XX:UseZGC启用ZGC低停顿GC算法 --XX:DisableExplicitGC禁止显式GC调用避免第三方库触发Full GC --Xlog:gc*开启GC日志便于排查GC问题。常见问题1. 锁主题分区数为什么必须为1Kafka的消费者组消费的最小单位是分区若分区数大于1多个消费者可以同时消费不同分区的消息导致多个实例同时拿到锁重复执行任务。2. ZGC有没有生产环境兼容性问题JDK 17及以上的ZGC已经正式生产可用仅JDK 11的ZGC为实验性版本存在未修复的Bug不推荐生产使用。若服务器内存紧张低于8G可选择G1GC同时将心跳间隔调整为1分钟避免GC停顿超过心跳间隔。3. 如果Kafka集群宕机怎么办可配置降级逻辑Kafka连接失败时自动切换为本地单实例执行同时发送告警通知运维保证任务至少执行一次避免业务中断。4. 任务执行超时导致锁提前释放怎么办锁过期时间必须设置为任务最大执行时长的1.5倍以上比如任务最大执行30分钟锁过期时间设为45分钟避免任务未执行完锁就释放。5. 心跳消息发送失败怎么办已配置Kafka生产者重试3次若仍然失败则主动释放锁避免锁被无效持有同时发送告警通知运维检查网络。适用边界与关键取舍适用场景本方案适合已有Kafka集群、任务执行频率低每日/每周/每月、执行时间长分钟级到小时级、重复执行资损风险高的场景比如对账、结算、数据同步、报表生成。不适用场景任务执行频率高每分钟及以上、执行时间短秒级的场景Kafka锁的网络开销过大不如ShedLock或Redis锁性价比高无Kafka集群的业务额外搭建Kafka集群的成本远高于收益不推荐使用。关键取舍锁可靠性 vs 开销配置acksall、副本数3保证锁消息不丢会略微增加抢锁延迟但对于低频任务完全可接受锁过期时间 vs 任务执行时长锁过期时间设得过短会导致任务未执行完锁释放设得过长会导致实例宕机后任务延迟执行需根据业务SLA权衡GC算法选择ZGC停顿时间短但内存占用比G1高10%左右内存充足的场景优先选择内存紧张时可选择G1GC并调小心跳间隔。容易踩坑的细节必须配置-XX:DisableExplicitGC禁止显式GC调用很多第三方缓存库会默认调用System.gc()触发Full GC导致长时间STWKafka生产者的linger.ms不要配置过大否则抢锁消息会批量延迟发送推荐设为5ms锁心跳发送需在独立线程中执行不要和任务执行线程共用避免任务业务逻辑阻塞导致心跳发送延迟Spring Boot 3的Scheduled默认使用单线程池若同时执行多个定时任务需配置TaskExecutor开启异步执行避免任务互相阻塞。总结本文提出的基于Kafka的分布式锁方案复用企业已有的Kafka集群实现了多实例定时任务的唯一执行控制同时通过JVM GC优化避免了锁误释放的问题相比传统DB锁、Redis锁方案具备运维成本低、全流程可追溯的优势适合中大型团队的定时任务治理场景。若团队无Kafka集群可根据实际场景选择ShedLock或Redis锁方案核心思路都是保证分布式场景下定时任务的幂等性。
返回列表