如何使用tokio-tungstenite的send_all方法发送WebSocket消息流
使用tokio-tungstenite发送消息流的正确姿势
我来帮你搞定这个问题!你遇到的核心问题是过时的流构造方法、错误类型不匹配,还有错误的类型注解。咱们一步步来修正:
首先,换掉过时的iter_ok
futures::stream::iter_ok早就被废弃了,现在官方推荐用futures::stream::iter来把集合转换成Stream。它的类型推断更友好,用法也更简洁。
解决错误类型不匹配问题
iter返回的Stream错误类型是Infallible(表示这个流永远不会产生错误),但tokio-tungstenite的Sink要求传入的Stream错误类型必须和Sink自身的错误类型(也就是tungstenite::Error)一致。所以我们需要用map_err把Infallible转换成一个合理的tungstenite::Error变体(比如ConnectionClosed,因为这个错误永远不会实际触发,选哪个都不影响逻辑)。
修正错误的类型注解
你之前给send_stream标注的tokio_tungstenite::WebSocketStream是完全错误的——WebSocketStream是同时实现Sink和Stream的WebSocket连接类型,而我们这里只需要一个普通的消息流,不需要这个类型。其实完全可以省略类型注解,让编译器自动推断,或者标注成impl Stream<Item = tungstenite::Message, Error = tungstenite::Error>。
完整的正确代码示例
use futures::stream::iter; use tokio_tungstenite::{connect_async, tungstenite}; use futures::{SinkExt, StreamExt}; async fn send_messages(url: &str) -> Result<(), tungstenite::Error> { // 建立WebSocket连接 let (ws_stream, _response) = connect_async(url).await?; let (mut sink, _stream) = ws_stream.split(); // 准备要发送的消息列表 let my_messages = vec![ tungstenite::Message::Text("message_1".to_string()), tungstenite::Message::Text("message_2".to_string()), ]; // 将消息列表转换成符合要求的Stream let send_stream = iter(my_messages) .map_err(|_| tungstenite::Error::ConnectionClosed); // 转换错误类型 // 发送整个消息流 sink.send_all(&send_stream).await?; Ok(()) }
额外说明
- 我用了
.await而不是.wait(),这是tokio异步环境下的标准做法,避免阻塞线程。如果一定要用同步阻塞的方式,也可以把.await换成.wait(),但要确保在合适的线程上下文里执行。 - 如果你不想手动处理错误类型,也可以用
iter(my_messages).map(Ok)把消息包装成Result类型,不过还是需要map_err把Infallible转成tungstenite::Error,不如上面的方式直接简洁。
内容的提问来源于stack exchange,提问作者james pearson
相关产品推荐
相关产品推荐

