调用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
相关产品推荐
相关产品推荐

