如何从BroadcastProcessor提取最后n个元素并通过REST接口返回?
问题描述
需要实现一个方法,从BroadcastProcessor中提取最多n个元素组成列表,经映射处理后通过HTTP返回。现有getTimesCompact方法无法正常工作,要么返回空对象,要么出现错误。相关代码如下:
BroadcastProcessor<HashMap<String, Long>> timesSink = BroadcastProcessor.create(); // 正常工作 public void addTimes(HashMap<String, Long> times) { this.timesSink.onNext(times); } // 正常工作 @GET @Path("/stream") @Produces(APPLICATION_JSON) public Multi<HashMap<String, Long>> getTimesStream() { return Multi.createBy().replaying().ofMulti(this.timesSink); } // 无法正常工作的方法 // n 是需要重放的元素数量 @GET @Path("/compact") @Produces(APPLICATION_JSON) public Uni<HashMap<String, Long>> getTimesCompact() { return Multi.createBy().replaying().upTo(n).ofMulti(this.timesSink) .collect().asList() .map(this::toCompactTimes); } private HashMap<String, Long> toCompactTimes(List<HashMap<String, Long>> times) { // 逻辑:将HashMap列表计算平均值 }
问题原因
核心问题在于collect().asList()的工作机制:它需要流发送onComplete信号才会返回收集到的列表,但BroadcastProcessor是热流,除非手动调用timesSink.onComplete()(这会关闭处理器,无法再接收新元素),否则重放的流永远不会结束,导致Uni一直处于等待状态,最终请求超时或返回空。
另外,如果缓存中的元素数量不足n个,collect().asList()会一直等待新元素产生,同样导致请求挂起。
解决方案
根据需求选择以下两种方案:
方案1:仅获取当前缓存中的元素(不等待新元素)
如果只需要当前缓存里已有的最多n个元素,不需要等待后续新增元素,可在发送完缓存元素后立即取消订阅,触发流完成:
@GET @Path("/compact") @Produces(APPLICATION_JSON) public Uni<HashMap<String, Long>> getTimesCompact() { return Multi.createBy().replaying().upTo(n).ofMulti(this.timesSink) // 发送完缓存元素后取消订阅,让流结束 .onSubscription().invoke(subscription -> { subscription.request(Long.MAX_VALUE); // 请求所有缓存元素 subscription.cancel(); // 发送完成后取消订阅,触发流结束 }) .collect().asList() // 处理空列表情况,避免后续方法抛出异常 .map(list -> list.isEmpty() ? new HashMap<>() : toCompactTimes(list)); }
方案2:等待凑够n个元素(含未来新增元素)
如果需要等待直到凑齐n个元素(包括后续新增的),需添加超时逻辑避免请求无限挂起:
import io.smallrye.mutiny.TimeoutException; import java.time.Duration; @GET @Path("/compact") @Produces(APPLICATION_JSON) public Uni<HashMap<String, Long>> getTimesCompact() { return Multi.createBy().replaying().upTo(n).ofMulti(this.timesSink) .take(n) // 拿到n个元素后立即停止 .collect().asList() .map(this::toCompactTimes) // 超时处理:10秒未凑齐n个元素则返回空HashMap .onFailure(TimeoutException.class).recoverWithItem(new HashMap<>()) .ifNoItem().after(Duration.ofSeconds(10)).fail(); }
额外注意事项
- 确保
n是正整数,若n初始化为0,会直接返回空列表。 - 检查
toCompactTimes方法是否处理空列表的情况,避免抛出空指针或其他异常。
内容的提问来源于stack exchange,提问作者jpkmiller
相关产品推荐
相关产品推荐

