- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我想将一些读取器/写入器转换为元素的记录管道。我设法按照this answer使用futures::stream::unfold
来使读者->流方向。但是,我在使用接收器->写入器时遇到了麻烦。我基本上是在寻找unfold
的逆函数。
我知道有AsyncWriterExt::into_sink
,但这仅在我可以生成所有字节以成批写入的情况下才有效。我还发现this answer建议在drain()
之后使用with()
。但是由于存在生命周期问题,此方法不起作用(FnMut
无法有效地存储对writer的引用,或者至少我没有做到这一点。
因此,我正在寻找的是一些我可以调用的函数,例如fold(initial_state, |element| {(writer.write(element).await, new_state)})
。您明白了(我希望)。
我也看到有async_codec
,但对我来说似乎有点过头了。同时,我求助于将所有写入存储为流,然后使用writer.into_sink().with_flat_map()
。但这真的很丑。
最佳答案
编辑:我显然不是唯一想要这样做的人,请参阅upstream implementation。 future 的用户(呵呵)将可以简单地使用futures::sink::unfold
。
好的,我鼓起勇气,根据futures::stream::unfold
一起砍了一些东西:
fn fold<T, F, Fut, Item, E>(init: T, f: F) -> FoldSink<T, F, Fut>
where
F: FnMut(T, Item) -> Fut,
Fut: Future<Output = Result<T, E>>
{
FoldSink {
f,
state: Some(init),
fut: None,
}
}
use pin_project::pin_project;
#[pin_project]
struct FoldSink<T, F, Fut> {
f: F,
state: Option<T>,
#[pin]
fut: Option<Fut>,
}
impl <T, Item, F, Fut, E> futures::sink::Sink<Item> for FoldSink<T, F, Fut>
where
F: FnMut(T, Item) -> Fut,
Fut: Future<Output = Result<T, E>>
{
type Error = E;
fn poll_ready(self: std::pin::Pin<&mut Self>, ctx: &mut std::task::Context<'_>) -> Poll<Result<(), E>> {
let mut this = self.project();
match this.fut.as_mut().as_pin_mut() {
Some(fut) => {
match fut.poll(ctx) {
Poll::Ready(Ok(new_state)) => {
this.fut.set(None);
*this.state = Some(new_state);
Poll::Ready(Ok(()))
},
Poll::Ready(Err(e)) => {
this.fut.set(None);
Poll::Ready(Err(e))
},
Poll::Pending => Poll::Pending,
}
},
None => {
Poll::Ready(Ok(()))
}
}
}
fn start_send(self: std::pin::Pin<&mut Self>, item: Item) -> Result<(), E> {
let mut this = self.project();
this.fut.set(Some((this.f)(this.state.take().expect("todo invalid state"), item)));
Ok(())
}
fn poll_flush(self: std::pin::Pin<&mut Self>, ctx: &mut std::task::Context<'_>) -> Poll<Result<(), E>> {
self.poll_ready(ctx)
}
fn poll_close(mut self: std::pin::Pin<&mut Self>, ctx: &mut std::task::Context<'_>) -> Poll<Result<(), E>> {
futures::ready!(self.as_mut().poll_ready(ctx))?;
let this = self.project();
this.state.take().unwrap();
Poll::Ready(Ok(()))
}
}
欢迎评论和改进!
关于rust - rust 铸就 future 的沉沦,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/65044743/
我正在通过这个示例https://www.rusoto.org/futures.html学习Rust和Rusoto 而且我发现许多代码已经过时了。所以我改变了这样的代码: use rusoto_cor
这是一个理论问题。我有一个服务可以调用来完成工作,但该服务可能无法完成所有工作,因此我需要调用第二个服务来完成它。 我想知道是否有办法在没有 Await.result 的情况下做类似的事情map 函数
这个问题是关于如何阅读 Rust 文档并提高我对 Rust 的理解,从而了解如何解决这个特定的编译器错误。 我读过 tokio docs并试验了许多 examples .在编写自己的代码时,我经常遇到
我有一个使用分页的 HTTP api,我想将它包装到一个通用的 Rust 流中,以便所有端点都可以使用相同的接口(interface),这样我就可以使用 Stream 附带的特征函数特征。 我收到了这
我正在查看 AKKA 的 Java Futures API,我看到了很多处理同一类型的多个 future 的方法,但我没有看到任何处理不同类型的 future 的方法。我猜我让事情变得更加复杂了。 无
环境:Akka 2.1,scala 版本 2.10.M6,JDK 1.7,u5 现在是我的问题: 我有: future1 = Futures.future(new Callable>(){...});
我有一些代码可以将请求提交给另一个线程,该线程可能会也可能不会将该请求提交给另一个线程。这会产生 Future> 的返回类型.是否有一些非令人发指的方法可以立即将其变成 Future等待整个 futu
如果我有以下代码: Future a = new Future(() { print('a'); return 1; }); Future b = new Future.error('Error!')
我一直试图简化我在 Scala 中做 future 的方式。我有一次收到了 Future[Option[Future[Option[Boolean]]但我在下面进一步简化了它。有没有更好的方法来简化这
Scala 中从 Future[Option[Future[Int]]] 转换的最干净的方法是什么?至 Future[Option[Int]] ?甚至有可能吗? 最佳答案 有两个嵌套Future s
使用下面的示例,future2 如何在 future1 完成后使用 future1 的结果(不阻塞 future3 从被提交)? from concurrent.futures import Proc
这两个类代表了并发编程的优秀抽象,因此它们不支持相同的 API 有点令人不安。 具体根据docs : asyncio.Future is almost compatible with concurre
我正在尝试使用 wasm_bindgen 实现 API 类使用异步调用。 #![allow(non_snake_case)] use std::future::Future; use serde::{
这个问题在这里已经有了答案: Futures / Success race (3 个回答) 去年关闭。 所有的 future 最终可能会成功(有些可能会失败),但我们希望第一个成功。并希望将这一结果表
我在练习asyncio在编写多线程代码多年之后。 注意到一些我觉得很奇怪的东西。都在 asyncio在 concurrent有一个Future目的。 from asyncio import Futur
如何将Future[Option[Future[Option[X]]]]转换为Future[Option[X]]? 如果它是 TraversableOnce 而不是 Option 我会使用 Futur
我正在尝试同时发送 HTTP 请求。为此,我使用 concurrent.futures 这是简单的代码: import requests from concurrent import futures
我们在 vertx 中使用 Futures 的例子如下: Future fetchVehicle = getUserBookedVehicle(routingContext, client);
下面的函数,取自 here : fn connection_for( &self, pool_key: PoolKey, ) -> impl Future>, ClientError>
我正在围绕Java库编写一个小的Scala包装器。 Java库有一个对象QueryExecutor,它公开了2种方法: execute(query):结果 asyncExecute(query):Li
我是一名优秀的程序员,十分优秀!