Spring Boot中@Scheduled定时任务如何在K8s集群仅单Pod执行(适配Firestore)
基于Firestore实现Spring Boot分布式定时任务单Pod执行方案
方案一:自行实现Firestore分布式锁
由于Shedlock暂不支持Firestore,最直接的解决方案是基于Firestore的事务特性手动实现分布式锁,确保同一时间只有一个Pod能获取锁并执行定时任务。
核心思路
利用Firestore的文档创建原子性和事务机制:
- 为每个定时任务定义唯一锁名称(比如
daily-report-job-lock) - 尝试创建包含过期时间、Pod标识的锁文档,只有第一个成功创建的Pod能获取锁
- 其他Pod尝试创建时会因文档已存在而失败,直接跳过任务
- 任务执行完成后主动释放锁;若Pod意外挂掉,锁会在过期时间后自动失效,不影响后续任务执行
具体实现步骤
1. 定义锁文档实体
import com.google.cloud.firestore.annotation.DocumentId; import java.time.Instant; public class DistributedLock { @DocumentId private String lockName; private String ownerPodId; private Instant expireTime; // 构造器、getter、setter省略 }
2. 实现分布式锁服务
import com.google.cloud.firestore.Firestore; import com.google.cloud.firestore.Transaction; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import java.time.Instant; import java.util.concurrent.TimeUnit; @Service public class FirestoreLockService { private final Firestore firestore; private final String podId; // 从K8s环境变量获取Pod唯一标识 public FirestoreLockService(Firestore firestore, @Value("${HOSTNAME}") String podId) { this.firestore = firestore; this.podId = podId; } /** * 尝试获取锁 * @param lockName 锁名称 * @param lockDuration 锁有效期(防止Pod挂死导致锁永久占用) * @return 是否获取成功 */ public boolean tryLock(String lockName, long lockDuration, TimeUnit timeUnit) { try { return firestore.runTransaction(transaction -> { var lockDocRef = firestore.collection("distributed-locks").document(lockName); var lockDoc = transaction.get(lockDocRef).get(); // 锁不存在或已过期,尝试获取锁 if (!lockDoc.exists() || lockDoc.toObject(DistributedLock.class).getExpireTime().isBefore(Instant.now())) { DistributedLock lock = new DistributedLock(); lock.setLockName(lockName); lock.setOwnerPodId(podId); lock.setExpireTime(Instant.now().plusMillis(timeUnit.toMillis(lockDuration))); transaction.set(lockDocRef, lock); return true; } return false; }).get(); } catch (Exception e) { // 异常情况下默认不获取锁,避免任务重复执行 return false; } } /** * 释放锁(仅锁持有者能释放) */ public void unlock(String lockName) { firestore.runTransaction(transaction -> { var lockDocRef = firestore.collection("distributed-locks").document(lockName); var lockDoc = transaction.get(lockDocRef).get(); if (lockDoc.exists()) { DistributedLock lock = lockDoc.toObject(DistributedLock.class); if (podId.equals(lock.getOwnerPodId())) { transaction.delete(lockDocRef); } } return null; }); } }
3. 改造定时任务
import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.util.concurrent.TimeUnit; @Component public class ScheduledTask { private final FirestoreLockService lockService; private static final String TASK_LOCK_NAME = "daily-report-job-lock"; public ScheduledTask(FirestoreLockService lockService) { this.lockService = lockService; } @Scheduled(cron = "0 0 1 * * ?") // 每天凌晨1点执行 public void executeDailyReport() { // 尝试获取锁,有效期设置为1小时(确保任务能执行完成) boolean lockAcquired = lockService.tryLock(TASK_LOCK_NAME, 1, TimeUnit.HOURS); if (!lockAcquired) { System.out.println("Pod " + System.getenv("HOSTNAME") + " 未获取到锁,跳过任务执行"); return; } try { // 执行核心任务逻辑 System.out.println("Pod " + System.getenv("HOSTNAME") + " 开始执行每日报表任务"); // ... 任务代码 } finally { // 无论任务成功还是失败,都释放锁 lockService.unlock(TASK_LOCK_NAME); } } }
方案二:结合Kubernetes CronJob(可选)
如果不想在代码层处理锁,也可以直接用Kubernetes的CronJob资源替代Spring的@Scheduled:
- 配置CronJob的
spec.jobTemplate.spec.parallelism=1和spec.concurrencyPolicy=Forbid,确保同一时间只有一个Job实例运行 - 这种方式不需要依赖Firestore,但需要调整部署架构,将定时任务从应用Pod中剥离出来
注意事项
- 锁有效期设置:必须大于任务的最长执行时间,防止任务未完成锁就过期,导致其他Pod重复执行
- 异常处理:任务执行过程中抛出异常时,一定要在
finally块中释放锁(或依赖锁的过期机制自动释放) - Pod ID获取:Kubernetes会自动为每个Pod设置
HOSTNAME环境变量,其值就是Pod的名称,可作为锁的持有者标识
内容的提问来源于stack exchange,提问作者Abhinandan Sharma
相关产品推荐
相关产品推荐

