自定义线程池-实现任务0丢失的处理策略

自定义线程池-实现任务0丢失的处理策略
设计一个线程池要求如下队列最大容量为10内存队列。当队列满了之后拒绝策略将新的任务写入数据库。从队列中取任务时若该队列为空能够从数据库中加载之前被拒绝的任务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 LinkedBlockingQueueRunnable { 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 QueueRunnableTask 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定期监控线程池的运行状态 // 包括活跃线程数、队列大小、已完成任务数等指标 } }