Go项目使用kafka-go客户端如何实现Kafka断连自动重连
kafka-go 客户端断连自动重连实现方案
github.com/segmentio/kafka-go 本身内置了基础的断连重试能力,绝大多数网络波动、broker临时重启、主从切换场景下,只要初始化参数配置正确,不需要额外写逻辑就能自动恢复连接。只有长时间网络中断、集群整体故障恢复这类极端场景,才需要加外层兜底逻辑。
优先配置内置重试参数
大部分人遇到的自动重连失效,都是初始化时没配全重试相关参数,用了默认的保守配置导致的。
消费者(Reader)核心配置
import ( "time" "github.com/segmentio/kafka-go" ) reader := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"broker1:9092", "broker2:9092", "broker3:9092"}, Topic: "your_topic", GroupID: "your_consumer_group", // 重试核心参数 MaxAttempts: 10, // 单次操作最大重试次数,默认值3,网络环境差可以适当调大 ReadBackoffMin: 100 * time.Millisecond, // 读失败最小退避间隔 ReadBackoffMax: 5 * time.Second, // 读失败最大退避间隔,内置指数退避逻辑 SessionTimeout: 30 * time.Second, // 消费者组会话超时,超过时长无心跳判定为断连 RebalanceTimeout: 30 * time.Second, Dialer: &kafka.Dialer{ Timeout: 10 * time.Second, // 建连超时 KeepAlive: 30 * time.Second, // 开启TCP KeepAlive,提前检测半开死连接 DualStack: true, }, })
生产者(Writer)核心配置
writer := &kafka.Writer{ Addr: kafka.TCP("broker1:9092", "broker2:9092", "broker3:9092"), Topic: "your_topic", Balancer: &kafka.LeastBytes{}, RequiredAcks: kafka.RequireOne, // 重试核心参数 MaxAttempts: 10, // 写操作最大重试次数 WriteBackoffMin: 100 * time.Millisecond, WriteBackoffMax: 10 * time.Second, BatchSize: 100, BatchTimeout: 10 * time.Millisecond, }
配完以上参数,几秒内的临时断连都会被内置逻辑自动处理,业务层感知不到错误。
极端故障场景兜底重连
如果断连时长超过内置重试的最大次数,内置逻辑会直接返回错误,这时候需要在外层加兜底,避免服务直接崩溃。
消费者兜底逻辑
消费者不要因为单次读报错就直接退出循环,先判断错误类型,不可恢复错误直接终止进程,可重试错误做退避后重建Reader实例即可:
import ( "context" "errors" "io" "log" "net" "time" "github.com/segmentio/kafka-go" ) // 指数退避控制,初始1s,最大30s var backoff = 1 * time.Second const maxBackoff = 30 * time.Second func nextBackoff() time.Duration { current := backoff backoff *= 2 if backoff > maxBackoff { backoff = maxBackoff } return current } func resetBackoff() { backoff = 1 * time.Second } func isRetryableErr(err error) bool { // 网络错误、超时、EOF都属于可重试错误 var netErr net.Error if errors.As(err, &netErr) { return true } var kafkaErr kafka.Error if errors.As(err, &kafkaErr) && kafkaErr.Temporary() { return true } if errors.Is(err, io.EOF) || errors.Is(err, context.DeadlineExceeded) { return true } return false } func StartConsumer() { for { reader := kafka.NewReader(/* 传入上面定义好的ReaderConfig */) for { ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) msg, err := reader.ReadMessage(ctx) cancel() if err != nil { if !isRetryableErr(err) { // 认证失败、权限不足、Topic不存在这类不可恢复错误,直接打日志退出 log.Fatalf("消费者遇到不可恢复错误,终止运行: %v", err) } log.Printf("连接Kafka异常,%s后准备重连: %v", nextBackoff(), err) time.Sleep(nextBackoff()) break } resetBackoff() // 这里写你的消息处理逻辑 handleMsg(msg) } // 销毁旧连接,避免脏状态影响 _ = reader.Close() } }
生产者兜底逻辑
生产者是无状态的,不需要重建实例,只要在写消息失败时判断错误类型,做退避重试即可:
func WriteMsg(ctx context.Context, msg kafka.Message) error { for { err := writer.WriteMessages(ctx, msg) if err == nil { resetBackoff() return nil } if !isRetryableErr(err) { return err } log.Printf("写Kafka失败,%s后重试: %v", nextBackoff(), err) select { case <-ctx.Done(): return ctx.Err() case <-time.After(nextBackoff()): } } }
注意事项
- 重连必须加退避逻辑,禁止出错后立刻死循环重连,否则网络故障时会产生大量无效连接,直接打满Kafka集群的连接数阈值
- 一定要区分可重试错误和不可重试错误,配置类错误重试没有意义,及时暴露问题比无效重试更重要
- Dialer必须开启KeepAlive,否则TCP半连接场景下客户端无法及时感知连接断开,会卡住很久才触发超时
- 不要把SessionTimeout设得太短,否则业务GC停顿、普通网络延迟都可能触发消费者组重平衡,反而加剧服务不稳定
内容的提问来源于stack exchange,提问作者young
相关产品推荐
相关产品推荐

