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

调用Kafka客户端Produce函数时出现context cancelled错误的原因排查

Kafka生产者推送记录时持续出现context cancelled错误排查

问题概述

使用github.com/twmb/franz-go/pkg/kgo包实现的Kafka生产者,向指定主题推送记录时持续触发context canceled错误,已通过Ping验证Broker连通性正常。

代码片段

ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
kafkaClient.Produce(ctx, record, func(record *kgo.Record, err error) {
        if err != nil {
            errorCallbackFn(err)
            return
        }
        successCallbackFn()
    })

已尝试排查步骤

  • 验证Broker URL正确性,Ping结果正常,排除地址问题
  • 延长context超时时间,问题未解决

错误信息及栈追踪

错误信息:

context canceled

错误栈:

kgo.(*Client).finishRecordPromise (producer.go:484)
github.com/twmb/franz-go/pkg/kgo kgo.(*producer).finishPromises
(producer.go:458) github.com/twmb/franz-go/pkg/kgo
kgo.(*producer).promiseBatch.func1 (producer.go:438)
github.com/twmb/franz-go/pkg/kgo runtime.goexit (asm_arm64.s:1172)
runtime

  • Async Stack Trace kgo.(*producer).promiseBatch (producer.go:438) github.com/twmb/franz-go/pkg/kgo

问题原因及解决方案

核心原因

kgo.Client.Produce是异步方法,调用后会立即返回,而当前代码中defer cancel()会在函数退出时立即取消context。此时异步回调尚未执行完成,Kafka客户端会因context被取消而触发context canceled错误,和设置的超时时间无关。

解决方案

方案1:使用WaitGroup等待回调执行完成

通过同步等待确保回调执行完毕后再取消context:

import "sync"

// ...

var wg sync.WaitGroup
wg.Add(1)
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()

kafkaClient.Produce(ctx, record, func(record *kgo.Record, err error) {
    defer wg.Done() // 回调完成后标记WaitGroup
    if err != nil {
        errorCallbackFn(err)
        return
    }
    successCallbackFn()
})

wg.Wait() // 等待回调执行完成

方案2:使用同步发送方法ProduceSync

如果业务允许同步发送,直接使用ProduceSync方法,它会阻塞直到消息发送完成或超时:

ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()

// 同步发送消息
results := kafkaClient.ProduceSync(ctx, record)
for _, res := range results {
    if res.Err != nil {
        errorCallbackFn(res.Err)
        continue
    }
    successCallbackFn()
}

额外排查点

如果上述方案仍未解决,可进一步检查:

  • Kafka主题是否存在,且Broker开启了自动创建主题的配置(auto.create.topics.enable=true)
  • 生产者配置的acks参数是否过高(比如acks=all但集群副本同步异常)
  • 网络是否存在隐性丢包(可通过抓包验证Kafka请求响应链路)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 01:50:23