【问题标题】:Kafka Listener rollback transaction on session timeoutKafka Listener 在会话超时时回滚事务
【发布时间】:2017-09-01 11:46:44
【问题描述】:

我正在使用 spring-kafka 版本 1.1.3 来使用来自主题的消息。在消费者配置中,自动提交设置为 truemax.poll.records10session.timeout.ms 与服务器协商为 10 秒。

收到一条消息后,我将其中的一部分保存到数据库中。我的数据库有时会很慢,这会导致 kafka 监听器的会话超时:

组 mygroup 的自动偏移提交失败:提交无法完成 因为该组已经重新平衡并将分区分配给 另一个成员。这意味着后续调用之间的时间 poll() 比配置的 session.timeout.ms 长,这 通常意味着轮询循环花费了太多时间消息 加工。您可以通过增加会话来解决这个问题 超时或通过减少 poll() 中返回的批处理的最大大小 使用 max.poll.records。

由于我无法增加服务器上的会话超时并且max.poll.records 已经下降到 10,我希望能够将我的数据库调用包装在一个事务中,如果是 kafka,它将回滚会话超时。

这可能吗?我该如何做到这一点?

很遗憾,我无法在文档中找到解决方案。

【问题讨论】:

  • 为什么要回滚数据库提交?如果您只需要将使用者设置为手动提交,然后捕获异常。
  • 什么异常?仅记录上述消息。据我所知,kafka session 超时没有抛出异常。

标签: spring apache-kafka spring-kafka


【解决方案1】:

您必须考虑升级到 Spring Kafka 1.2 和 Kafka 0.10.x。旧的 Apache Kafka 在心跳方面存在缺陷。因此,使用autoCommit 和慢速侦听器最终会导致意外的重新平衡,而您自己会遇到这样的问题。您使用的 Spring Kafka 版本的逻辑如下:

// if the container is set to auto-commit, then execute in the
// same thread
// otherwise send to the buffering queue
if (this.autoCommit) {
    invokeListener(records);
}
else {
    if (sendToListener(records)) {
        if (this.assignedPartitions != null) {
            // avoid group management rebalance due to a slow
            // consumer
            this.consumer.pause(this.assignedPartitions);
            this.paused = true;
            this.unsent = records;
        }
    }
}

因此,您可以考虑关闭autoCommit 并依赖默认开启的内置pause 功能。

【讨论】:

  • 谢谢!如果可行,我会尝试并接受答案。
  • @meva ,它对你有用吗?接受答案让其他人知道这是一个解决方案,这是一种很好的方式
【解决方案2】:

决定升级到 Kafka 0.11,因为它增加了事务支持(参见 Release Notes)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-03-30
    • 2015-02-18
    • 1970-01-01
    • 2011-04-27
    • 1970-01-01
    • 2021-05-24
    • 2012-07-21
    • 1970-01-01
    相关资源
    最近更新 更多