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

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)
        }
    }
}
关键改进点
  1. 统一超时控制:用单个带35秒超时的Context控制两个goroutine,避免重复的超时逻辑,确保超时信号能同时传递给接收和处理逻辑。
  2. 非阻塞Recv:将stream.Recv()放到独立goroutine中,通过channel传递结果,让主逻辑能响应超时信号,不再被Recv阻塞。
  3. 避免goroutine泄漏:通过defer group.Done()和Context的取消信号,确保无论超时还是正常完成,goroutine都能正确退出。
  4. 简化逻辑:移除了冗余的readWrite参数和循环发送channel的逻辑,单个缓冲channel足够传递一次决策。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 15:04:50