2 Star 0 Fork 1

fanfansong / gokafka

加入 Gitee
与超过 1200万 开发者一起发现、参与优秀开源项目,私有仓库也完全免费 :)
免费加入
该仓库未声明开源许可证文件(LICENSE),使用请关注具体项目描述及其代码上游依赖。
克隆/下载
example.go 1.36 KB
一键复制 编辑 原始数据 按行查看 历史
宋慧琳 提交于 2018-12-05 16:36 . 添加product的例子
package main
import (
"context"
"fmt"
"time"
kafka "gitee.com/fanfansong/gokafka/gokafka"
)
type ConsumerHandelr struct {
}
func (handler ConsumerHandelr) Handle(ctx context.Context, topic string, messages []byte) error {
fmt.Println("Topic:", topic)
fmt.Println("message:", string(messages))
return nil
}
func complexConsume() {
//kafkaUrl使用逗号分割
handler := ConsumerHandelr{}
consumer, err := kafka.NewConsumer("test", "192.168.150.221:9092")
if err != nil {
println(err)
return
}
groupId := "myTest"
err = consumer.AddGroup(groupId, handler)
if err != nil {
println(err)
return
}
defer consumer.Close()
ctx := context.Background()
topics := []string{"maggie_test", "test"}
go consumer.ConsumeTopicsByGroup(ctx, groupId, topics)
time.Sleep(time.Second * 10)
}
func simpleConsume() {
handler := ConsumerHandelr{}
consumer, err := kafka.NewConsumer("test", "192.168.150.221:9092")
if err != nil {
println(err)
return
}
defer consumer.Close()
ctx := context.Background()
go consumer.ConsumeTopic(ctx, "maggie_test", handler)
time.Sleep(time.Second * 10)
}
func syncProduct() {
producer, err := kafka.NewSyncProducer("192.168.150.221:9092")
if err != nil {
fmt.Println(err)
}
defer producer.Close()
producer.SendMessage("maggie_test", []byte("!!!!!!test!!!!!!!!"))
}
// func main() {
// syncProduct()
// return
// }
1
https://gitee.com/fanfansong/gokafka.git
git@gitee.com:fanfansong/gokafka.git
fanfansong
gokafka
gokafka
master

搜索帮助

53164aa7 5694891 3bd8fe86 5694891