带优先级的生产者-消费者模式:兼容旧Service接口的缓存实现问询
问题与解决方案:兼容旧Service接口的优先级队列实现
问题背景
你现有一个已实现的Service接口,基于OkHttp Client实现,定义如下:
public interface Service { CompletableFuture<ResponseCar> getResponseOfApple(String url, UUID id); CompletableFuture<ResponsePear> getResponseOfPear(String url, UUID id); }
现在新增需求:
- 后台Worker定期更新缓存所有Apple和Pear的REST API返回结果
- 当用户请求缓存中不存在的Pear时,该任务优先级高于后台更新任务
你已经设计了基于PriorityBlockingQueue的方案,定义了PrioritizedQueueElement和QueueElementType:
@RequiredArgsConstructor public class PrioritizedQueueElement extends Observable implements Comparable<PrioritizedQueueElement> { @NonNull @Getter private final Object argument; @Getter @Setter private Object returnValue; @NonNull @Getter private final QueueElementType priority; @Override public int compareTo(final PrioritizedQueueElement o) { if (o == null || this.priority == null || o.priority == null) { throw new NullPointerException("wrong value"); } return this.priority.compareTo(o.priority); } public <T> T getValue() { return (T) this.returnValue; } public void setReturnValue(final Object value) { this.returnValue = value; this.setChanged(); this.notifyObservers(); } } public enum QueueElementType { ONE_TIME_HIGH_PRIORITY, ONE_TIME_LOW_PRIORITY, RECURRENT_HIGH_PRIORITY, RECURRENT_LOW_PRIORITY }
现在需要解决两个核心问题:如何让新实现兼容旧Service接口?以及Observable如何转换为CompletableFuture?
解决方案
1. Observable转CompletableFuture的适配
你的思路本身没问题,只是需要把Observable的通知机制和CompletableFuture的完成逻辑绑定起来。具体做法是:
- 当创建
PrioritizedQueueElement时,同时创建一个对应的CompletableFuture - 给这个元素注册一个一次性Observer,避免内存泄漏,当元素的
returnValue被设置并触发notifyObservers时,完成这个Future - 额外处理任务执行失败的情况,比如在Worker执行API调用出错时,调用
future.completeExceptionally()
示例代码片段:
// 创建带Future的队列元素工具方法 public static <T> Pair<PrioritizedQueueElement, CompletableFuture<T>> createPrioritizedElement(Object args, QueueElementType priority) { PrioritizedQueueElement element = new PrioritizedQueueElement(args, priority); CompletableFuture<T> future = new CompletableFuture<>(); // 注册一次性Observer element.addObserver((o, arg) -> { if (arg != null && arg instanceof Exception) { future.completeExceptionally((Exception) arg); } else { future.complete(element.getValue()); } // 完成后立即移除Observer,避免内存泄漏 o.deleteObserver(this); }); return Pair.of(element, future); }
2. 兼容旧Service接口的具体实现
基于上面的适配逻辑,你可以实现新的Service实现类,同时整合优先级队列和后台Worker:
public class PrioritizedService implements Service { private final PriorityBlockingQueue<PrioritizedQueueElement> queue = new PriorityBlockingQueue<>(); private final OkHttpClient okHttpClient; private final Cache cache; // 后台Worker线程池 private final ExecutorService workerPool = Executors.newSingleThreadExecutor(); public PrioritizedService(OkHttpClient okHttpClient, Cache cache) { this.okHttpClient = okHttpClient; this.cache = cache; // 启动后台Worker循环处理队列 startWorker(); } private void startWorker() { workerPool.submit(() -> { while (!Thread.currentThread().isInterrupted()) { try { PrioritizedQueueElement element = queue.take(); // 根据元素类型和参数执行对应的API调用 handleQueueElement(element); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } catch (Exception e) { // 执行失败时通知Observer element.setChanged(); element.notifyObservers(e); } } }); } private void handleQueueElement(PrioritizedQueueElement element) throws IOException { Object args = element.getArgument(); if (args instanceof AppleTaskArgs) { AppleTaskArgs appleArgs = (AppleTaskArgs) args; ResponseCar response = callAppleApi(appleArgs.getUrl(), appleArgs.getId()); cache.putApple(appleArgs.getId(), response); element.setReturnValue(response); } else if (args instanceof PearTaskArgs) { PearTaskArgs pearArgs = (PearTaskArgs) args; ResponsePear response = callPearApi(pearArgs.getUrl(), pearArgs.getId()); cache.putPear(pearArgs.getId(), response); element.setReturnValue(response); } // 对于RECURRENT类型的任务,执行完成后重新加入队列(实现定期更新) if (element.getPriority().name().startsWith("RECURRENT_")) { // 可搭配ScheduledExecutorService实现延迟重新加入,避免高频循环 queue.add(element); } } // 实现旧接口的getResponseOfApple方法 @Override public CompletableFuture<ResponseCar> getResponseOfApple(String url, UUID id) { // 先查缓存 ResponseCar cached = cache.getApple(id); if (cached != null) { return CompletableFuture.completedFuture(cached); } // 缓存不存在,创建高优先级一次性任务 Pair<PrioritizedQueueElement, CompletableFuture<ResponseCar>> pair = createPrioritizedElement(new AppleTaskArgs(url, id), QueueElementType.ONE_TIME_HIGH_PRIORITY); queue.add(pair.getFirst()); return pair.getSecond(); } // 实现旧接口的getResponseOfPear方法 @Override public CompletableFuture<ResponsePear> getResponseOfPear(String url, UUID id) { // 先查缓存 ResponsePear cached = cache.getPear(id); if (cached != null) { return CompletableFuture.completedFuture(cached); } // 缓存不存在,创建最高优先级的一次性任务 Pair<PrioritizedQueueElement, CompletableFuture<ResponsePear>> pair = createPrioritizedElement(new PearTaskArgs(url, id), QueueElementType.ONE_TIME_HIGH_PRIORITY); queue.add(pair.getFirst()); return pair.getSecond(); } // 后台定期更新任务的启动方法(服务初始化时调用) public void startRecurrentUpdateTasks(List<String> appleUrls, List<String> pearUrls) { // 添加Apple的定期更新任务(低优先级) for (String url : appleUrls) { UUID taskId = UUID.randomUUID(); queue.add(new PrioritizedQueueElement(new AppleTaskArgs(url, taskId), QueueElementType.RECURRENT_LOW_PRIORITY)); } // 添加Pear的定期更新任务(低优先级) for (String url : pearUrls) { UUID taskId = UUID.randomUUID(); queue.add(new PrioritizedQueueElement(new PearTaskArgs(url, taskId), QueueElementType.RECURRENT_LOW_PRIORITY)); } } // 辅助类:封装任务参数 @RequiredArgsConstructor private static class AppleTaskArgs { @NonNull private final String url; @NonNull private final UUID id; // getter方法 public String getUrl() { return url; } public UUID getId() { return id; } } private static class PearTaskArgs { @NonNull private final String url; @NonNull private final UUID id; // getter方法 public String getUrl() { return url; } public UUID getId() { return id; } } // 原有OkHttp调用逻辑 private ResponseCar callAppleApi(String url, UUID id) throws IOException { // 你的OkHttp请求实现 } private ResponsePear callPearApi(String url, UUID id) throws IOException { // 你的OkHttp请求实现 } }
3. 可选优化:替换Observable为直接绑定CompletableFuture
Java的Observable是比较老旧的类,灵活性有限。你可以直接修改PrioritizedQueueElement,让它持有对应的CompletableFuture,代码会更简洁易维护:
@RequiredArgsConstructor public class PrioritizedQueueElement implements Comparable<PrioritizedQueueElement> { @NonNull @Getter private final Object argument; @NonNull @Getter private final QueueElementType priority; @Getter private final CompletableFuture<Object> future = new CompletableFuture<>(); @Override public int compareTo(final PrioritizedQueueElement o) { if (o == null || this.priority == null || o.priority == null) { throw new NullPointerException("wrong value"); } return this.priority.compareTo(o.priority); } public void complete(Object value) { this.future.complete(value); } public void completeExceptionally(Throwable ex) { this.future.completeExceptionally(ex); } }
这种方式下,创建元素时直接返回element.getFuture()即可,不需要Observer适配,逻辑更清晰。
总结
你的核心设计思路(用PriorityBlockingQueue处理优先级任务)是完全没问题的,只是需要把队列元素的结果通知机制和CompletableFuture做适配。上面的两种方案都可以实现旧接口的兼容,推荐第二种直接绑定Future的方式,代码更简洁,也避免了旧Observable带来的潜在问题。
内容的提问来源于stack exchange,提问作者akop
相关产品推荐
相关产品推荐

