如何实现流无数据2秒后自动终止的Rust异步函数?
问题分析与解决方案
你的实现逻辑存在问题,并非tokio::io::Lines的固有问题。以下是具体分析和正确实现:
问题原因
你的三种实现都复用了未完成的stream.next() Future:当定时器到期后,虽然你退出了循环,但之前启动的读取Future仍在等待stdin输入,导致程序无法立即终止,必须等下一次输入触发读取完成后才会退出。
- 尝试1:复用同一个
next_fut和timer,定时器到期break后,next_fut仍处于pending状态,持有stream引用,读取操作未被立即取消。 - 尝试2:复用同一个
next_fut,超时发生后该Future仍在等待输入,直到输入到来才会完成。 - 尝试3:
stream.timeout()的逻辑本身正确,但循环处理中,超时返回的Err会直接退出循环,不过若之前的读取Future未被正确取消,仍可能导致程序卡住。
正确实现
使用Tokio的select!宏,每次循环创建独立的读取和定时器Future,确保超时后能立即终止程序:
use std::time::Duration; use futures::stream::{Stream, StreamExt}; use tokio::time::sleep; const TWO_SECONDS: Duration = Duration::from_millis(2_000); async fn terminate_after_2_seconds_of_no_items<S, T>(mut stream: S) where S: Stream<Item = T> + Unpin, { loop { tokio::select! { item = stream.next() => { match item { Some(_) => { println!("received item, timer reset"); } None => { println!("stream ended, terminating"); break; } } } _ = sleep(TWO_SECONDS) => { println!("timer expired, too bad, quitting"); break; } } } }
实现说明
- 每次循环通过
select!同时等待两个异步任务:流的下一个项,以及2秒定时器。 - 收到流项时,打印信息并重新进入循环,自动重置定时器(每次循环创建新的
sleep实例)。 - 定时器到期时,立即break退出循环,函数返回。此时当前的读取Future会被丢弃,Tokio自动取消底层读取操作,程序直接终止,不会等待后续输入。
关于你提到的futures::stream::select无法复用定时器的问题,确实into_stream()会消耗Future无法重置,但Tokio的select!宏更适合这种场景,无需手动管理定时器重置逻辑。
内容的提问来源于stack exchange,提问作者virchau13
相关产品推荐
相关产品推荐

