Spring Kafka手动提交下AcknowledgingMessageListener报错排查
问题场景
在Spring Kafka中配置MANUAL_IMMEDIATE ack模式实现消费者手动提交偏移量,通过实现AcknowledgingMessageListener<String, String>接口编写自定义消息监听器,结合MethodKafkaListenerEndpoint完成监听器端点注册。
应用启动过程无任何异常,但消费者拉取到Kafka消息触发消费逻辑时,立即抛出java.lang.UnsupportedOperationException: Container should never call this错误,应用进程不会终止,错误仅在消息接收后触发,容器按照默认错误处理策略多次重试消费失败的偏移量,达到最大重试次数后持续打印错误日志,消息始终无法正常消费。
自定义监听器实现代码
import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.config.KafkaListenerEndpoint; import org.springframework.kafka.config.MethodKafkaListenerEndpoint; import org.springframework.kafka.listener.AcknowledgingMessageListener; import org.springframework.kafka.listener.MessageListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.lang.Nullable; import org.springframework.stereotype.Component; import java.util.concurrent.CompletableFuture; @Component public class MyCustomMessageListener extends CustomMessageListener { @Override @SneakyThrows public KafkaListenerEndpoint createKafkaListenerEndpoint(String name, String topic) { MethodKafkaListenerEndpoint<String, String> kafkaListenerEndpoint = createDefaultMethodKafkaListenerEndpoint(name, topic); kafkaListenerEndpoint.setBean(new MyMessageListener()); kafkaListenerEndpoint.setMethod(MyMessageListener.class.getMethod("onMessage", ConsumerRecord.class)); return kafkaListenerEndpoint; } @Slf4j private static class MyMessageListener implements AcknowledgingMessageListener<String, String> { /** * Invoked with data from kafka. * @param acknowledgment ack * @param data the data to be processed. */ @Override public void onMessage(ConsumerRecord<String, String> data, Acknowledgment acknowledgment) { log.info("My message listener got a new record: " + data); acknowledgment.acknowledge(); log.info("My message listener done processing record: " + data); } } }
应用错误日志
Caused by: java.lang.UnsupportedOperationException: Container should never call this at org.springframework.kafka.listener.AcknowledgingMessageListener.onMessage(AcknowledgingMessageListener.java:44) ~[spring-kafka-2.8.6.jar:2.8.6] at jdk.internal.reflect.GeneratedMethodAccessor35.invoke(Unknown Source) ~[na:na] at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[na:na] at java.base/java.lang.reflect.Method.invoke(Method.java:566) ~[na:na] at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.doInvoke(InvocableHandlerMethod.java:169) ~[spring-messaging-5.3.13.jar:5.3.13] at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.invoke(InvocableHandlerMethod.java:119) ~[spring-messaging-5.3.13.jar:5.3.13] at org.springframework.kafka.listener.adapter.HandlerAdapter.invoke(HandlerAdapter.java:56) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter.invokeHandler(MessagingMessageListenerAdapter.java:347) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter.onMessage(RecordMessagingMessageListenerAdapter.java:92) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter.onMessage(RecordMessagingMessageListenerAdapter.java:53) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeOnMessage(KafkaMessageListenerContainer.java:2645) ~[spring-kafka-2.8.6.jar:2.8.6] ... 11 common frames omitted 2022-06-07 11:15:28.908 INFO 32664 --- [antopicPE-0-C-1] o.a.k.clients.consumer.KafkaConsumer : [Consumer clientId=consumer-cleantopicPE-1, groupId=cleantopicPE] Seeking to offset 61 for partition consumerpe-0 2022-06-07 11:15:28.909 ERROR 32664 --- [antopicPE-0-C-1] o.s.k.l.KafkaMessageListenerContainer : Error handler threw an exception org.springframework.kafka.KafkaException: Seek to current after exception; nested exception is org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method 'public default void org.springframework.kafka.listener.AcknowledgingMessageListener.onMessage(org.apache.kafka.clients.consumer.ConsumerRecord<K, V>)' threw exception; nested exception is java.lang.UnsupportedOperationException: Container should never call this; nested exception is java.lang.UnsupportedOperationException: Container should never call this at org.springframework.kafka.listener.SeekUtils.seekOrRecover(SeekUtils.java:208) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.DefaultErrorHandler.handleRemaining(DefaultErrorHandler.java:133) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeErrorHandler(KafkaMessageListenerContainer.java:2682) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeRecordListener(KafkaMessageListenerContainer.java:2563) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeWithRecords(KafkaMessageListenerContainer.java:2433) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:2311) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:1982) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeIfHaveRecords(KafkaMessageListenerContainer.java:1366) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1357) ~[spring-kafka-2.8.6.jar:2.8.6] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1252) ~[spring-kafka-2.8.6.jar:2.8.6] at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) ~[na:na] at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) ~[na:na] at java.base/java.lang.Thread.run(Thread.java:834) ~[na:na] [剩余日志堆栈保持原始内容不变]
根因分析
- 反射获取监听器方法的代码存在错误:
setMethod传入的是MyMessageListener.class.getMethod("onMessage", ConsumerRecord.class),仅匹配单参数onMessage方法,最终拿到的是AcknowledgingMessageListener接口内置的默认单参数方法,这个方法的内置逻辑就是直接抛出UnsupportedOperationException: Container should never call this,完全没有执行自定义的双参数消费逻辑。 - 错误堆栈可直接印证问题:报错指向的方法是
public default void org.springframework.kafka.listener.AcknowledgingMessageListener.onMessage(org.apache.kafka.clients.consumer.ConsumerRecord<K, V>),就是接口自带的默认抛错方法,并非自行实现的接收ConsumerRecord+Acknowledgment双参数的方法。 MethodKafkaListenerEndpoint会严格按照传入的Method对象的参数签名做适配,传入单参数方法时,Spring会将其识别为不需要Ack参数的普通监听器,和配置的MANUAL_IMMEDIATE手动ack模式逻辑不匹配,最终触发异常。
修复方案
修改反射获取方法的代码,匹配双参数签名的onMessage方法即可:
// 错误写法(删除) // kafkaListenerEndpoint.setMethod(MyMessageListener.class.getMethod("onMessage", ConsumerRecord.class)); // 正确写法 kafkaListenerEndpoint.setMethod(MyMessageListener.class.getMethod("onMessage", ConsumerRecord.class, Acknowledgment.class));
修复后验证要点:
- 容器会正确识别监听器的手动ack参数签名,消息投递时会自动注入
Acknowledgment对象,不会再调用到接口的默认抛错方法。 - 提前确认监听器容器工厂的
ackMode已设置为ContainerProperties.AckMode.MANUAL_IMMEDIATE,和监听器内手动调用acknowledge()提交偏移量的逻辑匹配。
内容的提问来源于stack exchange,提问作者None for Nothing
相关产品推荐
相关产品推荐

