【问题标题】:How do I convert a futures_io::AsyncRead to rusoto::ByteStream?如何将 futures_io::AsyncRead 转换为 rusoto::ByteStream?
【发布时间】:2020-06-11 10:25:32
【问题描述】:

我正在尝试构建一个从 SFTP 服务器提取文件并将它们上传到 S3 的服务。

对于 SFTP 部分,我使用的是async-ssh2,它为我提供了一个实现futures::AsyncRead 的文件处理程序。由于这些 SFTP 文件可能非常大,因此我试图将这个 File 处理程序转换为可以使用 Rusoto 上传的 ByteStream。看起来ByteStream 可以用futures::Stream 初始化。

我的计划是在File 对象上实现Stream(基于代码here)以与Rusoto 兼容(代码在下面复制以供后代使用):

use core::pin::Pin;
use core::task::{Context, Poll};
use futures::{ready, stream::Stream};

pub struct ByteStream<R>(R);

impl<R: tokio::io::AsyncRead + Unpin> Stream for ByteStream<R> {
    type Item = u8;

    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context) -> Poll<Option<Self::Item>> {
        let mut buf = [0; 1];

        match ready!(Pin::new(&mut self.0).poll_read(cx, &mut buf)) {
            Ok(n) if n != 0 => Some(buf[0]).into(),
            _ => None.into(),
        }
    }
}

这是一个很好的方法吗?我看到了this question,但它似乎在使用tokio::io::AsyncRead。是否使用tokio 的规范方式来执行此操作?如果是这样,有没有办法从futures_io::AsyncRead 转换为tokio::io::AsyncRead

【问题讨论】:

    标签: rust rusoto


    【解决方案1】:

    这就是我进行转换的方式。我基于上面的代码,除了我使用了更大的缓冲区(8 KB)来减少网络调用的数量。

    use bytes::Bytes;
    use core::pin::Pin;
    use core::task::{Context, Poll};
    use futures::{ready, stream::Stream};
    use futures_io::AsyncRead;
    use rusoto_s3::StreamingBody;
    
    const KB: usize = 1024;
    
    struct ByteStream<R>(R);
    
    impl<R: AsyncRead + Unpin> Stream for ByteStream<R> {
        type Item = Result<Bytes, std::io::Error>;
    
        fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context) -> Poll<Option<Self::Item>> {
            let mut buf = vec![0_u8; 8 * KB];
    
            match ready!(Pin::new(&mut self.0).poll_read(cx, &mut buf[..])) {
                Ok(n) if n != 0 => Some(Ok(Bytes::from(buf))).into(),
                Ok(_) => None.into(),
                Err(e) => Some(Err(e)).into(),
            }
        }
    }
    

    允许我这样做:

    fn to_streamingbody(body: async_ssh2::File) -> Option<StreamingBody> {
        let stream = ByteStream(body);
        Some(StreamingBody::new(stream))
    }
    

    (注意rusoto::StreamingBodyrusoto::ByteStream 是别名)

    【讨论】:

    • 请注意,这将在每次调用 poll_next() 时分配并清零一个 8MB 缓冲区。最好避免这种费用,例如将buf 设为共享的thread local 缓冲区。另外,8MB 是probably much larger than you need
    • 哦,酷!我一直在尝试这样做,但不知道如何在没有不安全块的情况下绕过编译器错误。并感谢其他链接,将进行这些更改。
    • 更新为使用 8 KB 缓冲区大小;但是,我不确定如何在不复制的情况下使 thread_local 工作(这可能接近于需要一个新问题)。 Bytes::from(buf: Vec&lt;u8&gt;) 需要 Vec 的所有权,我们在从 RefCell 借用时没有(因此需要克隆)。如果我们将 buf 设为 [u8],我在 Bytes 中看到的唯一方法是 copy_from_slice,它再次复制它。
    • 转念一想may be impossible without GATs:“cramertj 对 [Stream] 提出的主要担忧是,就像迭代器一样,它总是将每个项目的所有权归还给它的调用者......实际上,如果许多流/迭代器实现可以拥有一些可以反复使用的内部存储,它们会更有效。例如,它们可能有一个内部缓冲区,当调用poll_next 时,它们会回馈(完成后) 对该缓冲区的引用。"
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-01-17
    • 2010-09-08
    • 2013-09-21
    • 2011-05-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多