【问题标题】:is there a way to read a InputStream asynchronously with Reactor or transform to bytes?有没有办法使用 Reactor 异步读取 InputStream 或转换为字节?
【发布时间】:2021-12-27 19:27:23
【问题描述】:

我正在尝试将文件上传到 S3,但 JVM 说我在代码片段中有一个线程阻塞方法调用,在调用 file.readAllBytes() 时线程不应该被阻塞,所以有没有办法有没有办法使方法与 Flux 或 Mono 异步?或任何其他方法来解决这个问题?

private Mono<Boolean> uploadFile(InputStream file, String bucket, String name) {
        try {
            return uploadAdapter.uploadObject(bucket,name,file.readAllBytes());
        } catch (IOException e) {
            return Mono.just(false);
        }
    }
@Override
    public Mono<Boolean> uploadObject(String bucketName, String objectKey, byte[] fileContent) {
        return Mono.fromFuture(
                        s3AsyncClient.putObject(configurePutObject(bucketName, objectKey),
                                AsyncRequestBody.fromBytes(fileContent)))
                .map(response -> response.sdkHttpResponse().isSuccessful());
    }

【问题讨论】:

  • 您无法真正从 InputStream 这样的同步 API 转换为异步 API。您可以通过让它在另一个线程中工作来“假装”它是异步的,但这并不能真正使其异步。

标签: java spring asynchronous project-reactor


【解决方案1】:

由于InputStream 是一个同步 API,您有 2 个选择,对于任何其他同步 API 都是如此:

  1. 切换到另一个 API。不开玩笑,这可能是解决许多问题的好方法。一般来说,反应式概念和异步非常常见,对于大多数需求,有一个替代的异步库可以做同样的事情。在您的情况下,您可以使用java.nio2,或Reactor-Netty library, which has great solutions for this use case 中的函数。
  2. 使用另一个调度程序。 Project reactor 建议您的所有异步调用都将在一个非阻塞调度程序上运行,而对于同步调用,请使用另一个(阻塞)线程每个请求的调度程序。您可以使用两个这样的调度程序:single()boundedElastic()。不同之处在于boundedElastic() 限制了您可以打开的线程数,因此您最终会将线程用作阻塞队列,这比single() 更安全。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2013-08-14
    • 2020-04-16
    • 1970-01-01
    • 1970-01-01
    • 2020-09-12
    • 2020-09-12
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多