代码拉取完成,页面将自动刷新
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
// }
此处可能存在不合适展示的内容,页面不予展示。您可通过相关编辑功能自查并修改。
如您确认内容无涉及 不当用语 / 纯广告导流 / 暴力 / 低俗色情 / 侵权 / 盗版 / 虚假 / 无价值内容或违法国家有关法律法规的内容,可点击提交进行申诉,我们将尽快为您处理。