Go gRPC双向流中Stream.Recv()阻塞致Context取消不触发的问题
问题分析
你的核心问题是gRPC的stream.Recv()是阻塞调用,无法直接被自定义Context的超时信号中断——它只响应流自身stream.Context()的取消(比如客户端主动断开)。你当前的代码尝试用Context超时,但没有将Recv的执行与超时信号绑定,导致超时后Recv仍然阻塞,无法退出当前迭代。
当前代码的几个关键问题:
receiveDriverDecision中用select包裹ctxchild.Done()和default,但default分支直接执行stream.Recv(),一旦进入分支就会阻塞,无法响应后续的超时信号。processDriverDecision中的time.After()触发cancel(),但receiveDriverDecision并未监听这个被取消的ctx,而是监听stream.Context()的超时,导致取消信号无法传递到Recv的goroutine。- 两个goroutine的超时逻辑重复,导致控制混乱。
解决方案
正确的做法是将stream.Recv()放到独立goroutine中,通过channel传递结果,然后在主逻辑中用select同时监听超时信号和Recv的结果,这样就能在超时后立即中断等待,退出当前迭代。
以下是修改后的核心代码:
1. 重构receiveDriverDecision函数
将Recv的执行与超时控制绑定,通过channel传递结果,避免阻塞:
func (app *driverServer) receiveDriverDecision(ctx context.Context, stream pb.DriverService_StreamDispatchOrdersServer, group *sync.WaitGroup, receivedChannel chan<- Received, orderID primitive.ObjectID, driverID primitive.ObjectID) { defer group.Done() recvChan := make(chan *pb.DriverDecisionRequest, 1) errChan := make(chan error, 1) // 在goroutine中执行Recv,避免阻塞主逻辑 go func() { ordr, err := stream.Recv() if err != nil { errChan <- err return } recvChan <- ordr.GetRequest() }() select { case <-ctx.Done(): log.Printf("RECEIVE DRIVER DECISION :: TIMED OUT OR CANCELLED") return case req := <-recvChan: // 处理收到的请求 recvd := Received{ Order: mappers.RPCToOrder(req.GetOrder()), DriverID: req.GetDriverId(), Decision: req.GetAcceptDecision(), } receivedChannel <- recvd log.Printf("RECEIVED DECISION WAS SENT") case err := <-errChan: if err == io.EOF { log.Printf("RECEIVE DRIVER DECISION :: END OF STREAM") } else { log.Printf("RECEIVE DRIVER DECISION :: ERROR: %v", err) } return } }
2. 简化入口循环的超时控制
在入口循环中创建带35秒超时的Context,传递给两个goroutine,避免重复的超时逻辑:
// 替换原来的ctx创建逻辑 ctx, cancel := context.WithTimeout(context.Background(), 35*time.Second) defer cancel() wg := new(sync.WaitGroup) wg.Add(2) received := make(chan Received, 1) // 缓冲channel避免阻塞 defer close(received) // 启动goroutine,传递带超时的ctx go app.receiveDriverDecision(ctx, stream, wg, received, order.ID, dID) go app.processDriverDecision(ctx, stream, order.ID, dID, order.DeliveryPhase, wg, received) wg.Wait()
3. 修改processDriverDecision函数
移除重复的time.After()超时逻辑,直接监听传入的ctx.Done(),同时处理收到的决策:
func (app *driverServer) processDriverDecision(ctx context.Context, stream pb.DriverService_StreamDispatchOrdersServer, orderID primitive.ObjectID, driverID primitive.ObjectID, phase models.OrderPhase, group *sync.WaitGroup, receivedChannel <-chan Received) { defer group.Done() conn, err := app.GetOrderClient() if err != nil { log.Printf("PROCESS DRIVER DECISION :: FAILED TO GET ORDER CLIENT: %v", err) return } orderClient := pb.NewOrderServiceClient(conn) select { case msg := <-receivedChannel: // 原有的决策处理逻辑保持不变 estimatedEarnings := (msg.Order.DriverPool.EstimatedEarnings + msg.Order.DriverPool.TipAmount) / 2 log.Printf("PROCESS DRIVER DECISION :: RECEIVED MESSAGE: %v", msg) orderAssigned, err := orderClient.GetOrderFromAssignedRecords(stream.Context(), &pb.GetOrderFromAssignedRequest{ OrderId: msg.Order.ID.Hex(), DeliveryPhase: pb.OrderPhase(msg.Order.DeliveryPhase), }) if err != nil { log.Printf("PROCESS DRIVER DECISION :: FAILED TO GET ORDER FROM ASSIGNED RECORDS: %v", err) return } if !orderAssigned.GetAssigned() { if msg.Decision { // 处理接单逻辑...(原代码保持不变) } } else { // 发送订单已被分配的响应...(原代码保持不变) } case <-ctx.Done(): log.Printf("PROCESS DRIVER DECISION :: TIMED OUT, NO RESPONSE FROM DRIVER") // 发送超时提示给客户端 err := stream.Send(&pb.DispatchOrderResponse{ Resp: &pb.DispatchOrderResponse_AcceptedResponse{ AcceptedResponse: &pb.AcceptedOrderResponse{ AssignedSuccessfully: false, Response: &pb.OrderPromptResponse{ Routes: []*pb.Route{}, EstimatedEarnings: 0, Distance: 0, Minutes: 0, }, }, }, }) if err != nil { log.Printf("PROCESS DRIVER DECISION :: FAILED TO SEND TIMEOUT RESPONSE: %v", err) } } }
关键改进点
- 统一超时控制:用单个带35秒超时的Context控制两个goroutine,避免重复的超时逻辑,确保超时信号能同时传递给接收和处理逻辑。
- 非阻塞Recv:将
stream.Recv()放到独立goroutine中,通过channel传递结果,让主逻辑能响应超时信号,不再被Recv阻塞。 - 避免goroutine泄漏:通过
defer group.Done()和Context的取消信号,确保无论超时还是正常完成,goroutine都能正确退出。 - 简化逻辑:移除了冗余的
readWrite参数和循环发送channel的逻辑,单个缓冲channel足够传递一次决策。
内容的提问来源于stack exchange,提问作者Patrick Cockrill
相关产品推荐
相关产品推荐

