【问题标题】:Kafka Consumer Unable To Resolve Listener Method Intermittently卡夫卡消费者无法间歇性地解析侦听器方法
【发布时间】:2021-08-28 13:40:46
【问题描述】:

我在 Kafka 消费者方面一直面临以下异常。令人惊讶的是,这个问题不一致,旧版本的代码(具有完全相同的配置,但有一些新的不相关功能)按预期工作。任何人都可以帮助确定可能导致此问题的原因吗?

[ERROR][938f3c68-f481-4224-b2c6-43af5fb27ada-0-C-1][org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer] - Error handler threw an exception
org.springframework.kafka.KafkaException: Seek to current after exception; nested exception is org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method could not be invoked with the incoming message
Endpoint handler details:
Method [public void com.mycompany.listener.KafkaBatchListener.onMessage(java.lang.Object,org.springframework.kafka.support.Acknowledgment)]
Bean [com.mycompany.listener.KafkaBatchListener@7a59780b]; nested exception is org.springframework.messaging.handler.invocation.MethodArgumentResolutionException: Could not resolve method parameter at index 0 in public void com.mycompany.listener.KafkaBatchListener.onMessage(java.util.List<org.apache.kafka.clients.consumer.ConsumerRecord<K, V>>,org.springframework.kafka.support.Acknowledgment): Could not resolve parameter [0] in public void com.mycompany.listener.KafkaBatchListener.onMessage(java.util.List<org.apache.kafka.clients.consumer.ConsumerRecord<K, V>>,org.springframework.kafka.support.Acknowledgment): No suitable resolver, failedMessage=GenericMessage [payload=[[B@21bc784f, MyPOJO(), [B@33bb5851], headers={kafka_offset=[4046, 4047, 4048], kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@4871203f, kafka_timestampType=[CREATE_TIME, CREATE_TIME, CREATE_TIME], kafka_receivedPartitionId=[0, 0, 0], kafka_receivedMessageKey=[[B@295620f1, MyPOJOKey(id=0), [B@5d3d6361], kafka_batchConvertedHeaders=[{myFirstHeader=[B@1f011689, myUUIDHeader=[B@7691bce8, myMetadataHeader=[B@6e585b63, myRequestIdHeader=[B@58c81ba2, myMetricsHeader=[B@4f6aeb6c, myTargetHeader=[B@34677895}, {myUUIDHeader=[B@1848ae39, myMetadataHeader=[B@c5b399, myRequestIdHeader=[B@186c1966, myMetricsHeader=[B@1740692e, myTargetHeader=[B@4a242499}, {myUUIDHeader=[B@67d01f3f, myMetadataHeader=[B@1f0f9d8a, myRequestIdHeader=[B@b928e5c, isLastMessage=[B@6079735b, myMetricsHeader=[B@7b7b18c, myTargetHeader=[B@64378f3d}], kafka_receivedTopic=[my_topic, my_topic, my_topic], kafka_receivedTimestamp=[1623420136620, 1623420137255, 1623420137576], kafka_acknowledgment=Acknowledgment for org.apache.kafka.clients.consumer.ConsumerRecords@7bc81d89, kafka_groupId=dev-consumer-grp}]
    at org.springframework.kafka.listener.SeekToCurrentBatchErrorHandler.handle(SeekToCurrentBatchErrorHandler.java:77) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.RecoveringBatchErrorHandler.handle(RecoveringBatchErrorHandler.java:124) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.ContainerAwareBatchErrorHandler.handle(ContainerAwareBatchErrorHandler.java:56) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeBatchErrorHandler(KafkaMessageListenerContainer.java:2010) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeBatchListener(KafkaMessageListenerContainer.java:1854) [spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeBatchListener(KafkaMessageListenerContainer.java:1720) [spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:1699) [spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeIfHaveRecords(KafkaMessageListenerContainer.java:1272) [spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1264) [spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1161) [spring-kafka-2.7.1.jar:2.7.1]
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) [?:?]
    at java.util.concurrent.FutureTask.run(FutureTask.java:264) [?:?]
    at java.lang.Thread.run(Thread.java:834) [?:?]
Caused by: org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method could not be invoked with the incoming message
Endpoint handler details:
Method [public void com.mycompany.listener.KafkaBatchListener.onMessage(java.lang.Object,org.springframework.kafka.support.Acknowledgment)]
Bean [com.mycompany.listener.KafkaBatchListener@7a59780b]; nested exception is org.springframework.messaging.handler.invocation.MethodArgumentResolutionException: Could not resolve method parameter at index 0 in public void com.mycompany.listener.KafkaBatchListener.onMessage(java.util.List<org.apache.kafka.clients.consumer.ConsumerRecord<K, V>>,org.springframework.kafka.support.Acknowledgment): Could not resolve parameter [0] in public void com.mycompany.listener.KafkaBatchListener.onMessage(java.util.List<org.apache.kafka.clients.consumer.ConsumerRecord<K, V>>,org.springframework.kafka.support.Acknowledgment): No suitable resolver, failedMessage=GenericMessage [payload=[[B@21bc784f, MyPOJO(), [B@33bb5851], headers={kafka_offset=[4046, 4047, 4048], kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@4871203f, kafka_timestampType=[CREATE_TIME, CREATE_TIME, CREATE_TIME], kafka_receivedPartitionId=[0, 0, 0], kafka_receivedMessageKey=[[B@295620f1, MyPOJOKey(id=0), [B@5d3d6361], kafka_batchConvertedHeaders=[{myFirstHeader=[B@1f011689, myUUIDHeader=[B@7691bce8, myMetadataHeader=[B@6e585b63, myRequestIdHeader=[B@58c81ba2, myMetricsHeader=[B@4f6aeb6c, myTargetHeader=[B@34677895}, {myUUIDHeader=[B@1848ae39, myMetadataHeader=[B@c5b399, myRequestIdHeader=[B@186c1966, myMetricsHeader=[B@1740692e, myTargetHeader=[B@4a242499}, {myUUIDHeader=[B@67d01f3f, myMetadataHeader=[B@1f0f9d8a, myRequestIdHeader=[B@b928e5c, isLastMessage=[B@6079735b, myMetricsHeader=[B@7b7b18c, myTargetHeader=[B@64378f3d}], kafka_receivedTopic=[my_topic, my_topic, my_topic], kafka_receivedTimestamp=[1623420136620, 1623420137255, 1623420137576], kafka_acknowledgment=Acknowledgment for org.apache.kafka.clients.consumer.ConsumerRecords@7bc81d89, kafka_groupId=dev-consumer-grp}]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.decorateException(KafkaMessageListenerContainer.java:2367) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeBatchOnMessage(KafkaMessageListenerContainer.java:2003) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeBatchOnMessageWithRecordsOrList(KafkaMessageListenerContainer.java:1973) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeBatchOnMessage(KafkaMessageListenerContainer.java:1925) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeBatchListener(KafkaMessageListenerContainer.java:1837) ~[spring-kafka-2.7.1.jar:2.7.1]
    ... 8 more
Caused by: org.springframework.messaging.handler.invocation.MethodArgumentResolutionException: Could not resolve method parameter at index 0 in public void com.mycompany.listener.KafkaBatchListener.onMessage(java.util.List<org.apache.kafka.clients.consumer.ConsumerRecord<K, V>>,org.springframework.kafka.support.Acknowledgment): Could not resolve parameter [0] in public void com.mycompany.listener.KafkaBatchListener.onMessage(java.util.List<org.apache.kafka.clients.consumer.ConsumerRecord<K, V>>,org.springframework.kafka.support.Acknowledgment): No suitable resolver
    at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.getMethodArgumentValues(InvocableHandlerMethod.java:145) ~[spring-messaging-5.2.12.RELEASE.jar:5.2.12.RELEASE]
    at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.invoke(InvocableHandlerMethod.java:116) ~[spring-messaging-5.2.12.RELEASE.jar:5.2.12.RELEASE]
    at org.springframework.kafka.listener.adapter.HandlerAdapter.invoke(HandlerAdapter.java:56) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter.invokeHandler(MessagingMessageListenerAdapter.java:339) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter.invoke(BatchMessagingMessageListenerAdapter.java:180) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter.onMessage(BatchMessagingMessageListenerAdapter.java:172) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter.onMessage(BatchMessagingMessageListenerAdapter.java:61) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeBatchOnMessage(KafkaMessageListenerContainer.java:1983) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeBatchOnMessageWithRecordsOrList(KafkaMessageListenerContainer.java:1973) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeBatchOnMessage(KafkaMessageListenerContainer.java:1925) ~[spring-kafka-2.7.1.jar:2.7.1]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeBatchListener(KafkaMessageListenerContainer.java:1837) ~[spring-kafka-2.7.1.jar:2.7.1]
    ... 8 more

我的应用使用以下内容:

  1. 自定义侦听器类com.mycompany.listener.KafkaBatchListener&lt;K, V&gt; 实现 org.springframework.kafka.listener.BatchAcknowledgingMessageListener&lt;K, V&gt;覆盖 onMessage(List&lt;ConsumerRecord&lt;K, V&gt;&gt; consumerRecords, Acknowledgment acknowledgment) 带有自定义标记注释@MyKafkaListener
  2. 一个自定义容器工厂,扩展 org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory&lt;K, V&gt; 并配置 setConsumerFactory(consumerFactory)setBatchErrorHandler(errorHandler)setBatchListener(true)ContainerProperties.setOnlyLogRecordMetadata(true)
  3. 一个SpringBoot @Configuration 类,实现 org.springframework.kafka.annotation.KafkaListenerConfigurer 并负责配置org.springframework.kafka.core.DefaultKafkaConsumerFactory&lt;K, V&gt;org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory&lt;K, V&gt;org.springframework.kafka.config.MethodKafkaListenerEndpoint&lt;String, String&gt;(由@MyKafkaListener 使用)
  4. 春季卡夫卡 2.7.1

其他查询: 即使设置了ContainerProperties.setOnlyLogRecordMetadata(true),异常堆栈跟踪仍然包含我省略的完整有效负载。知道为什么吗?

提前致谢!


更新:

  1. KafkaBatchListener
package com.mycompany.listener;

import java.util.List;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.listener.BatchAcknowledgingMessageListener;
import org.springframework.kafka.support.Acknowledgment;

public class KafkaBatchListener<K, V> implements BatchAcknowledgingMessageListener<K, V> {

    @Override
    @com.mycompany.listener.KafkaListener
    public void onMessage(final List<ConsumerRecord<K, V>> consumerRecords, final Acknowledgment acknowledgment) {

        // process batch using MyService<K, V>.process(consumerRecords)
        acknowledgment.acknowledge();
    }
}

  1. 自定义注释
package com.mycompany.listener;

import static java.lang.annotation.ElementType.METHOD;
import static java.lang.annotation.RetentionPolicy.RUNTIME;

import java.lang.annotation.Retention;
import java.lang.annotation.Target;

@Retention(RUNTIME)
@Target(METHOD)
public @interface KafkaListener {

}
  1. 监听器容器工厂
package com.mycompany.factory;

import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.listener.ContainerProperties;

import com.mycompany.errorhandler.ListenerContainerRecoveringBatchErrorHandler;

public class KafkaBatchListenerContainerFactory<K, V>
        extends ConcurrentKafkaListenerContainerFactory<K, V> {

    public KafkaBatchListenerContainerFactory(final DefaultKafkaConsumerFactory<K, V> consumerFactory,
            final ListenerContainerRecoveringBatchErrorHandler errorHandler, final int concurrency) {

        super.setConsumerFactory(consumerFactory);
        super.setBatchErrorHandler(errorHandler);
        super.setConcurrency(concurrency);
        super.setBatchListener(true);
        super.setAutoStartup(true);

        final ContainerProperties containerProperties = super.getContainerProperties();
        containerProperties.setAckMode(ContainerProperties.AckMode.MANUAL);
        containerProperties.setOnlyLogRecordMetadata(true);
    }

}
  1. 批处理错误处理程序
package com.mycompany.errorhandler;

import org.springframework.beans.factory.annotation.Value;
import org.springframework.kafka.listener.RecoveringBatchErrorHandler;
import org.springframework.stereotype.Component;
import org.springframework.util.backoff.FixedBackOff;

@Component
public class ListenerContainerRecoveringBatchErrorHandler extends RecoveringBatchErrorHandler {

    public ListenerContainerRecoveringBatchErrorHandler(
            @Value("${spring.kafka.consumer.properties.backOffMS:0}") final int backOffTimeMS,
            @Value("${spring.kafka.consumer.properties.retries:3}") final int retries) {

        super(new FixedBackOff(backOffTimeMS, retries));
    }

}
  1. Kafka 侦听器配置器
package com.mycompany.config;

import java.lang.reflect.Method;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import java.util.UUID;

import org.springframework.beans.factory.BeanFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.KafkaListenerConfigurer;
import org.springframework.kafka.config.KafkaListenerEndpointRegistrar;
import org.springframework.kafka.config.MethodKafkaListenerEndpoint;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.messaging.handler.annotation.support.DefaultMessageHandlerMethodFactory;

import com.mycompany.errorhandler.ListenerContainerRecoveringBatchErrorHandler;
import com.mycompany.factory.KafkaBatchListenerContainerFactory;
import com.mycompany.listener.KafkaBatchListener;

@Configuration
public class KafkaBatchListenerConfigurer<K, V> implements KafkaListenerConfigurer {

    private final List<KafkaBatchListener<K, V>> listeners;
    private final BeanFactory beanFactory;
    private final ListenerContainerRecoveringBatchErrorHandler errorHandler;
    private final int concurrency;

    @Autowired
    public KafkaBatchListenerConfigurer(final List<KafkaBatchListener<K, V>> listeners, final BeanFactory beanFactory,
            final ListenerContainerRecoveringBatchErrorHandler errorHandler,
            @Value("${spring.kafka.listener.concurrency:1}") final int concurrency) {
        this.listeners = listeners;
        this.beanFactory = beanFactory;
        this.errorHandler = errorHandler;
        this.concurrency = concurrency;
    }

    @Override
    public void configureKafkaListeners(final KafkaListenerEndpointRegistrar registrar) {

        final Method listenerMethod = lookUpBatchListenerMethod();

        listeners.forEach(listener -> {
            registerListenerEndpoint(listener, listenerMethod, registrar);
        });
    }

    private void registerListenerEndpoint(final KafkaBatchListener<K, V> listener, final Method listenerMethod,
            final KafkaListenerEndpointRegistrar registrar) {

        // final Map<String, Object> consumerConfig = get ConsumerConfig from a custom provider;
        registrar.setContainerFactory(createContainerFactory(consumerConfig));
        registrar.registerEndpoint(createListenerEndpoint(listener, listenerMethod, consumerConfig));
    }

    private KafkaBatchListenerContainerFactory<K, V> createContainerFactory(final Map<String, Object> consumerConfig) {

        final DefaultKafkaConsumerFactory<K, V> consumerFactory = new DefaultKafkaConsumerFactory<>(consumerConfig);

        final KafkaBatchListenerContainerFactory<K, V> containerFactory = new KafkaBatchListenerContainerFactory<>(
                consumerFactory, errorHandler, concurrency);

        return containerFactory;
    }

    private MethodKafkaListenerEndpoint<String, String> createListenerEndpoint(final KafkaBatchListener<K, V> listener,
            final Method listenerMethod, final Map<String, Object> consumerConfig) {

        final MethodKafkaListenerEndpoint<String, String> endpoint = new MethodKafkaListenerEndpoint<>();
        endpoint.setId(UUID.randomUUID().toString());
        endpoint.setBean(listener);
        endpoint.setMethod(listenerMethod);
        endpoint.setBeanFactory(beanFactory);
        endpoint.setGroupId("my-group-id");
        endpoint.setMessageHandlerMethodFactory(new DefaultMessageHandlerMethodFactory());

        // final String topicName = get TopicName for this key-value from a custom utility;
        endpoint.setTopics(topicName);

        final Properties properties = new Properties();
        properties.putAll(consumerConfig);
        endpoint.setConsumerProperties(properties);

        return endpoint;
    }

    private Method lookUpBatchListenerMethod() {
        return Arrays.stream(com.mycompany.listener.KafkaBatchListener.class.getMethods())
                .filter(m -> m.isAnnotationPresent(com.mycompany.listener.KafkaListener.class))
                .findAny()
                .orElseThrow(() -> new IllegalStateException(
                        String.format("[%s] class should have at least 1 method with [%s] annotation.",
                                com.mycompany.listener.KafkaBatchListener.class.getCanonicalName(),
                                com.mycompany.listener.KafkaListener.class.getCanonicalName())));
    }

}

【问题讨论】:

  • @GaryRussell 你介意看看这个问题吗?非常感谢!
  • 如果你的监听器已经实现了BatchAcknowledgingMessageListener,为什么还要用@KafkaListener注释呢? @KafkaListener 用于调用 POJO 端点;如果实现接口,直接传入容器即可;请显示您的侦听器和配置,以便我能更多地了解您要做什么。
  • @GaryRussell 我需要一个使用泛型的监听器。通过这种方式,我可以避免在 Spring 的 @KafkaListener 中对主题具有不同值的重复侦听器,并为每个键值对创建侦听器 bean。这些键值对与单个主题相关联,并且可以通过它们自己的键值通用服务实现来处理。使用代码配置更新问题。

标签: spring-boot apache-kafka kafka-consumer-api spring-kafka


【解决方案1】:

当您的侦听器已经实现消息侦听器接口之一时,您不需要所有标准的@KafkaListener 方法调用基础结构;而不是为每个监听器注册端点,只需从工厂为每个监听器创建一个容器并将监听器添加到容器属性中。

val container = containerFactory.createContainer("topic1");
container.getContainerProperties().set...
...
container.getContainerProperies().setMessageListener(myListenerInstance);
...
container.start();

【讨论】:

  • 你的意思是说我完全删除KafkaBatchListenerConfigurer.createListenerEndpoint()并在KafkaBatchListenerConfigurer.createContainerFactory()中添加建议的代码?
  • 是的;你不需要所有这些东西。它更简单,代码也更少。
  • Gary - 是否可以创建一个我自己的静态线程池并在所有主题的侦听器/容器之间共享它并在处理结束时仍然确认?尽量确保下游系统不会因大量数据而过载,因为每个主题都有自己的线程池,这意味着随着主题的增加,线程会越来越多。
  • 否;每个容器(并发)使用一个专用线程。该框架并非旨在跨侦听器共享线程;如果某个侦听器的线程饥饿时间超过max.poll.interval.ms,它将导致重新平衡。您可以(在某种程度上)通过设置idleBetweenPolls 来控制吞吐量。与max.poll.records一起,您可以控制每个集装箱的发货速度。您还可以暂停/恢复容器,这将暂停记录交付,同时让消费者保持活力。
  • Gary - 抱歉在 cmets 中继续查询... 需要使用某些记录,例如 5(依次发布到同一主题分区,标题指定 1of5、2of5 等)在同一批次中。这是为了确保它们都保存在下游数据库中以确保一致性。如果第 1 批仅消耗前 2 条消息(由于 max.poll.records 已满),则提交偏移量并且 JVM 崩溃;不完整的数据被写入。有没有办法避免这种情况?
猜你喜欢
  • 2019-03-27
  • 1970-01-01
  • 1970-01-01
  • 2019-07-03
  • 2018-05-05
  • 2021-08-22
  • 1970-01-01
  • 2021-04-02
  • 1970-01-01
相关资源
最近更新 更多