【发布时间】:2021-03-05 09:01:27
【问题描述】:
我有一个“代理”,它将二进制文件解析到缓冲区中,每当该缓冲区被填满时,通过 protobuf 消息将其发送到服务器,然后继续进行下一个二进制解析块,然后再次发送,等等。
在服务器上,我使用简单的net/conn 包来侦听代理连接并在while-for 循环中将其读取到缓冲区中。
当解析完成代理端时,它会在 protobuf 消息中发送一个terminate bool,表明这是最后一条消息,服务器可以继续处理接收到的全部数据。
但是,如果我将调试打印留在发送方,这会正常工作,这会使终端打印显着减慢通过 connection.Write() 发送后续 protobuf 消息的时间间隔。
如果我取消注释此记录器,那么它发送消息的速度太快,服务器处理的第一个传入消息是包含terminate 标志的一个数据包,例如,它没有收到实际有效负载,而是立即收到最后一条消息。
我知道 TCP 并没有真正区分不同的 []byte 数据包,这很可能是导致这种行为的原因。有没有更好的方法可以做到这一点,还有其他选择吗?
伪代码代理端:
buffer := make([]byte, 1024)
for {
n, ioErr := reader.Read(buffer)
if ioErr == io.EOF {
isPayloadFinal = true
// Create protobuf message
terminalMessage, err := CreateMessage_FilePackage(
2234,
protobuf.MessageType_PACKAGE,
make([]byte, 1),
isPayloadFinal,
)
// Send terminate message
sendProtoBufMessage(connection, terminalMessage)
break
}
// Create regular protobuf message
message, err := CreateMessage_FilePackage(
2234,
protobuf.MessageType_PACKAGE,
(buffer)[:n],
isPayloadFinal)
sendProtoBufMessage(connection, message)
}
伪代码服务器端:
buffer := make([]byte, 2048)
//var protoMessage protoBufMessage
for artifactReceived != true {
connection.SetReadDeadline(time.Now().Add(timeoutDuration))
n, _ := connection.Read(buffer)
decodedMessage := &protobuf.FileMessage{}
if err := proto.Unmarshal(buffer[:n], decodedMessage); err != nil {
log.Err(err).Msg("Error during unmarshalling")
}
if isPackageFinal := decodedMessage.GetIsTerminated(); isPackageFinal == true {
artifactReceived = true
log.Info().Msg("Artifact fully received")
/* Do stuff here */
break
}
// Handle partially arrived bytestream
handleProtoPackage(packageMessage, artifactPath)
} else {
fmt.Println("INVALID PROTOBUF MESSAGE")
}
}
以及供参考的proto文件:
message FilePackage{
int32 id = 1;
MessageType msgType = 2;
bytes payload = 3;
bool isTerminated = 4;
}
【问题讨论】:
标签: go tcp buffer protocol-buffers