ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

kafka-go 使用指南:用 Go 标准库风格的 API 连接 Kafka

kafka-go 使用指南:用 Go 标准库风格的 API 连接 Kafka kafka-go 使用指南用 Go 标准库风格的 API 连接 Kafka【免费下载链接】opencloud️ OpenCloud is the open source platform for file management, sharing and collaboration. Simple and sovereign.项目地址: https://gitcode.com/GitHub_Trending/op/opencloudkafka-go 是 Segment 团队开源的纯 Go Kafka 客户端它同时提供贴近 Kafka 协议的底层ConnAPI 与面向业务的高层Reader/WriterAPI并全面拥抱 Go 的context与io.Reader等标准库习惯。本指南基于仓库中 vendored 的 kafka-go v0.4.51 源码vendor/github.com/segmentio/kafka-go见 go.mod及其 README 编写读完你可以掌握底层连接的生产/消费、Topic 管理、消费者组与偏移量提交、多分区写入、TLS/SASL 加密认证以及压缩与日志等完整实战能力。为什么需要 kafka-go在 kafka-go 出现之前Go 社区可用的 Kafka 客户端各有明显短板sarama最流行但文档匮乏API 直接暴露 Kafka 协议的底层概念不支持 Go 的context且所有值都以指针传递造成大量动态内存分配、更频繁的 GC 和更高的内存占用confluent-kafka-go是基于 cgo 对librdkafka的封装所有使用它的 Go 代码都被迫引入 C 库依赖虽然文档更好但同样不支持 Go contextgoka面向特定用法把 Kafka 当作服务间消息总线而非有序事件日志且底层仍依赖 sarama。kafka-go 的定位正是弥补这些缺口它提供低层与高层两套 API并在设计上镜像 Go 标准库的概念context、io.Reader、net.Conn式接口让开发者可以像使用标准库一样自然地接入 Kafka。同时它是纯 Go 实现无 cgo 依赖跨平台编译部署更简单。版本兼容性官方测试覆盖 Kafka 0.10.1.0 至 2.7.1更新的 Kafka 通常也能工作但新增的协议特性可能尚未实现要求 Go 1.15 及以上版本本项目仓库锁定的是 v0.4.51见 go.mod。Conn底层连接 APIConn类型是 kafka-go 的核心它包装一条到 Kafka 服务器的原始网络连接对外暴露低层 API。由于贴近协议Conn也是构建Reader等高层抽象的理想积木。生产消息使用kafka.DialLeader直接连接到指定 topic 分区的 leader然后写入消息topic : my-topic partition : 0 conn, err : kafka.DialLeader(context.Background(), tcp, localhost:9092, topic, partition) if err ! nil { log.Fatal(failed to dial leader:, err) } conn.SetWriteDeadline(time.Now().Add(10*time.Second)) _, err conn.WriteMessages( kafka.Message{Value: []byte(one!)}, kafka.Message{Value: []byte(two!)}, kafka.Message{Value: []byte(three!)}, ) if err ! nil { log.Fatal(failed to write messages:, err) } if err : conn.Close(); err ! nil { log.Fatal(failed to close writer:, err) }消费消息同样用DialLeader连到 leader通过ReadBatch拉取一批消息可指定最小/最大字节数再用类似io.Reader的方式循环读取topic : my-topic partition : 0 conn, err : kafka.DialLeader(context.Background(), tcp, localhost:9092, topic, partition) if err ! nil { log.Fatal(failed to dial leader:, err) } conn.SetReadDeadline(time.Now().Add(10*time.Second)) batch : conn.ReadBatch(10e3, 1e6) // fetch 10KB min, 1MB max b : make([]byte, 10e3) // 10KB max per message for { n, err : batch.Read(b) if err ! nil { break } fmt.Println(string(b[:n])) } if err : batch.Close(); err ! nil { log.Fatal(failed to close batch:, err) } if err : conn.Close(); err ! nil { log.Fatal(failed to close connection:, err) }ReadBatch(minBytes, maxBytes)的两个参数分别向 broker 声明批次的最小与最大字节数batch.Read的语义与io.Reader完全一致读到错误即结束。管理 Topic自动创建当 broker 配置auto.create.topics.enabletruebitnami/kafka Docker 镜像中对应KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLEtrue时DialLeader会顺带创建目标 topicconn, err : kafka.DialLeader(context.Background(), tcp, localhost:9092, my-topic, 0) if err ! nil { panic(err.Error()) }显式创建若auto.create.topics.enablefalse则需要先连到任意 broker再通过Controller()找到集群 controller向 controller 发起CreateTopics请求topic : my-topic conn, err : kafka.Dial(tcp, localhost:9092) if err ! nil { panic(err.Error()) } defer conn.Close() controller, err : conn.Controller() if err ! nil { panic(err.Error()) } var controllerConn *kafka.Conn controllerConn, err kafka.Dial(tcp, net.JoinHostPort(controller.Host, strconv.Itoa(controller.Port))) if err ! nil { panic(err.Error()) } defer controllerConn.Close() topicConfigs : []kafka.TopicConfig{ { Topic: topic, NumPartitions: 1, ReplicationFactor: 1, }, } err controllerConn.CreateTopics(topicConfigs...) if err ! nil { panic(err.Error()) }通过非 leader 连接找 leader当已有一条非 leader 连接时可以借助Controller()定位 controller 地址再Dial到它适用于跨网络拓扑的场景避免对端地址不可达conn, err : kafka.Dial(tcp, localhost:9092) if err ! nil { panic(err.Error()) } defer conn.Close() controller, err : conn.Controller() if err ! nil { panic(err.Error()) } var connLeader *kafka.Conn connLeader, err kafka.Dial(tcp, net.JoinHostPort(controller.Host, strconv.Itoa(controller.Port))) if err ! nil { panic(err.Error()) } defer connLeader.Close()列出 topic用ReadPartitions()读取所有分区信息去重后即可得到 topic 列表conn, err : kafka.Dial(tcp, localhost:9092) if err ! nil { panic(err.Error()) } defer conn.Close() partitions, err : conn.ReadPartitions() if err ! nil { panic(err.Error()) } m : map[string]struct{}{} for _, p : range partitions { m[p.Topic] struct{}{} } for k : range m { fmt.Println(k) }Reader高层消费 APIReader面向“从单个 topic-partition 消费”这一典型场景自动处理重连与偏移量管理并借助 Go context 支持异步取消和超时。一个从分区 0、偏移量 42 开始消费的示例r : kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{localhost:9092, localhost:9093, localhost:9094}, Topic: topic-A, Partition: 0, MaxBytes: 10e6, // 10MB }) r.SetOffset(42) for { m, err : r.ReadMessage(context.Background()) if err ! nil { break } fmt.Printf(message at offset %d: %s %s\n, m.Offset, string(m.Key), string(m.Value)) } if err : r.Close(); err ! nil { log.Fatal(failed to close reader:, err) }务必优雅关闭Reader进程退出时必须调用Close()。Kafka 服务器需要一次优雅断开才会停止向已连接客户端推送消息若进程被 SIGINTCtrl-C或 SIGTERMdocker stop、Kubernetes 重启直接终止而未调用Close()同 topic 的新消费者新进程或新容器可能因旧会话未清理而出现连接延迟。生产代码应使用signal.Notify捕获退出信号并关闭 reader。消费指定时间范围内的消息SetOffsetAt可以把读取位置定位到某个历史时间点配合消息自身的Time字段即可实现时间窗消费例如回放过去一小时的数据startTime : time.Now().Add(-time.Hour) endTime : time.Now() batchSize : int(10e6) // 10MB r : kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{localhost:9092, localhost:9093, localhost:9094}, Topic: my-topic1, Partition: 0, MaxBytes: batchSize, }) r.SetOffsetAt(context.Background(), startTime) for { m, err : r.ReadMessage(context.Background()) if err ! nil { break } if m.Time.After(endTime) { break } // TODO: process message fmt.Printf(message at offset %d: %s %s\n, m.Offset, string(m.Key), string(m.Value)) } if err : r.Close(); err ! nil { log.Fatal(failed to close reader:, err) }ReaderConfig 关键字段以下字段及其默认值均可从 reader.go 的ReaderConfig源码确认字段作用默认值Brokersbroker 地址列表用于发现集群必填GroupID消费者组 ID设置后不能再设置Partition空GroupTopics组内多 topic 消费仅配合GroupID使用空Topic/Partition单 topic 单分区消费Partition与GroupID互斥空 / 0Dialer自定义拨号器TLS/SASL 在此注入默认拨号器QueueCapacity内部消息队列容量100MinBytes向 broker 声明的最小批次字节数低流量 topic 上设太大会延迟投递1MaxBytes向 broker 声明的最大批次字节数单条超限消息会被截断应大于最大消息体1MBMaxWait等待新数据的最大时长10sReadBatchTimeout从批次中取消息的超时10sReadLagInterval更新 lag 的间隔设负数可关闭 lag 上报—GroupBalancers组内分区分配策略优先级列表[Range, RoundRobin]HeartbeatInterval消费者组心跳间隔3sCommitInterval偏移量提交间隔为 0 时同步提交0PartitionWatchInterval/WatchPartitionChanges轮询并感知分区新增触发组内重平衡5s / falseSessionTimeout无心跳多久后 coordinator 判定消费者死亡并重平衡30sRebalanceTimeout重平衡时等待成员加入的时长30sJoinGroupBackoff加组失败后的重试间隔5sLogger/ErrorLogger运行日志与错误日志无消费者组与偏移量管理使用消费者组在ReaderConfig中指定GroupID即启用消费者组偏移量由 broker 托管r : kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{localhost:9092, localhost:9093, localhost:9094}, GroupID: consumer-group-id, Topic: topic-A, MaxBytes: 10e6, // 10MB }) for { m, err : r.ReadMessage(context.Background()) if err ! nil { break } fmt.Printf(message at topic/partition/offset %v/%v/%v: %s %s\n, m.Topic, m.Partition, m.Offset, string(m.Key), string(m.Value)) } if err : r.Close(); err ! nil { log.Fatal(failed to close reader:, err) }使用消费者组后ReadMessage会自动提交偏移量。同时存在以下已知限制(*Reader).SetOffset会返回错误组内偏移由 broker 管理不允许手动设置(*Reader).Offset恒返回-1(*Reader).Lag恒返回-1(*Reader).ReadLag会返回错误(*Reader).Stats返回的Partition为-1。显式提交偏移量若希望“处理成功后提交”改用FetchMessageCommitMessages组合实现 at-least-once 语义ctx : context.Background() for { m, err : r.FetchMessage(ctx) if err ! nil { break } fmt.Printf(message at topic/partition/offset %v/%v/%v: %s %s\n, m.Topic, m.Partition, m.Offset, string(m.Key), string(m.Value)) if err : r.CommitMessages(ctx, m); err ! nil { log.Fatal(failed to commit messages:, err) } }提交规则同一分区内CommitMessages提交的最高偏移量即为该分区的已提交偏移——例如FetchMessage取回分区内偏移 1、2、3 的消息只提交偏移 3 的消息也会一并提交 1 和 2。异步批量提交默认CommitMessages同步提交。对吞吐敏感的场景可在ReaderConfig中设置CommitInterval周期性批量刷入r : kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{localhost:9092, localhost:9093, localhost:9094}, GroupID: consumer-group-id, Topic: topic-A, MaxBytes: 10e6, // 10MB CommitInterval: time.Second, // flushes commits to Kafka every second })Writer高层生产 API大多数生产场景应使用Writer而非底层Conn它额外提供出错时自动重试与重连可配置的分区间消息分布策略Balancer同步或异步写入基于 context 的异步取消关闭时冲刷积压消息支持优雅停机发布前自动创建缺失的 topicAllowAutoTopicCreation注意该行为在 v0.4.30 之前是默认开启的。基础用法least-bytes 分布策略w : kafka.Writer{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), Topic: topic-A, Balancer: kafka.LeastBytes{}, } err : w.WriteMessages(context.Background(), kafka.Message{Key: []byte(Key-A), Value: []byte(Hello World!)}, kafka.Message{Key: []byte(Key-B), Value: []byte(One!)}, kafka.Message{Key: []byte(Key-C), Value: []byte(Two!)}, ) if err ! nil { log.Fatal(failed to write messages:, err) } if err : w.Close(); err ! nil { log.Fatal(failed to close writer:, err) }kafka.TCP(host:port, ...)用于构造AddrLeastBytes会把消息发送到当前负载最小的分区。发布前自动创建 topic设置AllowAutoTopicCreation: true后WriteMessages会在目标 topic 缺失时先尝试创建。由于创建后 leader 选举需要时间官方示例使用带超时 context 的重试循环处理LeaderNotAvailable与context.DeadlineExceededw : kafka.Writer{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), Topic: topic-A, AllowAutoTopicCreation: true, } messages : []kafka.Message{ {Key: []byte(Key-A), Value: []byte(Hello World!)}, {Key: []byte(Key-B), Value: []byte(One!)}, {Key: []byte(Key-C), Value: []byte(Two!)}, } var err error const retries 3 for i : 0; i retries; i { ctx, cancel : context.WithTimeout(context.Background(), 10*time.Second) defer cancel() // attempt to create topic prior to publishing the message err w.WriteMessages(ctx, messages...) if errors.Is(err, kafka.LeaderNotAvailable) || errors.Is(err, context.DeadlineExceeded) { time.Sleep(time.Millisecond * 250) continue } if err ! nil { log.Fatalf(unexpected error %v, err) } break } if err : w.Close(); err ! nil { log.Fatal(failed to close writer:, err) }同时写入多个 topic不设置Writer.Topic时可以逐条在Message.Topic上指定目标w : kafka.Writer{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), // NOTE: When Topic is not defined here, each Message must define it instead. Balancer: kafka.LeastBytes{}, } err : w.WriteMessages(context.Background(), // NOTE: Each Message has Topic defined, otherwise an error is returned. kafka.Message{Topic: topic-A, Key: []byte(Key-A), Value: []byte(Hello World!)}, kafka.Message{Topic: topic-B, Key: []byte(Key-B), Value: []byte(One!)}, kafka.Message{Topic: topic-C, Key: []byte(Key-C), Value: []byte(Two!)}, ) if err ! nil { log.Fatal(failed to write messages:, err) } if err : w.Close(); err ! nil { log.Fatal(failed to close writer:, err) }注意两种模式互斥设置了Writer.Topic就不要再在消息上设置Message.Topic反之亦然。Writer检测到这种歧义会直接返回错误。WriterConfig 关键字段来自 writer.go 的WriterConfig源码字段作用默认值Brokersbroker 地址列表为空时创建 Writer 会 panic必填Topic统一写入的 topic不设时每条消息必须自带Topic空Dialer自定义拨号器默认拨号器Balancer分区分布策略round-robinMaxAttempts单条消息最大投递尝试次数10QueueCapacity已废弃v0.4 后改为内存聚合批次—BatchSize发往单分区的目标批次消息数100BatchBytes单次请求最大字节数1048576BatchTimeout不完整批次的最大冲刷间隔1sReadTimeout/WriteTimeout读写超时10sRequiredAcks需要多少副本确认默认 -1 表示等全部副本-1Async设为 true 时WriteMessages不阻塞且忽略错误falseCompressionCodec压缩编码无Logger/ErrorLogger日志无与其他客户端的兼容性kafka-go 提供与主流客户端分区算法一致的 Balancer方便平滑迁移saramakafka.Hash等价于sarama.NewHashPartitionerkafka.ReferenceHash等价于sarama.NewReferenceHashPartitioner消息会路由到相同分区w : kafka.Writer{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), Topic: topic-A, Balancer: kafka.Hash{}, }librdkafka / confluent-kafka-gokafka.CRC32Balancer对应 librdkafka 默认的consistent_random策略w : kafka.Writer{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), Topic: topic-A, Balancer: kafka.CRC32Balancer{}, }Java 官方客户端kafka.Murmur2Balancer对应 Java 客户端默认分区器Java 中可直接指定 partitionkafka-go 不允许w : kafka.Writer{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), Topic: topic-A, Balancer: kafka.Murmur2Balancer{}, }消息压缩在Writer上设置Compression字段即可启用w : kafka.Writer{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), Topic: topic-A, Compression: kafka.Snappy, }Reader会通过检查消息属性自动识别压缩格式无需手动解压。补充说明v0.4 之前的版本要求显式导入压缩包以注册 codecv0.4 起压缩包导入已成为空操作默认即可读写压缩消息。TLS 与 SASL 安全连接TLS/SASL 通过DialerConn/Reader或TransportWriter/Client注入。若 Kafka 集群启用了 TLS 而客户端未配置 TLS通常会表现为难以排查的io.ErrUnexpectedEOF错误。TLS 配置Conn使用 Dialerdialer : kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, TLS: tls.Config{/* ...tls config... */}, } conn, err : dialer.DialContext(ctx, tcp, localhost:9093)Reader在 ReaderConfig 中指定 Dialerdialer : kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, TLS: tls.Config{/* ...tls config... */}, } r : kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{localhost:9092, localhost:9093, localhost:9094}, GroupID: consumer-group-id, Topic: topic-A, Dialer: dialer, })Writer直接构造或 NewWriter 两种方式// 直接构造 w : kafka.Writer{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), Topic: topic-A, Balancer: kafka.Hash{}, Transport: kafka.Transport{ TLS: tls.Config{}, }, } // 使用 kafka.NewWriter注意NewWriter 与 WriterConfig 已废弃未来版本会移除 dialer : kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, TLS: tls.Config{/* ...tls config... */}, } w : kafka.NewWriter(kafka.WriterConfig{ Brokers: []string{localhost:9092, localhost:9093, localhost:9094}, Topic: topic-A, Balancer: kafka.Hash{}, Dialer: dialer, })SASL 认证类型Plainsasl/plainmechanism : plain.Mechanism{ Username: username, Password: password, }SCRAMsasl/scrammechanism, err : scram.Mechanism(scram.SHA512, username, password) if err ! nil { panic(err) }Conn 使用 SASLmechanism, err : scram.Mechanism(scram.SHA512, username, password) if err ! nil { panic(err) } dialer : kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, SASLMechanism: mechanism, } conn, err : dialer.DialContext(ctx, tcp, localhost:9093)Reader 使用 SASLmechanism, err : scram.Mechanism(scram.SHA512, username, password) if err ! nil { panic(err) } dialer : kafka.Dialer{ Timeout: 10 * time.Second, DualStack: true, SASLMechanism: mechanism, } r : kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{localhost:9092, localhost:9093, localhost:9094}, GroupID: consumer-group-id, Topic: topic-A, Dialer: dialer, })Writer / Client 使用 SASL共享 TransportTransport 负责连接池等资源管理最佳实践是创建少量 Transport 并在应用内共享mechanism, err : scram.Mechanism(scram.SHA512, username, password) if err ! nil { panic(err) } sharedTransport : kafka.Transport{ SASL: mechanism, } w : kafka.Writer{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), Topic: topic-A, Balancer: kafka.Hash{}, Transport: sharedTransport, } client : kafka.Client{ Addr: kafka.TCP(localhost:9092, localhost:9093, localhost:9094), Timeout: 10 * time.Second, Transport: sharedTransport, }日志与可观测性在创建Reader/Writer时传入Logger与ErrorLoggerkafka.LoggerFunc适配普通日志函数即可获得内部运行可见性func logf(msg string, a ...interface{}) { fmt.Printf(msg, a...) fmt.Println() } // Reader r : kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{localhost:9092, localhost:9093, localhost:9094}, Topic: my-topic1, Partition: 0, Logger: kafka.LoggerFunc(logf), ErrorLogger: kafka.LoggerFunc(logf), }) // Writer w : kafka.Writer{ Addr: kafka.TCP(localhost:9092), Topic: topic, Logger: kafka.LoggerFunc(logf), ErrorLogger: kafka.LoggerFunc(logf), }此外Writer.Stats与Reader.Stats会返回带metric标签的结构化统计如kafka.writer.write.count、kafka.writer.message.count、kafka.writer.error.count等见 writer.go可接入 Prometheus 等指标系统。本地测试针对 Kafka 2.3.1 及更新版本某些历史用例可能因协议细微变化失败可设置环境变量KAFKA_SKIP_NETTEST1跳过网络相关测试用 Docker 在本地启动 Kafka项目自带 docker-compose.ymldocker-compose up -d运行测试KAFKA_VERSION2.3.1 \ KAFKA_SKIP_NETTEST1 \ go test -race ./...或清理缓存后执行go clean -cache make test小结kafka-go 以纯 Go、标准库风格的设计填补了 Go 生态 Kafka 客户端的空白底层Conn提供对协议的完整掌控Reader/Writer覆盖绝大多数生产消费场景消费者组、显式/批量偏移量提交、多 topic 写入、与 sarama/librdkafka/Java 兼容的分区算法、TLS/SASL 安全连接一应俱全。在 OpenCloud 项目中它作为间接依赖 vendored 在 vendor/github.com/segmentio/kafka-go版本 v0.4.51go.mod。需要接入 Kafka 消息流时建议优先使用Reader/Writer高层 API并在进程退出时确保优雅关闭连接与提交偏移量。【免费下载链接】opencloud️ OpenCloud is the open source platform for file management, sharing and collaboration. Simple and sovereign.项目地址: https://gitcode.com/GitHub_Trending/op/opencloud创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表