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

TCP大数据包偶发截断问题求助:60KB数据传输异常排查

TCP大数据传输异常排查求助

我搭建了TCP服务端与客户端用于测试数据包长度,设置上限为64KB。传输小数据时一切正常,但传输60000字节的大数据时,有时正常有时出现数据截断,还会显示大量0值。附上相关代码,恳请各位提供排查建议,非常感谢。

1. Client代码

package base

import (
    "fmt"
    "net"
)

type TcpClient struct {
    conn *Connection
}

func NewClient() *TcpClient {
    return &TcpClient{}
}

func (t *TcpClient) Connect(address string) error {
    cn, err := net.Dial("tcp", address)
    if err != nil {
        fmt.Printf("client dial failed|err:%v\n", err)
        return err
    }
    t.conn = NewConnection(cn)
    t.conn.Run()

    return nil
}

func (t *TcpClient) Write(msg *Message) {
    t.conn.writeQ <- msg
}

2. Server代码

package base

import (
    "fmt"
    "net"
)

type TcpServer struct {
    listener net.Listener
}

func NewTcpServer() *TcpServer {
    return &TcpServer{}
}

func (t *TcpServer) Listen(address string) error {
    lis, err := net.Listen("tcp", address)
    if err != nil {
        fmt.Printf("listen failed|err:%v\n", err)
        return err
    }

    t.listener = lis
    go t.acceptConn()

    return err
}

func (t *TcpServer) acceptConn() {
    for {
        conn, err := t.listener.Accept()
        if err != nil {
            fmt.Printf("accept failed|err:%v\n", err)
            return
        }

        go t.handleConn(conn)
    }
}

func (t *TcpServer) handleConn(c net.Conn) {
    newConn := NewConnection(c)
    newConn.Run()
}

3. Connection代码

package base

import (
    "bufio"
    "encoding/binary"
    "fmt"
    "net"
)

type Connection struct {
    conn   net.Conn
    writeQ chan *Message
    writer *bufio.Writer
    reader *bufio.Reader
}

func NewConnection(conn net.Conn) *Connection {
    return &Connection{
        conn:   conn,
        writeQ: make(chan *Message, 100),
        writer: bufio.NewWriterSize(conn, 1024*64),
        reader: bufio.NewReaderSize(conn, 1024*64),
    }
}

func (c *Connection) Run() {
    go c.loopRead()
    go c.loopWrite()
}

func (c *Connection) loopRead() {
    for {
        header, err := c.reader.Peek(HeaderLength)
        if err != nil {
            fmt.Printf("peek header failed|err:%v\n", err)
            return
        }

        dataLength := binary.BigEndian.Uint32(header[0:])
        fmt.Printf("read data to read|header:%v|length:%d\n", header, dataLength)

        msgBuff := make([]byte, HeaderLength+dataLength)
        if _, err = c.reader.Read(msgBuff); err != nil {
            fmt.Printf("read msg buff failed|err:%v\n", err)
            return
        }

        fmt.Printf("read total length:%d|data:%v\n", len(msgBuff), msgBuff[HeaderLength:])
    }
}

func (c *Connection) loopWrite() {
    for {
        select {
        case msg := <-c.writeQ:
            var header = make([]byte, HeaderLength) //data Length
            length := len(msg.Data)
            binary.BigEndian.PutUint32(header[0:], uint32(length))
            fmt.Printf("write header:%v|length:%d\n", header, length)
            if _, err := c.writer.Write(header); err != nil {
                fmt.Printf("write header failed|err:%v\n", err)
                return
            }

            if length > 0 {
                if _, err := c.writer.Write(msg.Data); err != nil {
                    fmt.Printf("write body failed|err:%v\n", err)
                    return
                }
            }

            if err := c.writer.Flush(); err != nil {
                fmt.Printf("flush failed|err:%v\n", err)
                return
            }
        }
    }
}

4. 公共定义代码

package base

const (
    Addr         = "127.0.0.1:8998"
    HeaderLength = 4
)

type Message struct {
    Length uint32
    Data   []byte
}

5. 主函数代码

客户端主函数

func main() {
    c := base.NewClient()
    if err := c.Connect(base.Addr); err != nil {
        fmt.Printf("client connect failed|err:%v\n", err)
        os.Exit(1)
    }

    for {
        msg := &base.Message{}
        for i := 0; i < 60000; i++ {
            msg.Data = append(msg.Data, 'a')
        }
        c.Write(msg)
        time.Sleep(time.Second * 10)
    }
}

服务端主函数

func main() {
    server := base.NewTcpServer()
    if err := server.Listen(base.Addr); err != nil {
        fmt.Printf("listen failed|err:%v\n", err)
        os.Exit(1)
    }
    for {

    }
}

排查建议

  • 确保读取完整数据:当前代码依赖bufio.Reader.Read一次性读取全部所需字节,但TCP是流式协议,Read可能返回部分数据,未读取的字节会留在缓冲区或网络中,导致msgBuff未被覆盖的部分为初始0值。建议使用io.ReadFull封装读取逻辑,保证读取到指定长度的字节:

    import "io"
    
    func (c *Connection) readFull(buf []byte) (n int, err error) {
        return io.ReadFull(c.reader, buf)
    }
    

    替换loopRead中的Read调用为readFull,如果读取长度不足会直接返回错误,避免数据不完整。

  • 优化Header读取逻辑:使用Peek获取Header后,后续读取会包含已Peek的Header,虽然逻辑上能凑够总长度,但换成直接读取Header的方式更清晰可靠,避免缓冲区数据重复处理的潜在问题:

    func (c *Connection) loopRead() {
        for {
            header := make([]byte, HeaderLength)
            if _, err := c.readFull(header); err != nil {
                fmt.Printf("read header failed|err:%v\n", err)
                return
            }
    
            dataLength := binary.BigEndian.Uint32(header)
            fmt.Printf("read data to read|header:%v|length:%d\n", header, dataLength)
    
            msgBody := make([]byte, dataLength)
            if _, err := c.readFull(msgBody); err != nil {
                fmt.Printf("read msg body failed|err:%v\n", err)
                return
            }
    
            fmt.Printf("read total length:%d|data:%v\n", len(msgBody), msgBody)
        }
    }
    
  • 添加关键节点日志:在客户端写入、服务端读取时打印实际处理的字节数,比如客户端打印len(msg.Data),服务端打印每次readFull返回的字节数,帮助定位问题出在发送还是接收环节。

  • 验证消息构造正确性:确认客户端每次发送的msg.Data长度确实为60000,避免因append操作意外导致长度不符。

内容的提问来源于stack exchange,提问作者XjChen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 18:55:13