【发布时间】:2020-06-05 04:09:15
【问题描述】:
对不起,如果这是一个菜鸟问题,我是 grpc 的服务器端流的新手。
我现在在服务器上的一个流到客户端的函数中拥有什么
req, err := http.NewRequest("GET", actualURL, nil)
//跳过一些行// res, _ := http.DefaultClient.Do(req)
// closing body
defer res.Body.Close()
body, err := ioutil.ReadAll(res.Body)
//跳过一些行//
// unmarshaling the xml data received from the GET request done above
xmlDataErr := xml.Unmarshal(body, &ArrData)
// creating a variable of a custom type and creating object of specific struct Flight
var flightArrData ArrFlights
currArrData := flightArrData.Arr.Body.Flight
log.Println("计算要发送到客户端的总行数", len(currArrData))
// 循环响应并将其发送到客户端
对于 i := range currArrData {
farr := &pb.FlightArrivalsResponse{
ArrivalResponse: &pb.ArrivalFlightData{
Hapt: currArrData[i].HApt,
Fltnr: currArrData[i].Fltnr,
Sdt: currArrData[i].Sdt,
Acreg: currArrData[i].Acreg,
Park: currArrData[i].Park,
EstD: currArrData[i].EstD,
Gate: currArrData[i].Gate,
AblkD: currArrData[i].AblkD,
ActD: currArrData[i].ActD,
Callsign: currArrData[i].Callsign,
},
}
senderr := stream.Send(farr)
// skipping some lines //
}
// 完成后返回 nil 返回零 }
我面临的问题 在客户端,我有接收此响应的函数,休眠 n 分钟并再次请求响应。
客户端在第一次调用时确实得到了预期的响应,但是对于每个后续调用都会发生一些奇怪的事情,这是我当前的问题,我尝试以下面每个后续调用的形式来说明:
从客户端调用 1 到服务器 --> 服务器返回 200 行
客户端睡眠 n 分钟
从客户端调用 2 到服务器 --> 服务器返回 400 行!! (基本上每行都是 200 + 200 的两倍)
客户端睡眠 n 分钟
从客户端调用 3 到服务器 --> 服务器返回 600 行!! (200+200+200)
总结
我确实会在客户端检查错误 == io.EOF,而现在要停止这种从服务器到客户端的响应堆叠的唯一方法是停止服务器并重新启动。
我不确定我在这里遗漏了什么,以确保我只发送我从 GET 请求中收到的实际和准确的响应。非常感谢任何提示。
更多信息
来自 proto 文件中 gRPC protobuffer def 的部分
rpc GetFlightArrivals (FlightArrivalsRequestURL) returns (stream FlightArrivalsResponse) {
};
上述gRPC服务端impl的完整代码
func (s *Server) GetFlightArrivals(url *pb.FlightArrivalsRequestURL, stream pb.Flight_GetFlightArrivalsServer) error {
cData := ch.ConfigProcessor() // fetches the initial part of the URL from config file
if url.ArrivalURL == "" {
actualURL = cData.FURL + "/arr/all" // adding the arrival endpoint
} else {
actualURL = url.ArrivalURL
}
// build new request to get arrival data
req, err := http.NewRequest("GET", actualURL, nil)
if err != nil {
log.Fatalln("Recheck URL or connectivity, failed to make REST call")
}
req.Header.Add("Cache-Control", "no-cache")
req.Header.Add("Accept", "text/plain")
req.Header.Add("Connection", "keep-alive")
req.Header.Add("app_id", cData.AppID)
req.Header.Add("app_key", cData.AppKey)
req.Header.Add("Content-Type", "application/xml")
res, _ := http.DefaultClient.Do(req)
// closing body
defer res.Body.Close()
body, err := ioutil.ReadAll(res.Body)
if err != nil {
log.Fatalln("Failed to get any response")
return err
}
// unmarshaling the xml data
xmlDataErr := xml.Unmarshal(body, &flightArrData)
if xmlDataErr != nil {
log.Fatalln("Failed to unmarshal arrival xml data see error, ", xmlDataErr)
return xmlDataErr
}
currArrData := flightArrData.Arr.Body.Flight
log.Println("Counting total arrivals in Finland", len(currArrData))
log.Println("Starting FlightDeparturesResponse for client")
for i := range currArrData {
farr := &pb.FlightArrivalsResponse{
ArrivalResponse: &pb.ArrivalFlightData{
Hapt: currArrData[i].HApt,
Fltnr: currArrData[i].Fltnr,
Sdt: currArrData[i].Sdt,
Acreg: currArrData[i].Acreg,
Park: currArrData[i].Park,
EstD: currArrData[i].EstD,
Gate: currArrData[i].Gate,
AblkD: currArrData[i].AblkD,
ActD: currArrData[i].ActD,
Callsign: currArrData[i].Callsign,
},
}
senderr := stream.Send(farr)
if senderr != nil {
log.Fatalln("Failed to stream arrival response to the client, see error ", senderr)
return senderr
}
}
currArrData = nil
log.Println("Attempting to empty the arrival data")
return nil
}
检查上述 gRPC impl 的测试用例
我只是启动一个测试 grpc 服务器并在这个测试用例中调用该 grpc。客户端的实现与接收数据相同。
func TestGetFlightArrivals(t *testing.T) {
const addr = "localhost:50051"
conn, err := grpc.Dial(addr, grpc.WithInsecure())
if err != nil {
t.Fatalf("Did not connect: #{err}")
}
defer conn.Close()
f := pb.NewFlightClient(conn)
t.Run("GetFlightArrivals", func(t *testing.T) {
var res, err = f.GetFlightArrivals(context.Background(), &pb.FlightArrivalsRequestURL{ArrivalURL: ""})
if err != nil {
t.Error("Failed to make the REST call", err)
}
for {
msg, err := res.Recv()
if err == io.EOF {
t.Log("Finished reading all the message")
break
}
if err != nil {
t.Error("Failed to receive response")
}
t.Log("Message from the server", msg.GetArrivalResponse())
}
})
}
【问题讨论】:
-
请发帖minimal reproducible example。目前,您的示例代码中有很大的漏洞(例如,您解组为
ArrData,但从未使用过该数据),这使得评论变得困难。 -
@brits 我会尽快回复。
-
您的解释还不够,但我敢打赌您使用的是指针而不是破坏它。您可能会一次又一次地使用相同的指针,重复您的行。
-
我已经用完整的代码和一个调用该 impl 的测试用例编辑了我的问题,这就是您在附加信息方面所寻找的吗?
-
@rubens21 我不确定您对复制的确切含义,我希望在服务器上发送和在客户端接收应该清空服务器端的堆栈。我现在在我的问题中添加了更多信息。
标签: go server client grpc grpc-go