自定义线程池-实现任务0丢失的处理策略
2026/8/2 20:16:48 网站建设 项目流程

设计一个线程池,要求如下:

  1. 队列最大容量为10(内存队列)。
  2. 当队列满了之后,拒绝策略将新的任务写入数据库。
  3. 从队列中取任务时,若该队列为空,能够从数据库中加载之前被拒绝的任务


(1)自定义Runnable接口继承Serializable实现可序列化

public interface SerializableTask extends Runnable,Serializable { }

(2) 自定义Runable任务序列化

public class CustomTask { public static String serializedTask(SerializableTask runnable){ try(ByteArrayOutputStream baos = new ByteArrayOutputStream(); ObjectOutputStream oos = new ObjectOutputStream(baos)) { // 序列化任务对象 oos.writeObject(runnable); return Base64.getEncoder().encodeToString(baos.toByteArray()); }catch (Exception e){ throw new RuntimeException("无法序列化"); } } public static SerializableTask deserialization(String serializedTask){ // 反序列化任务 byte[] data = Base64.getDecoder().decode(serializedTask); try (ByteArrayInputStream bais = new ByteArrayInputStream(data); ObjectInputStream ois = new ObjectInputStream(bais)) { SerializableTask task = (SerializableTask) ois.readObject(); return task; }catch (Exception e){ throw new RuntimeException("无法反序列化"); } } }

(3)自定义阻塞队列 (DatabaseBackedBlockingQueue)

  • 继承LinkedBlockingQueue并重写关键方法
  • take()方法逻辑
    • 优先从内存队列取任务
    • 队列为空时从数据库加载
    • 数据库也为空时阻塞等待新任务
  • offer()方法:队列未满时接受,满时返回false触发拒绝策略
public class CustomBlockQueue extends LinkedBlockingQueue<Runnable> { private RunnableTaskService runnableTaskService; public CustomBlockQueue(int maxLocalCapacity, RunnableTaskService runnableTaskService) { super(maxLocalCapacity); this.runnableTaskService = runnableTaskService; } @Override public Runnable take() throws InterruptedException { // 1. 优先检查本地队列 Runnable task = super.poll(); if (task != null){ return task; } // 2. 本地队列为空时尝试从数据库加载 while (true) { System.out.println(Thread.currentThread().getName()+"进入循环"); RunnableTask dbTask = runnableTaskService.loadTask(); if (dbTask != null) { String taskName = dbTask.getTaskName(); Runnable deserialization = CustomTask.deserialization(taskName); dbTask.setTaskState(1); runnableTaskService.updateRunnableTask(dbTask); return deserialization; } // 3. 数据库为空则等待新任务 if (isEmpty()) { task = super.take(); // 阻塞直到有新任务 return task; } } } }

(4)拒绝策略 (DatabaseRejectionHandler)

  • 实现RejectedExecutionHandler接口
  • 当内存队列满时将任务存入数据库
  • 任务存入后会被后续的take()方法加载执行
public class DatabaseRejectionHandler implements RejectedExecutionHandler { private RunnableTaskService runnableTaskService; public DatabaseRejectionHandler(RunnableTaskService runnableTaskService){ this.runnableTaskService = runnableTaskService; } @Override public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) { SerializableTask serializableTask = (SerializableTask) r; String s = CustomTask.serializedTask(serializableTask); RunnableTask build = RunnableTask.builder() .taskName(s) .taskState(0) .build(); int i = runnableTaskService.saveRunnableTask(build); System.out.println("保存到数据库中: " + i); } }

资源管理

