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

如何从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 08:45:39