有两种可能性:如果您想在客户端估计超时之前提前关闭连接,您可以使用具有合理持续时间的Flux::take。如果您想在响应中添加任意信息,请使用Flux::buffer 并将自定义消息添加到已缓冲的对象中。
为了展示这一点,请考虑以下示例:
@Slf4j
@RestController
@RequestMapping("api")
public class DemoController {
@GetMapping("flux-take")
public Flux<Result> fluxTake() {
Flux<Result> resultFlux = createFlux();
resultFlux.share().subscribe(result -> log.info("result={}", result));
return resultFlux.take(Duration.ofMillis(1500));
}
@GetMapping("flux-buffer")
public Mono<BufferedResult> fluxBuffer() {
Flux<Result> resultFlux = createFlux();
resultFlux.share().subscribe(result -> log.info("result={}", result));
return resultFlux.buffer(Duration.ofMillis(1500))
.next()
.map(results -> new BufferedResult("Request timed out, returning buffered results", results));
}
private Flux<Result> createFlux() {
return Flux.just(1, 2, 3, 4, 5, 6)
.delayElements(Duration.ofMillis(300))
.timed()
.map(timedId -> new Result(timedId.get(), timedId.elapsed().toMillis()));
}
@Data
@AllArgsConstructor
public static class Result {
private final Integer id;
private final Long elapsedMillis;
}
@Data
@AllArgsConstructor
public static class BufferedResult {
private final String message;
private final List<Result> results;
}
}
从客户的角度来看,结果是这样的:
$ http :8080/api/flux-take --timeout=2
HTTP/1.1 200 OK
Content-Type: application/json
transfer-encoding: chunked
[
{
"elapsedMillis": 304,
"id": 1
},
{
"elapsedMillis": 300,
"id": 2
},
{
"elapsedMillis": 304,
"id": 3
},
{
"elapsedMillis": 302,
"id": 4
}
]
$ http :8080/api/flux-buffer --timeout=2
HTTP/1.1 200 OK
Content-Length: 187
Content-Type: application/json
{
"message": "Request timed out, returning buffered results",
"results": [
{
"elapsedMillis": 301,
"id": 1
},
{
"elapsedMillis": 301,
"id": 2
},
{
"elapsedMillis": 302,
"id": 3
},
{
"elapsedMillis": 304,
"id": 4
}
]
}
请注意客户端的 2 秒超时和应用程序内的 1500 毫秒超时。客户端永远不会达到 2 秒超时,并且总是会收到正确的响应。
IntelliJ 将在 resultFlux.share().subscribe(...) 上显示提示 Calling 'subscribe' in non-blocking scope,但我想不出更好的方法来完成此任务。