返回文章列表

TCP 没有消息边界:用 Rust 写一个靠谱的拆包层

505·4 分钟阅读
Rust微服务

TCP 很可靠,但它只保证字节按顺序到达,不保证你写出去的一次就是对方读到的一次。

很多网络 demo 都会写成这样:客户端 write_all(b"hello"),服务端 read(&mut buf),然后打印 "hello"。本地跑得很好,换成真实流量就开始出现怪问题:

收到 hel
下一次收到 lo
再下一次收到 worldping

这不是 Tokio 的 bug,也不是网络不可靠。恰恰相反,这是 TCP 正常工作的样子。

TCP 是字节流,不是消息流。应用层协议必须自己定义消息边界。

前面 从 TCP 到 QUIC 里只提了一句“半包粘包”。这篇把这个问题单独拆开:一个能跑的 TCP demo,怎么变成一个不会乱包、不会被大包打爆、还能接进业务代码的协议层。

本文代码环境:

# Cargo.toml
[dependencies]
bytes = "1"
futures-util = { version = "0.3", features = ["sink"] }
tokio = { version = "1", features = ["full"] }
tokio-util = { version = "0.7", features = ["codec"] }

最大的误解:一次 read 就是一条消息

最常见的错误写法是这样:

use tokio::{io::AsyncReadExt, net::TcpStream};
 
async fn read_message(mut stream: TcpStream) -> std::io::Result<Vec<u8>> {
    let mut buf = vec![0; 1024];
    let n = stream.read(&mut buf).await?;
    buf.truncate(n);
    Ok(buf)
}

这段代码的问题不在语法,而在假设:它把“一次 read”当成了“一条消息”。

TCP 没有这个承诺。对端连续写两次:

write_all("hello")
write_all("world")

接收方可能看到很多种结果:

"helloworld"
"hel" + "loworld"
"hello" + "world"
"hell" + "oworld"

这些结果都合法。

所以讨论 TCP 拆包之前,先把这句话钉住:网络层只给你字节流,消息边界是应用协议的一部分。


换行符协议很直觉,但边界太脆

最容易想到的方案是用换行分隔:

SET name ferris\n
GET name\n
PING\n

这个设计适合人类调试,Redis RESP、HTTP/1 的很多部分都有类似影子。Rust 里也能用 LinesCodec 很快跑起来。

但它有几个限制:

问题 影响
payload 里不能随便出现换行 二进制数据很难处理
解析时要扫描分隔符 长消息会增加扫描成本
最大行长度必须限制 否则一行超大数据能撑爆内存
协议升级空间有限 后面加压缩、加 flags 会别扭

命令行调试协议可以用换行,服务间二进制协议优先用长度前缀。

换行协议不是不专业,只是它适合的场景不同。你要传 JSON 日志、聊天文本、简单命令,它很好。你要传 protobuf、压缩数据、图片块、批量消息,长度前缀更稳。


长度前缀:朴素,但最耐用

长度前缀协议的结构很简单:

┌──────────────┬──────────────────────┐
│ 4 字节长度    │ payload              │
│ u32 big-endian│ N bytes              │
└──────────────┴──────────────────────┘

发送方先写 payload 长度,再写 payload。接收方先攒够 4 字节,读出长度,再继续攒够完整 payload。

编码函数不用复杂:

use bytes::{BufMut, BytesMut};
 
const MAX_FRAME: usize = 1024 * 1024;
 
fn encode_frame(payload: &[u8], dst: &mut BytesMut) -> std::io::Result<()> {
    if payload.len() > MAX_FRAME {
        return Err(std::io::ErrorKind::InvalidData.into());
    }
 
    dst.reserve(4 + payload.len());
    dst.put_u32(payload.len() as u32);
    dst.extend_from_slice(payload);
    Ok(())
}

这里最重要的不是 put_u32,而是 MAX_FRAME

没有最大帧限制的长度前缀协议,等于把内存分配权交给对端。对方发一个长度 4GB,你如果照着分配 buffer,服务就没了。

协议设计里,长度字段一定要配最大值。


解码真正麻烦:数据可能还没来齐

编码是一次性写出去,解码不是一次性读回来。

解码器要能处理三种状态:

buffer 不足 4 字节           → 继续等
buffer 有长度,但 payload 不够 → 继续等
buffer 已经有完整帧          → 切出一条消息

BytesMut 写出来是这样:

use bytes::{Buf, BytesMut};
 
fn decode_frame(src: &mut BytesMut) -> std::io::Result<Option<BytesMut>> {
    if src.len() < 4 {
        return Ok(None);
    }
 
    let len = u32::from_be_bytes(src[..4].try_into().unwrap()) as usize;
    if len > MAX_FRAME {
        return Err(std::io::ErrorKind::InvalidData.into());
    }
    if src.len() < 4 + len {
        return Ok(None);
    }
 
    src.advance(4);
    Ok(Some(src.split_to(len)))
}

这段代码有两个点值得注意。

advance(4) 是把长度头丢掉。split_to(len) 是把 payload 从 buffer 前面切出来,剩下的字节留给下一次解码。

split_to 不是把数据复制一份,它只是调整 buffer 视图,成本很低。这就是 bytes crate 适合做网络协议的原因。


手写循环很容易把状态写散

如果不用 codec,服务端读循环大概会变成这样:

use tokio::io::AsyncReadExt;
 
