前往小程序,Get更优阅读体验!
立即前往
首页
学习
活动
专区
圈层
工具
发布
首页
学习
活动
专区
圈层
工具
社区首页 >专栏 >go: kafka 将group设置为最新

go: kafka 将group设置为最新

作者头像
超级大猪
发布2019-11-21 22:34:21
发布2019-11-21 22:34:21
2.1K00
代码可运行
举报
文章被收录于专栏:大猪的笔记大猪的笔记
运行总次数:0
代码可运行

有时,在确保group当前没有consumer的情况下,可以将这个group的偏移设置成最新,以保证下次启动时,group能从最新的消息消费。 代码:

代码语言:javascript
代码运行次数:0
运行
复制
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:%v\n",
                groupId, err.Error(),
                topics[0])
        }
    }()

    go func() {
        for ntf := range consumer.Notifications() {
            logging.Infof("consumer.Notification: groupId:%s Rebalanced: %+v;topic:%v\n",
                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
}
本文参与 腾讯云自媒体同步曝光计划,分享自作者个人站点/博客。
原始发表:2019-08-26 ,如有侵权请联系 cloudcommunity@tencent.com 删除

本文分享自 作者个人站点/博客 前往查看

如有侵权,请联系 cloudcommunity@tencent.com 删除。

本文参与 腾讯云自媒体同步曝光计划  ,欢迎热爱写作的你一起参与!

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档