分享
golang kafka
simbadan · · 21196 次点击 · · 开始浏览这是一个创建于 的文章,其中的信息可能已经有所发展或是发生改变。
golang kafka – hello world
https://github.com/Shopify/sarama
https://shopify.github.io/sarama/
consumer.go
package main
import (
"fmt"
"github.com/Shopify/sarama"
"log"
"os"
"strings"
"sync"
)
var (
wg sync.WaitGroup
logger = log.New(os.Stderr, "[srama]", log.LstdFlags)
)
func main() {
sarama.Logger = logger
consumer, err := sarama.NewConsumer(strings.Split("localhost:9092", ","), nil)
if err != nil {
logger.Println("Failed to start consumer: %s", err)
}
partitionList, err := consumer.Partitions("hello")
if err != nil {
logger.Println("Failed to get the list of partitions: ", err)
}
for partition := range partitionList {
pc, err := consumer.ConsumePartition("hello", int32(partition), sarama.OffsetNewest)
if err != nil {
logger.Printf("Failed to start consumer for partition %d: %s\n", partition, err)
}
defer pc.AsyncClose()
wg.Add(1)
go func(sarama.PartitionConsumer) {
defer wg.Done()
for msg := range pc.Messages() {
fmt.Printf("Partition:%d, Offset:%d, Key:%s, Value:%s", msg.Partition, msg.Offset, string(msg.Key), string(msg.Value))
fmt.Println()
}
}(pc)
}
wg.Wait()
logger.Println("Done consuming topic hello")
consumer.Close()
}
producer.go
package main
import (
"github.com/Shopify/sarama"
"log"
"os"
"strings"
)
var (
logger = log.New(os.Stderr, "[srama]", log.LstdFlags)
)
func main() {
sarama.Logger = logger
config := sarama.NewConfig()
config.Producer.RequiredAcks = sarama.WaitForAll
config.Producer.Partitioner = sarama.NewRandomPartitioner
msg := &sarama.ProducerMessage{}
msg.Topic = "hello"
msg.Partition = int32(-1)
msg.Key = sarama.StringEncoder("key")
msg.Value = sarama.ByteEncoder("你好, 世界!")
producer, err := sarama.NewSyncProducer(strings.Split("localhost:9092", ","), config)
if err != nil {
logger.Println("Failed to produce message: %s", err)
os.Exit(500)
}
defer producer.Close()
partition, offset, err := producer.SendMessage(msg)
if err != nil {
logger.Println("Failed to produce message: ", err)
}
logger.Printf("partition=%d, offset=%d\n", partition, offset)
}
有疑问加站长微信联系(非本文作者)
入群交流(和以上内容无关):加入Go大咖交流群,或添加微信:liuxiaoyan-s 备注:入群;或加QQ群:692541889
关注微信21196 次点击
上一篇:golang环境
添加一条新回复
(您需要 后才能回复 没有账号 ?)
- 请尽量让自己的回复能够对别人有帮助
- 支持 Markdown 格式, **粗体**、~~删除线~~、
`单行代码` - 支持 @ 本站用户;支持表情(输入 : 提示),见 Emoji cheat sheet
- 图片支持拖拽、截图粘贴等方式上传
收入到我管理的专栏 新建专栏
golang kafka – hello world
https://github.com/Shopify/sarama
https://shopify.github.io/sarama/
consumer.go
package main
import (
"fmt"
"github.com/Shopify/sarama"
"log"
"os"
"strings"
"sync"
)
var (
wg sync.WaitGroup
logger = log.New(os.Stderr, "[srama]", log.LstdFlags)
)
func main() {
sarama.Logger = logger
consumer, err := sarama.NewConsumer(strings.Split("localhost:9092", ","), nil)
if err != nil {
logger.Println("Failed to start consumer: %s", err)
}
partitionList, err := consumer.Partitions("hello")
if err != nil {
logger.Println("Failed to get the list of partitions: ", err)
}
for partition := range partitionList {
pc, err := consumer.ConsumePartition("hello", int32(partition), sarama.OffsetNewest)
if err != nil {
logger.Printf("Failed to start consumer for partition %d: %s\n", partition, err)
}
defer pc.AsyncClose()
wg.Add(1)
go func(sarama.PartitionConsumer) {
defer wg.Done()
for msg := range pc.Messages() {
fmt.Printf("Partition:%d, Offset:%d, Key:%s, Value:%s", msg.Partition, msg.Offset, string(msg.Key), string(msg.Value))
fmt.Println()
}
}(pc)
}
wg.Wait()
logger.Println("Done consuming topic hello")
consumer.Close()
}
producer.go
package main
import (
"github.com/Shopify/sarama"
"log"
"os"
"strings"
)
var (
logger = log.New(os.Stderr, "[srama]", log.LstdFlags)
)
func main() {
sarama.Logger = logger
config := sarama.NewConfig()
config.Producer.RequiredAcks = sarama.WaitForAll
config.Producer.Partitioner = sarama.NewRandomPartitioner
msg := &sarama.ProducerMessage{}
msg.Topic = "hello"
msg.Partition = int32(-1)
msg.Key = sarama.StringEncoder("key")
msg.Value = sarama.ByteEncoder("你好, 世界!")
producer, err := sarama.NewSyncProducer(strings.Split("localhost:9092", ","), config)
if err != nil {
logger.Println("Failed to produce message: %s", err)
os.Exit(500)
}
defer producer.Close()
partition, offset, err := producer.SendMessage(msg)
if err != nil {
logger.Println("Failed to produce message: ", err)
}
logger.Printf("partition=%d, offset=%d\n", partition, offset)
}