Quarkus环境下基于SourceId的多线程作业并行实现及故障恢复咨询
问题分析与解决方案
场景回顾
接收带SourceId的作业执行消息,核心约束:
- 同一
SourceId的作业同一时间仅允许一个运行,其余进入队列等待 - 作业拆分为多个小执行单元(Batch),每个Batch完成后更新作业进度
- 作业全部完成/出错时更新作业状态(完成/失败)
- 支持服务器崩溃重启后作业恢复,维持
SourceId级并发控制
现有代码方案的遗漏点
- 并发竞态问题:
PrimaryWorker中staging.checkIfEntryExistForSource()与createEntryForDataSource()是非原子操作,多线程提交同一SourceId作业时,可能重复创建队列和SecondaryWorker。 - 线程资源泄漏:每次提交作业都创建新的
SecondaryWorker线程,未用线程池管理,高并发下会耗尽系统资源。 - 错误处理缺失:没有处理Batch执行失败、作业中断的情况,无法将作业状态更新为失败,也没有重试机制。
- 进度更新逻辑错误:
BusinessJobWorker的updateDatabase()在所有Batch处理完才执行,不符合“每个单元完成后更新进度”的要求。 - 无持久化机制:所有队列和作业状态都在内存中,服务器崩溃后作业完全丢失,无法恢复。
- Batch执行逻辑缺失:
BatchWorker没有实现具体的Batch处理逻辑,也没有反馈处理结果给BusinessJobWorker来更新进度。 - 串行/并行逻辑矛盾:
canBeParallel()返回true时设置队列大小为1,反而限制了并行,逻辑完全倒置。
数据库更新的合理方式
- 进度更新:
- 用原子更新语句避免并发覆盖:
UPDATE job SET progress = progress + :batchProgress WHERE id = :jobId,无需事务,保证每个Batch完成后即时更新。 - 高并发场景下可使用乐观锁:在作业表增加
version字段,更新时带WHERE version = :currentVersion,更新失败则重试。
- 用原子更新语句避免并发覆盖:
- 状态更新:
- 保证状态转换的幂等性:
UPDATE job SET status = :newStatus, update_time = NOW() WHERE id = :jobId AND status = :oldStatus,避免重复更新状态。 - 作业最终状态(完成/失败)需与最后一次进度更新在同一事务中,确保进度与状态一致。
- 保证状态转换的幂等性:
- Batch状态持久化:新增
batch表记录每个Batch的状态(待处理/处理中/成功/失败),重启时可根据未完成的Batch恢复作业进度。
可简化开发的Java/Quarkus特性
Java 特性
- 线程池管理:用
ExecutorService的newSingleThreadExecutor()为每个SourceId创建专属线程池,自动保证同一SourceId作业串行执行,替代手动维护队列和线程的逻辑:ConcurrentHashMap<String, ExecutorService> sourceExecutors = new ConcurrentHashMap<>(); public void submit(String sourceId, Runnable job) { sourceExecutors.computeIfAbsent(sourceId, k -> Executors.newSingleThreadExecutor()) .submit(job); } - CompletableFuture:处理Batch并行执行,简化异步结果收集和进度更新:
List<CompletableFuture<Void>> batchFutures = batches.stream() .map(batch -> CompletableFuture.runAsync(() -> processBatch(batch), batchExecutor)) .collect(Collectors.toList()); // 等待所有Batch完成,或处理异常 CompletableFuture.allOf(batchFutures.toArray(new CompletableFuture[0])) .exceptionally(ex -> { updateJobStatus(jobId, FAILED); return null; }) .thenRun(() -> updateJobStatus(jobId, COMPLETED));
Quarkus 特性
- Panache ORM:简化数据库CRUD操作,自带乐观锁支持(
@Version注解),无需手动写SQL:@Entity public class Job extends PanacheEntity { public String sourceId; public JobStatus status; public int progress; @Version public Long version; } - @Async 与 Managed Executor:Quarkus提供的异步执行框架,自动管理线程池,无需手动创建
ExecutorService:@Inject ManagedExecutor executor; public CompletionStage<Void> submitJob(String sourceId) { return executor.runAsync(() -> processJob(sourceId)); } - Quarkus Reactive Messaging:若作业来自消息队列,可直接用
@Incoming注解接收消息,结合@Blocking处理同步业务逻辑,自动实现消息重试与死信队列。 - Startup 与 ApplicationScoped:用
@ApplicationScoped和@Startup注解管理全局的SourceId线程池或作业恢复逻辑,确保应用启动时自动加载未完成的作业。
崩溃恢复与多线程并发控制的专业实现意见
- 作业状态持久化:
- 所有作业和Batch的状态必须持久化到数据库,作业表需包含:
source_id、status(待处理/运行中/完成/失败)、progress、create_time、update_time。 - 应用启动时,扫描数据库中
status为待处理或运行中的作业,重新提交到对应SourceId的执行队列。
- 所有作业和Batch的状态必须持久化到数据库,作业表需包含:
- SourceId级并发控制:
- 单实例场景:用
ConcurrentHashMap<String, ExecutorService>为每个SourceId分配单线程池,天然保证串行执行。 - 多实例场景:使用分布式锁,锁定
SourceId后再启动作业,避免多个实例同时执行同一SourceId的作业:RLock lock = redissonClient.getLock("source-lock:" + sourceId); try { if (lock.tryLock(5, TimeUnit.SECONDS)) { // 检查作业状态,避免重复执行 Job job = Job.find("sourceId = ?1 and status in (?2, ?3)", sourceId, RUNNING, QUEUED).firstResult(); if (job == null) { // 启动作业 } } } finally { if (lock.isHeldByCurrentThread()) { lock.unlock(); } }
- 单实例场景:用
- 作业中断处理:
- 为每个作业分配唯一标识,在内存中维护作业与线程的映射,支持主动中断作业。
- 作业执行时定期检查线程中断状态,若中断则更新作业状态为失败并退出。
内容的提问来源于stack exchange,提问作者Yashdeep Hinge
相关产品推荐
相关产品推荐