async fn read_loop(mut stream: tokio::net::TcpStream) -> std::io::Result<()> {
    let mut buf = BytesMut::with_capacity(4096);
    loop {
        let n = stream.read_buf(&mut buf).await?;
        if n == 0 {
            break;
        }
        while let Some(frame) = decode_frame(&mut buf)? {
            handle_frame(frame).await?;
        }
    }
    Ok(())
}

这段逻辑能工作,但业务一多就容易散:

  • 读 socket 的逻辑
  • 缓冲区管理
  • 拆包
  • 单帧处理
  • 错误处理
  • 写响应

全塞在一个循环里,后面加心跳、认证、压缩、trace,就会越来越难看。

更好的边界是:拆包层只负责字节和消息之间的转换,业务层只处理完整消息。


用 tokio-util 写成 Codec

Tokio 生态里,这个边界通常用 DecoderEncoder 表达。

先定义一个 codec:

use bytes::{Bytes, BytesMut};
use tokio_util::codec::{Decoder, Encoder};
 
struct FrameCodec;

解码就是把刚才的函数塞进 trait:

impl Decoder for FrameCodec {
    type Item = BytesMut;
    type Error = std::io::Error;
 
    fn decode(&mut self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
        decode_frame(src)
    }
}

编码也同理:

impl Encoder<Bytes> for FrameCodec {
    type Error = std::io::Error;
 
    fn encode(&mut self, item: Bytes, dst: &mut BytesMut) -> Result<(), Self::Error> {
        encode_frame(&item, dst)
    }
}

有了 codec,就可以把 TcpStream 变成一个 Stream + Sink

use futures_util::{SinkExt, StreamExt};
use tokio_util::codec::Framed;
 
async fn handle_socket(stream: tokio::net::TcpStream) -> std::io::Result<()> {
    let mut framed = Framed::new(stream, FrameCodec);
    while let Some(frame) = framed.next().await {
        let frame = frame?;
        framed.send(frame.freeze()).await?;
    }
    Ok(())
}

这里业务看到的已经不是“任意长度的 TCP 字节片段”,而是一条条完整 frame。

这就是 Framed 的价值:把字节流的状态机藏在 codec 里,把业务逻辑还给业务层。


能用 LengthDelimitedCodec 就别急着手写

如果协议就是标准的长度前缀,tokio-util 已经给了现成实现:LengthDelimitedCodec

最简单的 echo server 可以这么写:

use futures_util::{SinkExt, StreamExt};
use tokio_util::codec::{Framed, LengthDelimitedCodec};
 
async fn echo(stream: tokio::net::TcpStream) -> std::io::Result<()> {
    let codec = LengthDelimitedCodec::builder()
        .max_frame_length(1024 * 1024)
        .new_codec();
    let mut framed = Framed::new(stream, codec);
 
    while let Some(frame) = framed.next().await {
        framed.send(frame?.freeze()).await?;
    }
    Ok(())
}

我个人的偏好是:

场景 选择
只有长度前缀 LengthDelimitedCodec
有 magic/version/flags/checksum 自定义 Decoder + Encoder
文本命令协议 LinesCodec 或自己限制行长
复杂协议、多阶段握手 手写状态机,但边界仍然保持 codec 化

不要为了“可控”一上来就手写所有东西。LengthDelimitedCodec 已经处理了缓冲、半包、粘包和最大帧限制,能覆盖很多服务间协议的第一版。


协议头不要一开始就设计太大

如果需要比长度前缀多一点信息,可以加一个很小的头:

magic  version  flags  length
2B     1B       1B     4B

含义可以这样定:

字段 作用
magic 快速判断是不是自己的协议
version 给未来升级留入口
flags 标记压缩、加密、消息类型
length payload 长度

但别一上来就把协议头设计成几十个字段。

协议字段加上去容易,删掉很难。特别是服务间协议,一旦多语言客户端接入,兼容性就会变成长期成本。

第一版协议只放真正影响解析边界的字段。 业务字段放 payload 里,用 protobuf、JSON、MessagePack 都可以。


拆包层必须处理恶意输入

很多 demo 只考虑“正常客户端”,生产服务不能这样。

拆包层至少要处理这些情况:

输入 应对
长度超过上限 立即报错并断开
长度为 0 明确允许还是拒绝
只发头不发 body 依赖读超时关闭连接
持续发垃圾字节 尽快失败,不进入业务层
payload 编码错误 协议层成功,业务解码失败

这里有个边界很重要:拆包层不应该理解业务字段。

它只需要保证“每次交给业务层的都是一条完整 payload”。payload 里面是不是合法 JSON、protobuf 字段缺不缺,那是下一层的事。

这个边界分清楚,错误处理也会清楚:

TCP read error       → 连接层
frame too large      → 协议层
invalid protobuf     → 序列化层
unknown command      → 业务层

结论

Rust 写 TCP 服务,真正难的不是 TcpListener::bind,而是你有没有认真处理消息边界。

可以直接记这几条:

  • TCP 是字节流,不是消息流
  • 一次 read 不等于一条消息
  • 长度前缀要配最大帧限制
  • BytesMut 适合做拆包缓冲,split_to 能低成本切帧
  • 简单长度前缀优先用 LengthDelimitedCodec
  • 自定义协议也应该实现成 Decoder + Encoder
  • 拆包层只处理边界,不处理业务语义

只要用 TCP 写自定义协议,就必须先设计分帧。

分帧不是优化项,也不是“后面再补”的细节。它是协议能不能成立的地基。没有这层,后面的认证、心跳、压缩、重试、流控,都会建在一个会随机错位的字节流上。