  • 核心/最大线程数根据容器资源动态调整
  • 线程工厂添加命名前缀(便于监控)
  • 保活时间控制闲置线程销毁
// 5. 监控线程池状态 ScheduledExecutorService monitor = Executors.newSingleThreadScheduledExecutor(); monitor.scheduleAtFixedRate(() -> { System.out.println("\n[监控] 活跃线程: " + executor.getActiveCount() + " | 队列大小: " + executor.getQueue().size() + " | 总完成任务: " + executor.getCompletedTaskCount()); }, 1, 2, TimeUnit.SECONDS);

自定义线程工厂

// 自定义线程工厂 static class NamedThreadFactory implements ThreadFactory { private final AtomicInteger counter = new AtomicInteger(1); private final String namePrefix; public NamedThreadFactory(String namePrefix) { this.namePrefix = namePrefix; } @Override public Thread newThread(Runnable r) { return new Thread(r, namePrefix + "-" + counter.getAndIncrement()); } }

测试

public class RunnableTaskServiceImpl implements RunnableTaskService { // 任务队列,用于存储待处理的任务 private final Queue<RunnableTask> localQueue = new LinkedList<>(); private AtomicInteger count = new AtomicInteger(); public int count(){ return count.get(); } /** * 保存可运行任务到队列 * * @param runnableTask 待保存的任务对象 * @throws IllegalArgumentException 如果任务对象为null */ @Override public void saveRunnableTask(RunnableTask runnableTask) { // NPE检查:验证任务对象不为空 if (runnableTask == null) { return; } // 将任务添加到队列中 localQueue.add(runnableTask); count.incrementAndGet(); } /** * 从队列中加载并移除一个任务 * * @return 队列中的第一个任务,如果队列为空则返回null */ @Override public RunnableTask loadTask() { // 从队列头部取出任务(如果队列为空,poll()会返回null) RunnableTask runnableTask = localQueue.poll(); // 可以根据业务需求添加日志 if (runnableTask == null) { System.out.println("[警告] 队列中没有可加载的任务"); } count.decrementAndGet(); return runnableTask; } /** * 处理所有任务的主流程 * 创建线程池并提交任务进行并发执行 */ @Override public void handlerAllTask() { ThreadPoolExecutor threadPoolExecutor = null; try { // 步骤1:初始化自定义阻塞队列 // 队列容量为10,当队列满时会触发自定义的处理逻辑 CustomBlockQueue customBlockQueue = new CustomBlockQueue(10, this); // 步骤2:初始化数据库拒绝策略处理器 // 当线程池和队列都满时,使用此处理器将任务持久化到数据库 DatabaseRejectionHandler databaseRejectionHandler = new DatabaseRejectionHandler(this); // 步骤3:创建线程池 // 核心线程数:4,最大线程数:4,空闲线程存活时间:0秒 // 使用自定义阻塞队列和拒绝策略 threadPoolExecutor = new ThreadPoolExecutor( 4, // 核心线程数 4, // 最大线程数 0, // 空闲线程存活时间 TimeUnit.SECONDS, // 时间单位 customBlockQueue, // 工作队列 databaseRejectionHandler // 拒绝策略 ); // NPE检查:确保线程池创建成功 Objects.requireNonNull(threadPoolExecutor, "线程池创建失败"); // 步骤4:批量提交任务 int totalTasks = 50; System.out.println("开始提交任务, 总数: " + totalTasks); // 循环创建并提交任务 for (int i = 1; i <= totalTasks; i++) { final int taskId = i; // 创建可序列化的任务对象 SerializableTask serializableTask = () -> { // 打印当前执行任务的线程名称和任务ID System.out.println(Thread.currentThread().getName() + " 执行任务: " + taskId); }; // 提交任务到线程池执行 try { threadPoolExecutor.execute(serializableTask); } catch (Exception e) { // 捕获任务提交时的异常(如线程池已关闭) System.err.println("任务 " + taskId + " 提交失败: " + e.getMessage()); e.printStackTrace(); } } System.out.println("所有任务已提交完成"); // 步骤5:优雅关闭线程池(可选) // 注意:根据业务需求决定是否需要等待任务完成 // threadPoolExecutor.shutdown(); // if (!threadPoolExecutor.awaitTermination(60, TimeUnit.SECONDS)) { // threadPoolExecutor.shutdownNow(); // } } catch (Exception e) { // 捕获整个流程中的异常,确保程序不会崩溃 System.err.println("任务处理流程发生异常: " + e.getMessage()); e.printStackTrace(); // 如果线程池已创建,尝试关闭 if (threadPoolExecutor != null) { try { threadPoolExecutor.shutdownNow(); } catch (Exception shutdownException) { System.err.println("线程池关闭失败: " + shutdownException.getMessage()); } } } // 注释说明:监控线程池状态的代码(已注释) // 可以使用ScheduledExecutorService定期监控线程池的运行状态 // 包括活跃线程数、队列大小、已完成任务数等指标 } }

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询