You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

带优先级的生产者-消费者模式:兼容旧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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 06:51:05