01、介绍

Apache Kafka 是一款开源的消息引擎系统。它在项目中的作用主要是削峰填谷和解耦。本文我们只介绍 Apache Kafka 的 Golang 客户端库 Sarama。Sara( U J T T E jma 是 MIT 许可的b | ( v Apache Kafk5 . N R ` T } (a 0.8 及更高版本的 Golang 客户端库。

如果读者朋友对 Apache Kafka 服务端还不了解,建议先阅读官方文档中的入门部分,本文使用的版本是 Apac# . t 2 U Phe Kafka 2.8。

02、生产者

我们可以使用 Sarau k z I C 4 Zma 库的 Asyn| J 9cProducer 或 Sy( T _ }ncProducer 生产消息。在大多数情况下首选使用 AsyncProducer 生产消息。它Q | Q #通过一个 channel 接收消息,并在后台尽可能高效的异步生产消F m Q息。

SyncProducer 发送 Kafka 消息后阻塞,直到接收到 ACK 确认。SyncProducer 有两个警告:它通常效率较低,并且实际的耐用性保证取决于 Producer.R& ] * q ? dequiredAcks 的配置值。在某些配置中W Y { P : R P,有时仍会丢失由 SyncProducer 确认的消息,但是使用比较简单。

为了读者朋友们容易理解,本文我们介绍 SyncProduceU x D r . % ? X fr 作为生产者的使用方式。如果读者朋友想了解 AsyncProdq 6 B X \ : K { `ucer 作为生产者的使用方式,请参考官方文档。

使用 SyncProducer 作为生产者的示例代码:

  1. funcsendMessage(brokerAddr[]string,config*sarama.Config,topicstring,valuesarama.Encoder){
  2. producer,err:=sar] y s Sama.NewSyn! K G = E o g !cProducer(brokero T + J . G )Addr,config)
  3. iferr!=nil{
  4. fmt.Prp F Iintln(err)
  5. return
  6. }
  7. deferfunc(){
  8. iferr=producer.Close();G q 4 ] n C yerr!=nil{
  9. fL I f R - L \mt.Println(err)
  10. return
  11. }
  12. }()
  13. msg:=&sarama.ProducerMessage{
  14. Topic:topic,
  15. Value:value,
  16. }
  17. partiu k S C 8tion,offsK * X Cet,err:=producer.SendMessage(msg)
  18. iferr!=nil{
  19. fmt.Println(err)
  20. returny # #
  21. }
  22. fmt.Print` w 0 a O ` vf("partition:%do1 O g ] Tffset:%d\n",pe M \ ^ @artition,offset)
  23. }

阅读上面这段代码,我们调用 NewSyncProducer() 创建一l 7 ` 2 1 [ F个新的 SyncProK g j C $ducer,给定 broker 地址和配置信息Y 5 C { = } F 6 @。调用 SendMessJ \ Page() 生产给定的消息,并且仅在生产成功或失败时返P , g S C回。它将返回分区(Partition)和生产的消息的偏移量(Offset),如果消息生产失败@ P ; 8 4 N %,则返回错误。

需要注意的是,为了避免泄露,必须在生产者上调用 Close(),因为当它超出范围时,可X T k # 7 \能不会自动垃圾回收。

03、消费者

我们可以使用 Sarama 库的消费者 Consc L _umer 或消费者组 ConsumerGroup AK j # Z DPI 消费消息。为了读者朋友们容易理解,本文我们介绍使用 Consumer 消费消息。

Consumer 管理 PartitionConsumers,该 Partim w * Z & w 0 [tionConsumers 处理来自 brokers 的 Kafka 消息。

Consumer 消费消息的示例代码:

  1. fu: c Y , = \ H 2 ]ncconsume) , H Dr(brokenAddr[]string,topicstring,partitionint32,offsetint64){
  2. consumer,err:=sai u \ Erama.NewConsumer(brokenAddr,nil)
  3. iferr!=nil{
  4. fmt.Println(err)
  5. return
  6. }
  7. deferfunck ^ n 5(){
  8. ifer\ 6 I or=consumer.Clo` * ` \ F qse();err!=nil{
  9. fmt.Println(err)
  10. return
  11. }
  12. }()
  13. partitionCoQ P d p { 3nsumer,err:=consumer.ConsumePa: ; B M t # p Zrtition(topic,partitA 3 . ] i 0 : =ion,offset)
  14. iferr!=nil{
  15. fmt.Println(err)
  16. return
  17. }
  18. deferfunc()m k g M ( N{
  19. iferr=partitionConsumer.Close();err!=nil{
  20. fmt.Println(err)
  21. return
  22. }
  23. }()
  24. formsg:=rangepartitionConsumer.Messages(){C i ? U M V l w #
  25. fmt.Printf("partition:%doffset:%dkey:%sval:%s\n",msg.Partition,msg.Offset,msg.Key,msg.Value)
  26. }
  27. }

阅读上面这i O K u #段代码,我们调用 NewConsumer() 创建一个新的 consumer,给定 broker 地址和配置信% V :息。调用 ConsumePartition() 创建 ParY ! R = xtitionConsumer,给定 topic、partition 和 offset。PartitionConsumer 处理来自给定 topic 和 partition 的g c ( x c 0 / : Kafka 消息。

需要注意的是,为了防止泄露,% o K w g 3 w K必须调用 consumer 和 partitionConsumer 的 Close(),因为当它超出范围时,可能不会自动垃圾回q ( C E收。

04、总结

本文主要介绍如` H + \ j o Q x O何使用 Apache Kafka 的 Golang 语言客户端库 Sarama 生产和消费 Kafk0 * – @a 消~ f F ; ` K G n _息。关于生产者和6 f P ? ` E _ ; w消费者,分别列举了一个简单示例。除此之外,Sarama 库还提供了很多其它 Api,感兴趣的读者朋友可以阅读官方文档了解更多。

【责任编辑:未丽燕 TEL:(010)68476606】

点赞 0

发表评论

您的电子邮箱地址不会被公开。 必填项已用*标注