Kol*_*oDS 4 future message-queue rust
我正在尝试编写一个简单的期货-rs mpsc 队列用法示例:
extern crate futures; // v0.1 (old)
use futures::{Sink, Stream};
use futures::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel::<i32>(1000);
let handle = thread::spawn(move || {
tx.clone().send(1);
tx.clone().send(2);
tx.clone().send(3);
});
let mut rx = rx.map(|x| {
println!("stream: {}", x);
x * x
});
handle.join().unwrap();
rx.poll().unwrap();
}
Run Code Online (Sandbox Code Playgroud)
但它不会向控制台输出任何内容(我希望它打印stream: 1,stream: 2和stream: 3)。我也试图取代rx.poll().unwrap()有rx.wait(),但它仍然什么也不输出。而且我在 futures-rs 文档中没有找到任何使用示例。我究竟做错了什么?
这是强烈建议,编译器会告诉你阅读警告和错误消息。这是带有编译器的静态类型语言的一大好处:
warning: unused result which must be used: futures do nothing unless polled, #[warn(unused_must_use)] on by default
--> src/main.rs:11:9
|
11 | tx.clone().send(1);
| ^^^^^^^^^^^^^^^^^^^
warning: unused result which must be used: futures do nothing unless polled, #[warn(unused_must_use)] on by default
--> src/main.rs:12:9
|
12 | tx.clone().send(2);
| ^^^^^^^^^^^^^^^^^^^
warning: unused result which must be used: futures do nothing unless polled, #[warn(unused_must_use)] on by default
--> src/main.rs:13:9
|
13 | tx.clone().send(3);
| ^^^^^^^^^^^^^^^^^^^
Run Code Online (Sandbox Code Playgroud)
我不是期货专家,但这编译时没有警告并打印所有三个值:
extern crate futures; // 0.1.23
use futures::{sync::mpsc, Async, Future, Sink, Stream};
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel(1000);
let handle = thread::spawn(move || {
tx.send(1)
.and_then(|tx| tx.send(2))
.and_then(|tx| tx.send(3))
.wait()
.expect("Unable to send");
});
let mut rx = rx.map(|x| x * x);
handle.join().unwrap();
while let Ok(Async::Ready(Some(v))) = rx.poll() {
println!("stream: {}", v);
}
}
Run Code Online (Sandbox Code Playgroud)
and_then用于在前一个值之后发送每个后续值。wait用于阻塞生成的线程,直到所有内容都成功发送。该poll方法用于从队列中获取值,直到它用完为止。有多种方法可能会失败,我将它们全部忽略,只专注于成功案例。
| 归档时间: |
|
| 查看次数: |
2962 次 |
| 最近记录: |