Sarama v1.61.0 发布:消费者支持协同再平衡,新增 KIP-848 ConsumerGroupDescribe API

Sarama v1.61.0 发布:消费者支持协同再平衡,新增 KIP-848 ConsumerGroupDescribe API

3小时前
duguying
13 阅读
来源:GitHub Discussions

IBM 维护的 Go Kafka 客户端库 Sarama 发布 v1.61.0,新增 cooperative-sticky 协同再平衡分配器与 KIP-848 ConsumerGroupDescribe API(API key 69),并修复 broker 注销、roundrobin 死循环、异步生产者竞态等一批稳定性问题。

相关标签

#Sarama#Kafka#Go#IBM

发布概览

IBM 维护的 Go 语言 Kafka 客户端库 Sarama 正式发布 v1.61.0(2026-09-22)。本次版本的绝对主角是消费者再平衡协议:新增协同再平衡(Cooperative Rebalancing)支持,补齐 KIP-848 定义的 ConsumerGroupDescribe API,并集中修复了生产者、消费者与协议层的一批稳定性问题。

这一版共合入 18 个 PR,其中 8 位贡献者是首次为 Sarama 提交代码。

核心亮点

1. 消费者支持协同再平衡(Cooperative Rebalancing)

feat(consumer): support cooperative rebalancing(#3696,由维护者 dnwe 提交)为消费者组引入了新的 cooperative-sticky 分配器。两者最大的差别在于再平衡期间的会话生命周期:

  • eager(原有行为):任何一次再平衡都会终止当前 session 与这次 Consume 调用,所有分区先全部交还再重新分配,组内出现全量停顿。调用方必须在无限循环里反复调用 Consume 才能拿到新会话。
  • cooperative(新增):再平衡变成对现有 session 的增量更新——保留的 claim 继续消费,只有被回收的分区会关闭其 Messages() 通道;Setup() 在首次加入时执行一次,Cleanup() 在成员离组时才执行;Consume 调用与 session context 跨越多次再平衡保持存活。

实现细节上,新增的 consumeCooperative / reconcileAssignment 会对比新旧分配并做差集(diffClaims),仅撤销被回收的 claim,随后在 Consumer.Group.Rebalance.Timeout 上限内等待对应的 ConsumeClaim 返回;超时会以 ErrRebalanceTimedOut 触发。若撤销 claim 导致需要二次加入组,会在同一 session 内自动做一次 follow-up join。

示例程序(#3760)同步支持了两种新取值,并用 cooperative-sticky-upgrade 给出了滚动升级的安全路径——先同时提供两个分配器,让老实例和新实例能在过渡期内协商,全部升级完成后再切换为纯 cooperative:

// 直接使用协同再平衡
config.Consumer.Group.Rebalance.GroupStrategies = []sarama.BalanceStrategy{
    sarama.NewBalanceStrategyCooperativeSticky(),
}

// 滚动升级过渡期:同时提供 cooperative-sticky 与 range
config.Consumer.Group.Rebalance.GroupStrategies = []sarama.BalanceStrategy{
    sarama.NewBalanceStrategyCooperativeSticky(),
    sarama.NewBalanceStrategyRange(),
}

示例还新增了 -assignor 参数:range, roundrobin, sticky, cooperative-sticky, cooperative-sticky-upgrade

2. 新增 KIP-848 ConsumerGroupDescribe API(API key 69)

feat: support KIP-848 ConsumerGroupDescribe API(#3748,首次贡献者 shilohlee98)实现了 Kafka 下一代消费者再平衡协议 KIP-848 中的消费者组描述接口:

  • 新增 API key 69 的 v0 / v1 请求与响应编解码器
  • 新增 Broker.ConsumerGroupDescribe 方法,并接入版本协商、mocks 与单元测试;
  • 支持 Kafka 3.7+,在 Kafka 4.0+ 上支持 v1 成员类型。

对应的需求来自 issue #3747。对于需要在客户端侧直接观察消费者组成员、状态与分配情况的使用方,这比依赖旧版 DescribeGroups 的字段语义更贴合新协议。

3. 一批稳定性修复

消费与生产两侧各修掉几个可能引发线上故障的问题:

  • Broker 注销越界fix(client): only deregister the broker registered under that ID(#3752)——按 ID 注销时不再误伤其它 broker 注册项,避免连接泄漏与元数据错乱。
  • RoundRobin 死循环fix(consumer): stop roundrobin balancer looping on an unsubscribed topic(#3740)——订阅关系变化时,roundrobin 分配器不再陷入无限循环。
  • 生产者响应排空fix(producer): drain pending responses immediately once channel closes(#3729)——通道关闭后立即排空待处理响应,避免关闭流程挂起。
  • 异步生产者竞态fix: prevent race in asyncProducer retryHandler(#3733)——消除重试处理中的 data race,对开启 -race 的 CI 是实质改善。
  • 协议层报错更明确fix(protocol): return ErrUnsupportedVersion for absent broker API keys(#3741)——broker 未暴露对应 API key 时明确返回 ErrUnsupportedVersion,而不是抛出含糊错误。

4. 依赖与工程质量

  • 依赖更新:klauspost/compress 1.19.2 → 1.20.0、pierrec/lz4/v4 4.1.29 → 4.1.30、golangci-lint v2.13.2、vearutop/teststat v0.1.29,以及 golang-x 组的 3 项更新。
  • 测试隔离:fix(test): capture logs without replacing the global Logger(#3737)——测试不再替换全局 Logger,避免并行用例相互污染。
  • 互操作验证:test(functional): cover cooperative Java interoperability(#3754)——新增功能性测试,验证 Go 客户端与 Java 客户端在协同再平衡下的互操作性。
  • 文档:docs: clarify producer batching configuration(#3712)——讲清生产者批量发送相关配置项的语义。

发布产物

  • 发布标签:v1.61.0,上一版本为 v1.60.2
  • 完整变更对比:v1.60.2...v1.61.0
  • 作为纯 Go 库,本版本以 Go Module 形式分发,无独立二进制产物
  • 首次贡献者 8 人:@shilohlee98、@hsdfat、@lenamonj、@navi89mai、@krvladislav、@dloebl、@official-burak、@wxldonnaBJ

关于 Sarama

Sarama 是 IBM 维护的、采用 MIT 许可证的 Apache Kafka Go 客户端库,也是 Go 生态中接入 Kafka 使用最广泛的选择之一。除核心库外,还提供用于测试的 mocks 子包、examples 示例目录,以及 tools 目录下用于测试、诊断与埋点的命令行工具。

兼容性方面,Sarama 提供 「2 个版本 + 2 个月」的兼容承诺:同时支持最新的两个 Kafka 与 Go 稳定版本,并为更早版本保留两个月宽限期(更旧的 Kafka 版本通常仍可工作)。API 稳定性遵循 Go 的语义化版本(semver)规范。

获取与升级

升级到 v1.61.0:

go get github.com/IBM/sarama@v1.61.0
go mod tidy

注意再平衡语义变化。 采用 cooperative 策略后,Consume 不会再因每次再平衡而返回,ConsumeClaim 必须在被回收分区的 Messages() 通道关闭时及时退出,否则会等待到 Consumer.Group.Rebalance.Timeout 并报 ErrRebalanceTimedOut。若你的现有代码把「Consume 返回」当作再平衡信号,需要按新的生命周期改写。建议先在预发环境用小流量验证,再按 cooperative-sticky-upgradecooperative-sticky 的两步走完成生产升级。

返回列表

更多资讯

想看更多优质内容?

浏览所有资讯