如何在同一TcpStream上使用不同Codec?Framed构造会消耗该流
解决方案
一、实现多Codec数据传输的两种方案
方案1:自定义支持多类型的Codec
最直接的方式是编写自定义Codec,让它兼容文本行和二进制数据的编解码。先定义枚举区分数据类型,再实现Encoder和Decoder trait:
use bytes::Bytes; use tokio_util::codec::{Decoder, Encoder}; use std::io; #[derive(Debug)] enum Message { Text(String), Binary(Bytes), } struct MultiCodec; impl Decoder for MultiCodec { type Item = Message; type Error = io::Error; fn decode(&mut self, src: &mut bytes::BytesMut) -> Result<Option<Self::Item>, Self::Error> { // 自定义协议头:首字节标记数据类型,0x00为文本,0x01为二进制 if src.len() < 1 { return Ok(None); } let flag = src[0]; match flag { 0x00 => { src.split_to(1); let mut lines_codec = tokio_util::codec::LinesCodec::new(); lines_codec.decode(src).map(|opt| opt.map(Message::Text)) } 0x01 => { if src.len() < 5 { // 1字节标记 + 4字节数据长度 return Ok(None); } let len = u32::from_be_bytes(src[1..5].try_into()?) as usize; if src.len() < 5 + len { return Ok(None); } src.split_to(1); src.split_to(4); let data = src.split_to(len).freeze(); Ok(Some(Message::Binary(data))) } _ => Err(io::Error::new(io::ErrorKind::InvalidData, "未知消息类型")), } } } impl Encoder<Message> for MultiCodec { type Error = io::Error; fn encode(&mut self, item: Message, dst: &mut bytes::BytesMut) -> Result<(), Self::Error> { match item { Message::Text(s) => { dst.extend_from_slice(&[0x00]); let mut lines_codec = tokio_util::codec::LinesCodec::new(); lines_codec.encode(s, dst) } Message::Binary(b) => { dst.extend_from_slice(&[0x01]); let len = b.len() as u32; dst.extend_from_slice(&len.to_be_bytes()); dst.extend_from_slice(&b); Ok(()) } } } }
使用时直接用该Codec包装TcpStream,即可同时处理文本和二进制数据:
let mut framed = Framed::new(stream, MultiCodec); // 发送文本 framed.send(Message::Text("Hello client".to_string())).await?; // 发送二进制文件数据 let file_data = Bytes::from(std::fs::read("image.png")?); framed.send(Message::Binary(file_data)).await?;
方案2:用Arc共享TcpStream(不推荐)
若不想自定义Codec,可将TcpStream用Arc<Mutex>包裹,实现多Framed实例共享。但需注意同一时间只能有一个Framed执行读写,否则会出现数据混乱:
use std::sync::{Arc, Mutex}; let stream = Arc::new(Mutex::new(stream)); // 处理文本 let text_stream = Arc::clone(&stream); tokio::spawn(async move { let mut lines = Framed::new(text_stream.lock().unwrap(), LinesCodec::new()); lines.send("Text message").await.unwrap(); }); // 处理二进制 let binary_stream = Arc::clone(&stream); tokio::spawn(async move { let mut bytes_framed = Framed::new(binary_stream.lock().unwrap(), BytesCodec::new()); bytes_framed.send(Bytes::from_static(b"binary data")).await.unwrap(); });
这种方式需要手动同步读写逻辑,容易引发问题,优先推荐方案1。
二、优雅关闭TcpStream和Framed实例
优雅关闭需按以下步骤执行:
- 停止发送:调用Framed的
close()方法,关闭写半部分并发送结束信号(若Codec支持); - 处理剩余数据:等待Framed的
next()返回None,确保所有已接收数据处理完毕; - 关闭流:Drop Framed实例会自动关闭流,也可手动调用
shutdown()确保双工关闭。
示例代码(基于方案1):
let mut framed = Framed::new(stream, MultiCodec); // 触发关闭 framed.close().await?; // 等待剩余数据处理完成 while let Some(_msg) = framed.next().await {} // 手动关闭流(可选,drop framed也会自动执行) if let Some(stream) = framed.into_inner() { stream.shutdown(std::net::Shutdown::Both)?; }
若使用Arc
// 假设有AtomicBool类型的关闭标志,先通知所有子任务停止 let mut guard = stream.lock().unwrap(); guard.shutdown(std::net::Shutdown::Both)?;
内容的提问来源于stack exchange,提问作者realzhujunhao
相关产品推荐
相关产品推荐

