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

如何在Go的kgo.Produce回调中返回错误并终止sendAll函数

解决kgo.Produce异步回调中返回错误并终止sendAll函数的问题

kgo.Produce是异步调用的,回调函数会在后台goroutine中执行,没法直接在外层sendAll函数里返回错误或终止循环。下面提供两种可行的解决方案:

方案一:使用同步发送接口ProduceSync

如果业务场景允许同步等待发送结果,直接用ProduceSync可以简化逻辑,出错时直接返回并终止循环:

func (k *buffer) sendAll() error {
    for _, record := range k.buffer {
        // 同步等待发送结果,First()获取第一个错误(如果有)
        _, err := k.client.ProduceSync(ctx, record).First()
        if err != nil {
            fmt.Printf("record had a produce error: %v\n", err)
            return err // 直接返回错误,循环终止
        }
    }
    return nil
}

方案二:异步发送结合Context与错误通道

如果必须保留异步发送的模式,需要通过context取消后续发送,并利用通道传递错误:

func (k *buffer) sendAll() error {
    errCh := make(chan error, 1) // 缓冲通道避免回调阻塞
    ctx, cancel := context.WithCancel(ctx)
    defer cancel() // 函数退出时确保取消上下文

    for _, record := range k.buffer {
        // 先检查是否已有错误,有则直接返回
        select {
        case err := <-errCh:
            return err
        default:
        }

        k.client.Produce(ctx, &record, func(r *kgo.Record, err error) {
            if err != nil {
                fmt.Printf("record had a produce error: %v\n", err)
                // 尝试发送错误到通道,非阻塞避免回调挂起
                select {
                case errCh <- err:
                default:
                }
                cancel() // 取消上下文,阻止后续发送操作
            }
        })
    }

    // 等待所有发送完成或收到错误
    select {
    case err := <-errCh:
        return err
    case <-ctx.Done():
        // 若上下文被取消,优先返回发送错误
        if err := <-errCh; err != nil {
            return err
        }
        return ctx.Err()
    }
}

注意事项

  • 代码中的k.client是你初始化的kgo生产者客户端,需确保正确传入。
  • 异步方案中,错误通道设置缓冲为1即可,因为我们只需要捕获第一个错误就终止流程。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 12:39:27