You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在同一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实例

优雅关闭需按以下步骤执行:

  1. 停止发送:调用Framed的close()方法,关闭写半部分并发送结束信号(若Codec支持);
  2. 处理剩余数据:等待Framed的next()返回None,确保所有已接收数据处理完毕;
  3. 关闭流: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 15:10:16