go: kafka 將group設置為最新
- 2019 年 11 月 21 日
- 筆記
有時,在確保group當前沒有consumer的情況下,可以將這個group的偏移設置成最新,以保證下次啟動時,group能從最新的消息消費。 程式碼:
func initKafka() (err error) { defer checkerr.MarkPanic(&err) topics := []string{"hello"} groupId := "superpig" addr := []string{"127.0.0.1:9092"} client, err := sarama.NewClient(addr, nil) checkerr.CheckError(err) partionIds, err := client.Partitions(topics[0]) checkerr.CheckError(err) defer client.Close() time.Sleep(time.Second) config := cluster.NewConfig() config.Consumer.Return.Errors = true config.Group.Return.Notifications = true config.Version = sarama.V2_0_0_0 config.Consumer.Offsets.Initial = sarama.OffsetNewest consumer, err := cluster.NewConsumer(addr, groupId, topics, config) defer consumer.Close() // 這兩個go 必不可少,否則不正常,我也不知道為啥 go func() { for err := range consumer.Errors() { logging.Errorf("consumer.Error: groupId:%s:Error: %s;topic:%vn", groupId, err.Error(), topics[0]) } }() go func() { for ntf := range consumer.Notifications() { logging.Infof("consumer.Notification: groupId:%s Rebalanced: %+v;topic:%vn", groupId, ntf, topics[0]) } }() time.Sleep(time.Second * 5) for id := range partionIds { lastoffset, err := client.GetOffset(topics[0], int32(id), sarama.OffsetNewest) checkerr.CheckError(err) // 必須 lastoffset - 1,否則offset被設置成0 logging.Infof("lastoffset:%v", lastoffset) consumer.MarkPartitionOffset(topics[0], int32(id), lastoffset-1, "") err = consumer.CommitOffsets() checkerr.CheckError(err) } return nil }