【问题标题】:Kafka Listeners stop reading from topics after a few hours几个小时后,Kafka 听众停止阅读主题
【发布时间】:2021-07-26 09:14:45
【问题描述】:

我一直在开发的一个应用程序开始在我们的暂存和生产环境中引起问题,这似乎是由于 Kafka 听众在应用程序启动几个小时后不再从分配的主题中读取任何内容。

该应用程序在云铸造环境中运行,它有 13 个@KafkaListener,根据给定的模式从多个主题中读取。主题的数量是相等的(应用程序上的每个用户使用该模式为 13 个听众中的每一个创建自己的主题)。主题有 3 个分区。还使用了自动缩放,至少有 2 个应用程序实例同时运行。其中一个主题的负载比其他主题更重,每秒接收 1 到 200 条消息。每条消息的处理时间很短,因为我们接收批次,而处理部分只继续将批次写入数据库。

如前所述,当前的问题是它在启动后工作了一段时间,然后突然听众不再接收消息。日志中没有明显的错误或警告。创建了一个临时端点,其中 KafkaListenerEndpointRegistry 用于查看侦听器容器,它们似乎都在运行并分配了适当的分区。对容器执行 .stop() 和 .start() 会导致处理一批额外的消息,然后不会再处理其他任何事情。

以下是使用的配置:

@Bean
public ConsumerFactory<String, String> consumerFactory(){
    return new DefaultKafkaConsumerFactory<>(kafkaConfig.getConfiguration());
}

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(){
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setBatchListener(true);
    factory.setConcurrency(3);
    factory.getContainerProperties().setPollTimeout(5000);
    factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
}

kafkaConfig 设置如下设置:

PARTITION_ASSIGNMENT_STRATEGY_CONFIG: RoundRobinAssignor
MAX_POLL_INTERVAL_MS_CONFIG: 60000
MAX_POLL_RECORDS_CONFIG: 10
MAX_PARTITION_FETCH_BYTES_CONFIG: Integer.MAX_VALUE
ENABLE_AUTO_COMMIT_CONFIG: false
METADATA_MAX_AGE_CONFIG: 15000
REQUEST_TIMEOUT_MS_CONFIG: 30000
HEARTBEAT_INTERVAL_MS_CONFIG: 15000
SESSION_TIMEOUT_MS_CONFIG: 60000

另外,每个监听器都在自己的类中,并且监听方法如下:

@KafkaListener(id="<patternName>-container", topicPattern = "<patternName>.*", groupId = "<patternName>Group")
public void listen(@Payload List<String> payloads,
                   @Header(KafkaHeaders.RECEIVED_TOPIC) String topics,
                   Acknowledgement acknowledgement){
    //processPayload...
    acknowledgement.acknowledge();
}

spring-kakfa版本是2.7.4。

此配置是否存在可以解决问题的问题?我最近尝试了多次更改但没有成功,更改这些配置设置,在类级别移动 @KafkaListener 注释,在停止读取时重新启动侦听器容器,甚至异步完成对消息的所有处理并确认消息当它们被侦听器方法拾取时。没有错误或警告日志,由于每秒发送的消息量,我无法看到任何有助于调试日志的内容。我们还有另一个应用程序在相同的环境中运行相同的设置,但只有 3 个侦听器(不同的主题模式),不会出现此问题。它处于类似的负载下,因为这 3 个侦听器接收到的消息正在输出到主题,从而导致存在问题的应用程序负载过大。

我非常感谢任何帮助或指示我还能做什么,因为这个问题严重阻碍了我们的生产。如果我错过了可以提供帮助的内容,请告诉我。

谢谢。

【问题讨论】:

  • 大多数此类问题是由于侦听器线程卡在用户代码中;发生这种情况时进行线程转储以查看线程在做什么。
  • 感谢您的提示,在对线程转储进行彻底分析后,我发现问题确实是由于最近添加的另一个服务导致线程无限期地停止处理。修复后,问题就消失了。

标签: spring-kafka


【解决方案1】:

大多数此类问题是由于侦听器线程卡在用户代码中;发生这种情况时进行线程转储以查看线程在做什么。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2012-12-20
    • 1970-01-01
    • 2022-12-13
    • 2014-09-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多